featurebase/index.go
Ben Johnson 523cf5bc0e
refactor main into pilosa.Server
This commit refactors most of the code in `cmd/pilosa` to
`pilosa.Server` so that it can be reused in long running cluster
testing.
2016-05-13 14:39:09 -06:00

250 lines
5.3 KiB
Go

package pilosa
import (
"errors"
"fmt"
"os"
"path/filepath"
"sort"
"sync"
)
// Index represents a container for fragments.
type Index struct {
mu sync.Mutex
remoteMax uint64
// Databases by name.
dbs map[string]*DB
// Data directory path.
Path string
}
// NewIndex returns a new instance of Index.
func NewIndex() *Index {
return &Index{
dbs: make(map[string]*DB),
remoteMax: 0,
}
}
// Open initializes the root data directory for the index.
func (i *Index) Open() error {
if err := os.MkdirAll(i.Path, 0777); err != nil {
return err
}
// Open path to read all database directories.
f, err := os.Open(i.Path)
if err != nil {
return err
}
defer f.Close()
fis, err := f.Readdir(0)
if err != nil {
return err
}
for _, fi := range fis {
if !fi.IsDir() {
continue
}
db := NewDB(i.DBPath(filepath.Base(fi.Name())), filepath.Base(fi.Name()))
if err := db.Open(); err != nil {
return fmt.Errorf("open db: name=%s, err=%s", db.Name(), err)
}
i.dbs[db.Name()] = db
}
return nil
}
// Close closes all open fragments.
func (i *Index) Close() error {
for _, db := range i.dbs {
db.Close()
}
return nil
}
// SliceN returns the highest slice across all frames.
func (i *Index) SliceN() uint64 {
i.mu.Lock()
defer i.mu.Unlock()
sliceN := i.remoteMax
for _, db := range i.dbs {
if n := db.SliceN(); n > sliceN {
sliceN = n
}
}
return sliceN
}
// Schema returns schema data for all databases and frames.
func (i *Index) Schema() []*DBInfo {
var a []*DBInfo
for _, db := range i.DBs() {
di := &DBInfo{Name: db.Name()}
for _, frame := range db.Frames() {
di.Frames = append(di.Frames, &FrameInfo{
Name: frame.Name(),
})
}
a = append(a, di)
}
return a
}
// DBPath returns the path where a given database is stored.
func (i *Index) DBPath(name string) string { return filepath.Join(i.Path, name) }
// DB returns the database by name.
func (i *Index) DB(name string) *DB {
i.mu.Lock()
defer i.mu.Unlock()
return i.db(name)
}
func (i *Index) db(name string) *DB { return i.dbs[name] }
// DBs returns a list of all databases in the index.
func (i *Index) DBs() []*DB {
i.mu.Lock()
defer i.mu.Unlock()
a := make([]*DB, 0, len(i.dbs))
for _, db := range i.dbs {
a = append(a, db)
}
sort.Sort(dbSlice(a))
return a
}
// CreateDBIfNotExists returns a database by name.
// The database is created if it does not already exist.
func (i *Index) CreateDBIfNotExists(name string) (*DB, error) {
i.mu.Lock()
defer i.mu.Unlock()
return i.createDBIfNotExists(name)
}
func (i *Index) createDBIfNotExists(name string) (*DB, error) {
if name == "" {
return nil, errors.New("database name required")
}
// Return database if it exists.
if db := i.db(name); db != nil {
return db, nil
}
// Otherwise create a new database.
db := NewDB(i.DBPath(name), name)
if err := db.Open(); err != nil {
return nil, err
}
i.dbs[db.Name()] = db
return db, nil
}
// Frame returns the frame for a database and name.
func (i *Index) Frame(db, name string) *Frame {
d := i.DB(db)
if d == nil {
return nil
}
return d.Frame(name)
}
// CreateFrameIfNotExists returns the frame for a database & name.
// The frame is created if it doesn't already exist.
func (i *Index) CreateFrameIfNotExists(db, name string) (*Frame, error) {
d, err := i.CreateDBIfNotExists(db)
if err != nil {
return nil, err
}
return d.CreateFrameIfNotExists(name)
}
// Fragment returns the fragment for a database, frame & slice.
func (i *Index) Fragment(db, frame string, slice uint64) *Fragment {
f := i.Frame(db, frame)
if f == nil {
return nil
}
return f.fragment(slice)
}
// CreateFragmentIfNotExists returns the fragment for a database, frame & slice.
// The fragment is created if it doesn't already exist.
func (i *Index) CreateFragmentIfNotExists(db, frame string, slice uint64) (*Fragment, error) {
f, err := i.CreateFrameIfNotExists(db, frame)
if err != nil {
return nil, err
}
return f.CreateFragmentIfNotExists(slice)
}
func (i *Index) SetMax(newmax uint64) {
i.mu.Lock()
defer i.mu.Unlock()
i.remoteMax = newmax
}
// IndexSyncer is an active anti-entropy tool that compares the local index
// with a remote index based on block checksums and resolves differences.
type IndexSyncer struct {
Index *Index
Host string
Cluster *Cluster
}
// SyncIndex compares the index on host with the local index and resolves differences.
func (s *IndexSyncer) SyncIndex() error {
sliceN := s.Index.SliceN()
// Iterate over schema in sorted order.
for _, di := range s.Index.Schema() {
for _, fi := range di.Frames {
for slice := uint64(0); slice <= sliceN; slice++ {
// Ignore slices that this host doesn't own.
if !s.Cluster.OwnsSlice(s.Host, slice) {
continue
}
// Sync fragment if own it.
if err := s.syncFragment(di.Name, fi.Name, slice); err != nil {
return fmt.Errorf("sync error: db=%s, frame=%s, slice=%d, err=%s", di.Name, fi.Name, slice, err)
}
}
}
}
return nil
}
// syncFragment synchronizes
func (s *IndexSyncer) syncFragment(db, frame string, slice uint64) error {
// Ensure fragment exists locally.
f, err := s.Index.CreateFragmentIfNotExists(db, frame, slice)
if err != nil {
return err
}
// Sync fragments together.
fs := FragmentSyncer{
Fragment: f,
Host: s.Host,
Cluster: s.Cluster,
}
if err := fs.SyncFragment(); err != nil {
return err
}
return nil
}