diff --git a/cluster_internal_test.go b/cluster_internal_test.go index b9e714326..059053423 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -96,7 +96,7 @@ func newIndexWithTempPath(name string) *Index { if err != nil { panic(err) } - index, err := NewIndex(path, name, DefaultPartitionN) + index, err := NewIndex(NewHolder(DefaultPartitionN), path, name) if err != nil { panic(err) } diff --git a/field.go b/field.go index b746e1800..f82f03543 100644 --- a/field.go +++ b/field.go @@ -31,7 +31,6 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/v2/internal" - "github.com/pilosa/pilosa/v2/logger" "github.com/pilosa/pilosa/v2/pql" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/stats" @@ -117,10 +116,6 @@ type Field struct { // Shards with data on any node in the cluster, according to this node. remoteAvailableShards *roaring.Bitmap - logger logger.Logger - - snapshotQueue snapshotQueue - translateStore TranslateStore // Instantiates new translation stores @@ -338,17 +333,17 @@ func OptFieldTypeBool() FieldOption { // that it's of the type `OptFieldType*`). This means // this function couldn't be used to set, for example, // `FieldOptions.Keys`. -func NewField(path, index, name string, opts FieldOption) (*Field, error) { +func NewField(holder *Holder, path, index, name string, opts FieldOption) (*Field, error) { err := validateName(name) if err != nil { return nil, errors.Wrap(err, "validating name") } - return newField(path, index, name, opts) + return newField(holder, path, index, name, opts) } // newField returns a new instance of field (without name validation). -func newField(path, index, name string, opts FieldOption) (*Field, error) { +func newField(holder *Holder, path, index, name string, opts FieldOption) (*Field, error) { // Apply functional option. fo := FieldOptions{} err := opts(&fo) @@ -372,7 +367,7 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) { remoteAvailableShards: roaring.NewBitmap(), - logger: logger.NopLogger, + holder: holder, OpenTranslateStore: OpenInMemTranslateStore, } @@ -448,7 +443,7 @@ func (f *Field) loadAvailableShards() error { } // some other problem: if err != nil { - f.logger.Printf("available shards file present but unreadable, discarding: %v", err) + f.holder.Logger.Printf("available shards file present but unreadable, discarding: %v", err) err = os.Remove(path) if err != nil { return errors.Wrap(err, "deleting corrupt available shards list") @@ -457,7 +452,7 @@ func (f *Field) loadAvailableShards() error { } bm := roaring.NewBitmap() if err = bm.UnmarshalBinary(buf); err != nil { - f.logger.Printf("available shards file corrupt, discarding: %v", err) + f.holder.Logger.Printf("available shards file corrupt, discarding: %v", err) err = os.Remove(path) if err != nil { return errors.Wrap(err, "deleting corrupt available shards list") @@ -548,17 +543,17 @@ func (f *Field) Options() FieldOptions { func (f *Field) Open() error { if err := func() (err error) { // Ensure the field's path exists. - f.logger.Debugf("ensure field path exists: %s", f.path) + f.holder.Logger.Debugf("ensure field path exists: %s", f.path) if err := os.MkdirAll(f.path, 0777); err != nil { return errors.Wrap(err, "creating field dir") } - f.logger.Debugf("load meta file for index/field: %s/%s", f.index, f.name) + f.holder.Logger.Debugf("load meta file for index/field: %s/%s", f.index, f.name) if err := f.loadMeta(); err != nil { return errors.Wrap(err, "loading meta") } - f.logger.Debugf("load available shards for index/field: %s/%s", f.index, f.name) + f.holder.Logger.Debugf("load available shards for index/field: %s/%s", f.index, f.name) if err := f.loadAvailableShards(); err != nil { return errors.Wrap(err, "loading available shards") } @@ -570,17 +565,17 @@ func (f *Field) Open() error { } // Apply the field options loaded from meta (or set via setOptions()). - f.logger.Debugf("apply options for index/field: %s/%s", f.index, f.name) + f.holder.Logger.Debugf("apply options for index/field: %s/%s", f.index, f.name) if err := f.applyOptions(f.options); err != nil { return errors.Wrap(err, "applying options") } - f.logger.Debugf("open views for index/field: %s/%s", f.index, f.name) + f.holder.Logger.Debugf("open views for index/field: %s/%s", f.index, f.name) if err := f.openViews(); err != nil { return errors.Wrap(err, "opening views") } - f.logger.Debugf("open row attribute store for index/field: %s/%s", f.index, f.name) + f.holder.Logger.Debugf("open row attribute store for index/field: %s/%s", f.index, f.name) if err := f.rowAttrStore.Open(); err != nil { return errors.Wrap(err, "opening attrstore") } @@ -607,7 +602,7 @@ func (f *Field) Open() error { return err } - f.logger.Debugf("successfully opened field index/field: %s/%s", f.index, f.name) + f.holder.Logger.Debugf("successfully opened field index/field: %s/%s", f.index, f.name) return nil } func blockingWriteAvailableShards(fieldPath string, availableShardBytes []byte) { @@ -748,7 +743,7 @@ fileLoop: <-fieldQueue }() name := filepath.Base(fi.Name()) - f.logger.Debugf("open index/field/view: %s/%s/%s", f.index, f.name, fi.Name()) + f.holder.Logger.Debugf("open index/field/view: %s/%s/%s", f.index, f.name, fi.Name()) view := f.newView(f.viewPath(name), name) if err := view.open(); err != nil { return fmt.Errorf("opening view: view=%s, err=%s", view.name, err) @@ -770,7 +765,7 @@ fileLoop: } view.rowAttrStore = f.rowAttrStore - f.logger.Debugf("add index/field/view to field.viewMap: %s/%s/%s", f.index, f.name, view.name) + f.holder.Logger.Debugf("add index/field/view to field.viewMap: %s/%s/%s", f.index, f.name, view.name) mu.Lock() f.viewMap[view.name] = view mu.Unlock() @@ -1183,14 +1178,10 @@ func (f *Field) createViewIfNotExistsBase(name string) (*view, bool, error) { } func (f *Field) newView(path, name string) *view { - view := newView(path, f.index, f.name, name, f.options) - view.logger = f.logger + view := newView(f.holder, path, f.index, f.name, name, f.options) view.rowAttrStore = f.rowAttrStore view.stats = f.Stats view.broadcaster = f.broadcaster - if f.snapshotQueue != nil { - view.snapshotQueue = f.snapshotQueue - } return view } diff --git a/field_internal_test.go b/field_internal_test.go index db4a3b5e7..296e54be0 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -205,7 +205,7 @@ func NewTestField(t *testing.T, opts FieldOption) *TestField { if err != nil { t.Fatal(err) } - field, err := NewField(path, "i", "f", opts) + field, err := NewField(NewHolder(DefaultPartitionN), path, "i", "f", opts) if err != nil { t.Fatal(err) } @@ -235,7 +235,7 @@ func (f *TestField) Reopen() error { } path, index, name := f.Path(), f.Index(), f.Name() - f.Field, err = NewField(path, index, name, OptFieldTypeDefault()) + f.Field, err = NewField(NewHolder(DefaultPartitionN), path, index, name, OptFieldTypeDefault()) if err != nil { return err } @@ -730,7 +730,7 @@ func TestDecimalField_MinMaxBoundaries(t *testing.T) { }, } { t.Run("minmax"+strconv.Itoa(i), func(t *testing.T) { - _, err := NewField("no-path", "i", "f", OptFieldTypeDecimal(test.scale, test.min, test.max)) + _, err := NewField(NewHolder(DefaultPartitionN), "no-path", "i", "f", OptFieldTypeDecimal(test.scale, test.min, test.max)) if err != nil && test.expErr { if !strings.Contains(err.Error(), "is not supported") { t.Fatal(err) diff --git a/field_test.go b/field_test.go index 40b64728a..ab24b466f 100644 --- a/field_test.go +++ b/field_test.go @@ -144,7 +144,7 @@ func TestField_NameRestriction(t *testing.T) { if err != nil { panic(err) } - field, err := pilosa.NewField(path, "i", ".meta", pilosa.OptFieldTypeDefault()) + field, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", ".meta", pilosa.OptFieldTypeDefault()) if field != nil { t.Fatalf("unexpected field name %s", err) } @@ -177,13 +177,13 @@ func TestField_NameValidation(t *testing.T) { panic(err) } for _, name := range validFieldNames { - _, err := pilosa.NewField(path, "i", name, pilosa.OptFieldTypeDefault()) + _, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", name, pilosa.OptFieldTypeDefault()) if err != nil { t.Fatalf("unexpected field name: %s %s", name, err) } } for _, name := range invalidFieldNames { - _, err := pilosa.NewField(path, "i", name, pilosa.OptFieldTypeDefault()) + _, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", name, pilosa.OptFieldTypeDefault()) if err == nil { t.Fatalf("expected error on field name: %s", name) } diff --git a/fragment.go b/fragment.go index 3258d6f00..f70e9e71c 100644 --- a/fragment.go +++ b/fragment.go @@ -106,6 +106,10 @@ type fragment struct { field string view string shard uint64 + + // parent holder, used to find snapshot queue, etc. + holder *Holder + // debugging tool: addresses of current and previous maps prevdata, currdata struct{ from, to uintptr } @@ -152,12 +156,10 @@ type fragment struct { mutexVector vector stats stats.StatsClient - - snapshotQueue snapshotQueue } // newFragment returns a new instance of Fragment. -func newFragment(path, index, field, view string, shard uint64, flags byte) *fragment { +func newFragment(holder *Holder, path, index, field, view string, shard uint64, flags byte) *fragment { f := &fragment{ path: path, index: index, @@ -168,11 +170,10 @@ func newFragment(path, index, field, view string, shard uint64, flags byte) *fra CacheType: DefaultCacheType, CacheSize: DefaultCacheSize, - Logger: logger.NopLogger, + holder: holder, MaxOpN: defaultFragmentMaxOpN, - stats: stats.NopStatsClient, - snapshotQueue: defaultSnapshotQueue, + stats: stats.NopStatsClient, } f.snapshotCond = sync.Cond{L: &f.mu} return f @@ -188,13 +189,13 @@ func (f *fragment) Open() error { if err := func() error { // Initialize storage in a function so we can close if anything goes wrong. - f.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) + f.holder.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) if err := f.openStorage(true); err != nil { return errors.Wrap(err, "opening storage") } // Fill cache with rows persisted to disk. - f.Logger.Debugf("open cache for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) + f.holder.Logger.Debugf("open cache for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) if err := f.openCache(); err != nil { e2 := f.closeStorage() if e2 != nil { @@ -214,7 +215,7 @@ func (f *fragment) Open() error { return err } - f.Logger.Debugf("successfully opened index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) + f.holder.Logger.Debugf("successfully opened index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) return nil } @@ -273,7 +274,7 @@ func (f *fragment) importStorage(data []byte, file *os.File, newGen generation, } return false, fmt.Errorf("unmarshal storage: file=%s, err=%s", file.Name(), err) } - f.Logger.Printf("warning: unmarshal storage, file=%s, err=%v", file.Name(), err) + f.holder.Logger.Printf("warning: unmarshal storage, file=%s, err=%v", file.Name(), err) trunc, ok := cause.(roaring.FileShouldBeTruncatedError) if ok { // generation code looks for a FileShouldBeTruncatedError @@ -299,7 +300,7 @@ func (f *fragment) applyStorage(data []byte, file *os.File, newGen generation, m if file != nil { fi, err := file.Stat() if err != nil { - f.Logger.Printf("trying to apply new storage to existing bitmap, stat failed: %v", err) + f.holder.Logger.Printf("trying to apply new storage to existing bitmap, stat failed: %v", err) } if err == nil && fi != nil && fi.Size() == 0 { return f.emptyStorage(file) @@ -359,7 +360,7 @@ func (f *fragment) openStorage(unmarshalData bool) error { storageOp = f.applyStorage } var err error - f.gen, err = newGeneration(f.gen, f.path, unmarshalData, storageOp, f.Logger) + f.gen, err = newGeneration(f.gen, f.path, unmarshalData, storageOp, f.holder.Logger) if f.gen != nil { scratchData := f.gen.Bytes() f.prevdata = f.currdata @@ -407,7 +408,7 @@ func (f *fragment) openCache() error { // Unmarshal cache data. var pb internal.Cache if err := proto.Unmarshal(buf, &pb); err != nil { - f.Logger.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err) + f.holder.Logger.Printf("error unmarshaling cache data, skipping: path=%s, err=%s", path, err) return nil } @@ -435,13 +436,13 @@ func (f *fragment) Close() error { func (f *fragment) close() error { // Flush cache if closing gracefully. if err := f.flushCache(); err != nil { - f.Logger.Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path) + f.holder.Logger.Printf("fragment: error flushing cache on close: err=%s, path=%s", err, f.path) return errors.Wrap(err, "flushing cache") } // Close underlying storage. if err := f.closeStorage(); err != nil { - f.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path) + f.holder.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path) return errors.Wrap(err, "closing storage") } @@ -688,7 +689,8 @@ func (f *fragment) unprotectedSetRow(row *Row, rowID uint64) (changed bool, err f.rowCache.Add(rowID, nil) // Snapshot storage. - f.snapshotQueue.Enqueue(f) + f.holder.SnapshotQueue.Enqueue(f) + f.stats.Count("setRow", 1, 1.0) return changed, nil } @@ -728,7 +730,7 @@ func (f *fragment) unprotectedClearRow(rowID uint64) (changed bool, err error) { f.rowCache.Add(rowID, nil) // Snapshot storage. - f.snapshotQueue.Enqueue(f) + f.holder.SnapshotQueue.Enqueue(f) return changed, nil } @@ -2006,11 +2008,11 @@ func (f *fragment) importPositions(set, clear []uint64, rowSet map[uint64]struct // we got an error. it's possible that the error indicates that something went wrong. mappedIn, mappedOut, unmappedIn, errs, e2 := f.storage.SanityCheckMapping(f.currdata.from, f.currdata.to) if errs != 0 { - f.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v", + f.holder.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v", f.path, mappedIn, mappedOut, unmappedIn, errs, e2) if f.prevdata.from != f.currdata.from { mappedIn, mappedOut, unmappedIn, errs, e2 = f.storage.SanityCheckMapping(f.prevdata.from, f.prevdata.to) - f.Logger.Printf("with previous map, storage would have %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v", + f.holder.Logger.Printf("with previous map, storage would have %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total, last %v", mappedIn, mappedOut, unmappedIn, errs, e2) } } @@ -2164,7 +2166,7 @@ func (f *fragment) importValue(columnIDs []uint64, values []int64, bitDepth uint // in theory, this should probably have been queued anyway, but if enough // of the bits matched existing bits, we'll be under our opN estimate, and // we want to ensure that the snapshot happens. - return f.snapshotQueue.Immediate(f) + return f.holder.SnapshotQueue.Immediate(f) } // importRoaring imports from the official roaring data format defined at @@ -2252,7 +2254,7 @@ func (f *fragment) incrementOpN(changed int) { f.opN += changed f.ops++ if f.opN > f.MaxOpN { - f.snapshotQueue.Enqueue(f) + f.holder.SnapshotQueue.Enqueue(f) } } @@ -2286,7 +2288,7 @@ func (f *fragment) snapshot() (err error) { // we can't see the actual values that were used to generate this, probably. if e2.Error() == "runtime error: invalid memory address or nil pointer dereference" { mappedIn, mappedOut, unmappedIn, errs, _ := f.storage.SanityCheckMapping(f.currdata.from, f.currdata.to) - f.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total", + f.holder.Logger.Printf("transaction failed on %s. storage has %d mapped in range, %d mapped out of range, %d unmapped in range, %d errors total", f.path, mappedIn, mappedOut, unmappedIn, errs) } } else { @@ -2306,7 +2308,7 @@ func (f *fragment) snapshot() (err error) { func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err error) { // nolint: interfacer completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.field, f.view, f.shard) start := time.Now() - defer track(start, completeMessage, f.stats, f.Logger) + defer track(start, completeMessage, f.stats, f.holder.Logger) // Create a temporary file to snapshot to. snapshotPath := f.path + snapshotExt @@ -3096,7 +3098,7 @@ func (s *fragmentSyncer) syncBlockFromPrimary(id int) error { // the primary node. nodes := s.Cluster.shardNodes(f.index, f.shard) if s.Node.ID != nodes[0].ID { - f.Logger.Debugf("non-primary replica expecting sync from primary: %s, index=%s, field=%s, shard=%d", nodes[0].ID, f.index, f.field, f.shard) + f.holder.Logger.Debugf("non-primary replica expecting sync from primary: %s, index=%s, field=%s, shard=%d", nodes[0].ID, f.index, f.field, f.shard) return nil } diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 5912eb132..2d8c27b73 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1477,7 +1477,7 @@ func BenchmarkFragment_Blocks(b *testing.B) { } // Open the fragment specified by the path. - f := newFragment(*FragmentPath, "i", "f", viewStandard, 0, 0) + f := newFragment(NewHolder(DefaultPartitionN), *FragmentPath, "i", "f", viewStandard, 0, 0) if err := f.Open(); err != nil { b.Fatal(err) } @@ -2006,7 +2006,7 @@ func BenchmarkFragment_Snapshot(b *testing.B) { b.ReportAllocs() // Open the fragment specified by the path. - f := newFragment(*FragmentPath, "i", "f", viewStandard, 0, 0) + f := newFragment(NewHolder(DefaultPartitionN), *FragmentPath, "i", "f", viewStandard, 0, 0) if err := f.Open(); err != nil { b.Fatal(err) } @@ -2121,7 +2121,7 @@ func BenchmarkImportRoaring(b *testing.B) { // care whether this succeeds, // but if it's happening we want // it to be done. - _ = f.snapshotQueue.Await(f) + _ = defaultSnapshotQueue.Await(f) f.Clean(b) b.Fatalf("import error: %v", err) } @@ -2162,7 +2162,7 @@ func BenchmarkImportRoaringConcurrent(b *testing.B) { err := frags[j].importRoaringT(data[j], false) // error unimportant if it happened, but we want // any snapshots to have finished. - _ = frags[j].snapshotQueue.Await(frags[j]) + _ = defaultSnapshotQueue.Await(frags[j]) return err }) } @@ -2203,7 +2203,7 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) { if err != nil { b.Fatalf("importing roaring: %v", err) } - err = frags[j].snapshotQueue.Immediate(frags[j]) + err = defaultSnapshotQueue.Immediate(frags[j]) if err != nil { b.Fatalf("snapshot after import: %v", err) } @@ -2214,7 +2214,7 @@ func BenchmarkImportRoaringUpdateConcurrent(b *testing.B) { j := j eg.Go(func() error { err := frags[j].importRoaringT(updata, false) - err2 := frags[j].snapshotQueue.Await(frags[j]) + err2 := defaultSnapshotQueue.Await(frags[j]) if err == nil { err = err2 } @@ -2282,7 +2282,7 @@ func BenchmarkImportRoaringUpdate(b *testing.B) { if err != nil { b.Errorf("import error: %v", err) } - err = f.snapshotQueue.Immediate(f) + err = defaultSnapshotQueue.Immediate(f) if err != nil { b.Errorf("snapshot after import error: %v", err) } @@ -2292,7 +2292,7 @@ func BenchmarkImportRoaringUpdate(b *testing.B) { f.Clean(b) b.Errorf("import error: %v", err) } - err = f.snapshotQueue.Await(f) + err = defaultSnapshotQueue.Await(f) if err != nil { b.Errorf("snapshot after import error: %v", err) } @@ -2391,7 +2391,7 @@ func BenchmarkImportIntoLargeFragment(b *testing.B) { } origF.Close() fi.Close() - nf := newFragment(fi.Name(), "i", "f", viewStandard, 0, 0) + nf := newFragment(NewHolder(DefaultPartitionN), fi.Name(), "i", "f", viewStandard, 0, 0) err = nf.Open() if err != nil { b.Fatalf("opening fragment: %v", err) @@ -2428,7 +2428,7 @@ func BenchmarkImportRoaringIntoLargeFragment(b *testing.B) { } origF.Close() fi.Close() - nf := newFragment(fi.Name(), "i", "f", viewStandard, 0, 0) + nf := newFragment(NewHolder(DefaultPartitionN), fi.Name(), "i", "f", viewStandard, 0, 0) err = nf.Open() if err != nil { b.Fatalf("opening fragment: %v", err) @@ -2607,17 +2607,23 @@ func (f *fragment) sanityCheck(t testing.TB) { func (f *fragment) Clean(t testing.TB) { f.mu.Lock() - err := f.snapshotQueue.Await(f) - f.mu.Unlock() - if err != nil { - t.Fatalf("snapshot failed before sanity check: %v", err) - } - f.sanityCheck(t) - if f.storage != nil && f.storage.Source != nil { - if f.storage.Source.Dead() { - t.Fatalf("cleaning up fragment %s, source %s, source already dead", f.path, f.storage.Source.ID()) + // we need to ensure that we unlock the mutex before terminating + // the clean operation, but we need it held during the sanity + // check or else, in some cases, the background snapshot queue + // can decide to pick it up. + func() { + defer f.mu.Unlock() + err := defaultSnapshotQueue.Await(f) + if err != nil { + t.Fatalf("snapshot failed before sanity check: %v", err) } - } + f.sanityCheck(t) + if f.storage != nil && f.storage.Source != nil { + if f.storage.Source.Dead() { + t.Fatalf("cleaning up fragment %s, source %s, source already dead", f.path, f.storage.Source.ID()) + } + } + }() errc := f.Close() // prevent double-closes of generation during testing. f.gen = nil @@ -2626,10 +2632,6 @@ func (f *fragment) Clean(t testing.TB) { if errc != nil || errf != nil { t.Fatal("cleaning up fragment: ", errc, errf, errp) } - if f.snapshotQueue != nil { - f.snapshotQueue.Stop() - f.snapshotQueue = nil - } // not all fragments have cache files if errp != nil && !os.IsNotExist(errp) { t.Fatalf("cleaning up fragment cache: %v", errp) @@ -2649,10 +2651,6 @@ func (f *fragment) CleanKeep(t testing.TB) { if errc != nil { t.Fatal("closing fragment: ", errc, errp) } - if f.snapshotQueue != nil { - f.snapshotQueue.Stop() - f.snapshotQueue = nil - } // not all fragments have cache files if errp != nil && !os.IsNotExist(errp) { t.Fatalf("cleaning up fragment cache: %v", errp) @@ -2680,12 +2678,11 @@ func mustOpenFragmentFlags(index, field, view string, shard uint64, cacheType st cacheType = DefaultCacheType } - f := newFragment(file.Name(), index, field, view, shard, flags) + f := newFragment(NewHolder(DefaultPartitionN), file.Name(), index, field, view, shard, flags) f.CacheType = cacheType f.RowAttrStore = &memAttrStore{ store: make(map[uint64]map[string]interface{}), } - f.snapshotQueue = newSnapshotQueue(1, 1, nil) if err := f.Open(); err != nil { panic(err) @@ -3179,7 +3176,7 @@ func TestUnionInPlaceMapped(t *testing.T) { // it's used only in computation of things that usually don't go to // disk, which is why we handle this specially in testing and not // generically. - err = f.snapshotQueue.Immediate(f) + err = defaultSnapshotQueue.Immediate(f) if err != nil { t.Fatalf("snapshot after union-in-place: %v", err) } @@ -3332,7 +3329,6 @@ func TestImportClearRestart(t *testing.T) { cols: []uint64{1, 1, 1, 1, 1, 1}, }, } - for i, test := range tests { for _, maxOpN := range []int{0, 10000} { t.Run(fmt.Sprintf("%dMaxOpN%d", i, maxOpN), func(t *testing.T) { @@ -3385,7 +3381,7 @@ func TestImportClearRestart(t *testing.T) { check(t, f, exp) - f2 := newFragment(f.path, "i", "f", viewStandard, 0, 0) + f2 := newFragment(NewHolder(DefaultPartitionN), f.path, "i", "f", viewStandard, 0, 0) f2.MaxOpN = maxOpN f2.CacheType = f.CacheType @@ -3419,7 +3415,7 @@ func TestImportClearRestart(t *testing.T) { check(t, f2, exp) - f3 := newFragment(f2.path, "i", "f", viewStandard, 0, 0) + f3 := newFragment(NewHolder(DefaultPartitionN), f2.path, "i", "f", viewStandard, 0, 0) f3.MaxOpN = maxOpN f3.CacheType = f.CacheType diff --git a/generation_test.go b/generation_test.go index 75444fcf7..7df71fc98 100644 --- a/generation_test.go +++ b/generation_test.go @@ -39,14 +39,14 @@ func TestGenerationPanic(t *testing.T) { } prevData = f.gen.(*mmapGeneration).data f.mu.Lock() - _ = f.snapshotQueue.Immediate(f) + _ = defaultSnapshotQueue.Immediate(f) f.mu.Unlock() runtime.GC() for i := 0; i < (f.MaxOpN / 2); i++ { _, _ = f.setBit(0, uint64(i*32)+23) } f.mu.Lock() - f.snapshotQueue.Await(f) + defaultSnapshotQueue.Await(f) f.mu.Unlock() runtime.GC() newData := f.gen.(*mmapGeneration).data diff --git a/holder.go b/holder.go index de3cd0701..d06920cd3 100644 --- a/holder.go +++ b/holder.go @@ -77,9 +77,8 @@ type Holder struct { // The interval at which the cached row ids are persisted to disk. cacheFlushInterval time.Duration - Logger logger.Logger - - snapshotQueue snapshotQueue + Logger logger.Logger + SnapshotQueue SnapshotQueue // Instantiates new translation stores OpenTranslateStore OpenTranslateStoreFunc @@ -169,6 +168,8 @@ func NewHolder(partitionN int) *Holder { translationSyncer: NopTranslationSyncer, Logger: logger.NopLogger, + + SnapshotQueue: defaultSnapshotQueue, } } @@ -213,11 +214,6 @@ func (h *Holder) Open() error { return errors.Wrap(err, "reading directory") } - // Run snapshots asynchronously. The snapshotQueue will have a background - // task associated with it which flushes it and waits until this channel - // is closed, so we should always close this channel when done. - h.snapshotQueue = newSnapshotQueue(10, 2, h.Logger) - for _, fi := range fis { // Skip files or hidden directories. if !fi.IsDir() || strings.HasPrefix(fi.Name(), ".") { @@ -261,18 +257,25 @@ func (h *Holder) Open() error { h.Logger.Printf("open holder: complete") - // Periodically flush cache. - h.wg.Add(1) - go func() { defer h.wg.Done(); h.monitorCacheFlush() }() - h.Stats.Open() - h.snapshotQueue.ScanHolder(h) h.opened.Close() return nil } +// Activate runs the background tasks relevant to keeping a holder in a stable +// state, such as scanning it for needed snapshots, or flushing caches. This +// is separate from opening because, while a server would nearly always want +// to do this, other use cases (like consistency checks of a data directory) +// need to avoid it even getting started. +func (h *Holder) Activate() { + // Periodically flush cache. + h.wg.Add(2) + go func() { defer h.wg.Done(); h.monitorCacheFlush() }() + go func() { defer h.wg.Done(); h.SnapshotQueue.ScanHolder(h, h.closing) }() +} + // checkForeignIndex is a check before applying a foreign // index to a field; if the index is not yet available, // (because holder is still opening and may not have opened @@ -313,10 +316,6 @@ func (h *Holder) Close() error { return errors.Wrap(err, "closing index") } } - if h.snapshotQueue != nil { - h.snapshotQueue.Stop() - h.snapshotQueue = nil - } // Reset opened in case Holder needs to be reopened. h.opened.mu.Lock() @@ -587,19 +586,16 @@ func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) { } func (h *Holder) newIndex(path, name string) (*Index, error) { - index, err := NewIndex(path, name, h.partitionN) + index, err := NewIndex(h, path, name) if err != nil { return nil, err } - index.logger = h.Logger index.Stats = h.Stats.WithTags(fmt.Sprintf("index:%s", index.Name())) index.broadcaster = h.broadcaster index.newAttrStore = h.NewAttrStore index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data")) - index.snapshotQueue = h.snapshotQueue index.OpenTranslateStore = h.OpenTranslateStore index.translationSyncer = h.translationSyncer - index.holder = h return index, nil } diff --git a/index.go b/index.go index 481d22cf4..4a7ee3f5f 100644 --- a/index.go +++ b/index.go @@ -27,7 +27,6 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/v2/internal" - "github.com/pilosa/pilosa/v2/logger" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/stats" "github.com/pkg/errors" @@ -46,9 +45,6 @@ type Index struct { trackExistence bool existenceFld *Field - // Partitions used by translation. - partitionN int - // Fields by name. fields map[string]*Field @@ -60,9 +56,6 @@ type Index struct { broadcaster broadcaster Stats stats.StatsClient - logger logger.Logger - snapshotQueue snapshotQueue - // Passed to field for foreign-index lookup. holder *Holder @@ -76,24 +69,23 @@ type Index struct { } // NewIndex returns a new instance of Index. -func NewIndex(path, name string, partitionN int) (*Index, error) { +func NewIndex(holder *Holder, path, name string) (*Index, error) { err := validateName(name) if err != nil { return nil, errors.Wrap(err, "validating name") } return &Index{ - path: path, - name: name, - partitionN: partitionN, - fields: make(map[string]*Field), + path: path, + name: name, + fields: make(map[string]*Field), newAttrStore: newNopAttrStore, columnAttrs: nopStore, broadcaster: NopBroadcaster, Stats: stats.NopStatsClient, - logger: logger.NopLogger, + holder: holder, trackExistence: true, translateStores: make(map[int]TranslateStore), @@ -155,18 +147,18 @@ func (i *Index) OpenWithTimestamp() error { return i.open(true) } func (i *Index) open(withTimestamp bool) (err error) { // Ensure the path exists. - i.logger.Debugf("ensure index path exists: %s", i.path) + i.holder.Logger.Debugf("ensure index path exists: %s", i.path) if err := os.MkdirAll(i.path, 0777); err != nil { return errors.Wrap(err, "creating directory") } // Read meta file. - i.logger.Debugf("load meta file for index: %s", i.name) + i.holder.Logger.Debugf("load meta file for index: %s", i.name) if err := i.loadMeta(); err != nil { return errors.Wrap(err, "loading meta file") } - i.logger.Debugf("open fields for index: %s", i.name) + i.holder.Logger.Debugf("open fields for index: %s", i.name) if err := i.openFields(withTimestamp); err != nil { return errors.Wrap(err, "opening fields") } @@ -181,15 +173,15 @@ func (i *Index) open(withTimestamp bool) (err error) { return errors.Wrap(err, "opening attrstore") } - i.logger.Debugf("open translate store for index: %s", i.name) + i.holder.Logger.Debugf("open translate store for index: %s", i.name) var g errgroup.Group var mu sync.Mutex - for partitionID := 0; partitionID < i.partitionN; partitionID++ { + for partitionID := 0; partitionID < i.holder.partitionN; partitionID++ { partitionID := partitionID g.Go(func() error { - store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.partitionN) + store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.holder.partitionN) if err != nil { return errors.Wrapf(err, "opening index translate store: partition=%d", partitionID) } @@ -239,7 +231,7 @@ fileLoop: defer func() { <-indexQueue }() - i.logger.Debugf("open field: %s", fi.Name()) + i.holder.Logger.Debugf("open field: %s", fi.Name()) mu.Lock() fld, err := i.newField(i.fieldPath(filepath.Base(fi.Name())), filepath.Base(fi.Name())) if withTimestamp { @@ -257,7 +249,7 @@ fileLoop: if err := fld.Open(); err != nil { return fmt.Errorf("open field: name=%s, err=%s", fld.Name(), err) } - i.logger.Debugf("add field to index.fields: %s", fi.Name()) + i.holder.Logger.Debugf("add field to index.fields: %s", fi.Name()) mu.Lock() i.fields[fld.Name()] = fld mu.Unlock() @@ -512,17 +504,13 @@ func (i *Index) createField(name string, opt *FieldOptions) (*Field, error) { } func (i *Index) newField(path, name string) (*Field, error) { - f, err := newField(path, i.name, name, OptFieldTypeDefault()) + f, err := newField(i.holder, path, i.name, name, OptFieldTypeDefault()) if err != nil { return nil, err } - f.logger = i.logger f.Stats = i.Stats f.broadcaster = i.broadcaster f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data")) - if i.snapshotQueue != nil { - f.snapshotQueue = i.snapshotQueue - } f.OpenTranslateStore = i.OpenTranslateStore return f, nil } diff --git a/index_internal_test.go b/index_internal_test.go index 026607da4..5b5a5b5a3 100644 --- a/index_internal_test.go +++ b/index_internal_test.go @@ -25,7 +25,7 @@ func mustOpenIndex(opt IndexOptions) *Index { if err != nil { panic(err) } - index, err := NewIndex(path, "i", DefaultPartitionN) + index, err := NewIndex(NewHolder(1), path, "i") if err != nil { panic(err) } diff --git a/index_test.go b/index_test.go index 63e9b6179..bfebee7f7 100644 --- a/index_test.go +++ b/index_test.go @@ -242,7 +242,7 @@ func TestIndex_InvalidName(t *testing.T) { if err != nil { panic(err) } - index, err := pilosa.NewIndex(path, "ABC", pilosa.DefaultPartitionN) + index, err := pilosa.NewIndex(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "ABC") if err == nil { t.Fatalf("should have gotten an error on index name with caps") } diff --git a/roaring/roaring.go b/roaring/roaring.go index 3eb5ecef2..f79b20c5b 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -2513,8 +2513,8 @@ type BitmapInfo struct { OpDetails []OpInfo BitCount uint64 ContainerCount int - Containers []ContainerInfo - OpContainers []ContainerInfo + Containers []ContainerInfo // The containers found in the bitmap originally + OpContainers []ContainerInfo // The containers resulting from ops log changes. } // Iterator represents an iterator over a Bitmap. diff --git a/roaring/unmarshal_binary.go b/roaring/unmarshal_binary.go index 59f65a117..ee566236a 100644 --- a/roaring/unmarshal_binary.go +++ b/roaring/unmarshal_binary.go @@ -96,8 +96,9 @@ func (b *Bitmap) UnmarshalBinary(data []byte) (err error) { return nil } -// unmarshalPilosaRoaring treats data as being encoded in Pilosa's 64 bit -// roaring format and decodes it into b. +// InspectBinary reads a roaring bitmap, plus a possible ops log, +// and reports back on the contents, including distinguishing between +// the original ops log and the post-ops-log contents. func InspectBinary(data []byte) (info BitmapInfo, err error) { if data == nil { return info, errors.New("no roaring bitmap provided") diff --git a/server.go b/server.go index 131a6971c..75f5b2024 100644 --- a/server.go +++ b/server.go @@ -66,9 +66,10 @@ type Server struct { // nolint: maligned extensions []*ext.ExtensionInfo // External - systemInfo SystemInfo - gcNotifier GCNotifier - logger logger.Logger + systemInfo SystemInfo + gcNotifier GCNotifier + logger logger.Logger + snapshotQueue SnapshotQueue nodeID string uri URI @@ -533,6 +534,9 @@ func (s *Server) UpAndDown() error { func (s *Server) Open() error { s.logger.Printf("open server") + // Start background monitoring. + s.snapshotQueue = newSnapshotQueue(10, 2, s.logger) + // Log startup err := s.holder.logStartup() if err != nil { @@ -560,6 +564,8 @@ func (s *Server) Open() error { if err := s.holder.Open(); err != nil { return errors.Wrap(err, "opening Holder") } + // bring up the background tasks for the holder. + s.holder.Activate() if err := s.cluster.setNodeState(nodeStateReady); err != nil { return errors.Wrap(err, "setting nodeState") } @@ -571,7 +577,6 @@ func (s *Server) Open() error { // buffered channel. s.cluster.listenForJoins() - // Start background monitoring. s.wg.Add(3) go func() { defer s.wg.Done(); s.monitorAntiEntropy() }() go func() { defer s.wg.Done(); s.monitorRuntime() }() @@ -596,6 +601,11 @@ func (s *Server) Close() error { if s.holder != nil { errh = s.holder.Close() } + if s.snapshotQueue != nil { + s.holder.SnapshotQueue = nil + s.snapshotQueue.Stop() + s.snapshotQueue = nil + } // prefer to return holder error over cluster // error. This order is somewhat arbitrary. It would be better if we had // some way to combine all the errors, but probably not important enough to diff --git a/snapshotqueue.go b/snapshotqueue.go index 373e5c538..5dcc5dbf8 100644 --- a/snapshotqueue.go +++ b/snapshotqueue.go @@ -16,6 +16,7 @@ package pilosa import ( "fmt" + "math/bits" "os" "sync" "sync/atomic" @@ -39,12 +40,22 @@ import ( // Await, Enqueue, and Immediate should be called only with the fragment lock // held. // -// ScanHolder spawns a new goroutine. You don't need to use `go` on it. -type snapshotQueue interface { +// If you create a queue, it should get stopped at some point. The +// atomicSnapshotQueue implementation used as defaultSnapshotQueue has +// a Start function which will tell you whether it actually started a +// queue. This logic exists because in a normal server case, you probably +// want the queue to be shut down as part of server shutdown, but if you're +// running cluster tests, you probably want to start and shop the queue as +// part of the test, not stop it when any server terminates. +// +// It's less likely to be desireable to start/stop individual queues, +// because fragments use the defaultSnapshotQueue anyway. This design +// needs revisiting. +type SnapshotQueue interface { Immediate(*fragment) error Enqueue(*fragment) Await(*fragment) error - ScanHolder(*Holder) + ScanHolder(*Holder, chan struct{}) Stop() } @@ -64,20 +75,18 @@ func (q *queuelessSnapshotQueue) Immediate(f *fragment) error { return f.snapshot() } -func (q *queuelessSnapshotQueue) ScanHolder(h *Holder) { +func (q *queuelessSnapshotQueue) ScanHolder(h *Holder, done chan struct{}) { } func (q *queuelessSnapshotQueue) Stop() { } -// defaultSnapshotQueue is the fallback to use if none is available, -// and currently uses queueless -- it runs all snapshots immediately. -var defaultSnapshotQueue *queuelessSnapshotQueue +var defaultSnapshotQueue = &queuelessSnapshotQueue{} // newSnapshotQueue makes a new snapshot queue, of depth N, with // w worker threads. -func newSnapshotQueue(n int, w int, l logger.Logger) snapshotQueue { - sq := prioritySnapshotQueue{normal: make(chan snapshotRequest, n), urgent: make(chan snapshotRequest), background: make(chan snapshotRequest), done: make(chan struct{}), logger: l} +func newSnapshotQueue(n int, w int, l logger.Logger) SnapshotQueue { + sq := prioritySnapshotQueue{normal: make(chan snapshotRequest, n), urgent: make(chan snapshotRequest), background: make(chan snapshotRequest), done: make(chan struct{}), maxOpN: 10000, logger: l} if sq.logger == nil { sq.logger = logger.NewStandardLogger(os.Stderr) } @@ -108,9 +117,11 @@ type prioritySnapshotQueue struct { done chan struct{} mu sync.RWMutex scanWG, workerWG sync.WaitGroup + maxOpN int + observedOpN [16]int stats struct { - enqueued uint64 - skipped uint64 + enqueued uint32 + skipped uint32 } } @@ -192,7 +203,9 @@ func (sq *prioritySnapshotQueue) Stop() { sq.urgent = nil close(sq.background) sq.background = nil - if sq.stats.skipped > 0 || sq.stats.enqueued > 1 { + enqueued := atomic.LoadUint32(&sq.stats.enqueued) + skipped := atomic.LoadUint32(&sq.stats.skipped) + if skipped > 0 || enqueued > 1 { sq.logger.Printf("snapshot queue: enqueued %d, skipped %d\n", sq.stats.enqueued, sq.stats.skipped) } } @@ -217,10 +230,10 @@ func (sq *prioritySnapshotQueue) Enqueue(f *fragment) { // try to enqueue snapshot select { case sq.normal <- snapshotRequest{frag: f, when: time.Now()}: - atomic.AddUint64(&sq.stats.enqueued, 1) + atomic.AddUint32(&sq.stats.enqueued, 1) return default: - atomic.AddUint64(&sq.stats.skipped, 1) + atomic.AddUint32(&sq.stats.skipped, 1) f.snapshotPending = false return } @@ -230,6 +243,11 @@ func (sq *prioritySnapshotQueue) Enqueue(f *fragment) { // held. Await waits on a condition variable inside f, associated with the // fragment's lock, so this does not conflict with the lock being used for // snapshots. +// +// Note that workers don't stop just because the queue's been stopped; only +// the background scanner is stopped. So an Await shouldn't block forever +// even if the queue gets shut down. If you're reading this, possibly that +// analysis is incorrect. func (sq *prioritySnapshotQueue) Await(f *fragment) (err error) { for f.snapshotPending { f.snapshotCond.Wait() @@ -281,9 +299,19 @@ func (sq *prioritySnapshotQueue) needsSnapshot(f *fragment) bool { if f.snapshotPending { return false } - if f.opN > f.MaxOpN { + if f.opN > sq.maxOpN { return true } + // aka "log2(n) + 1", or 0 for n==0 + pow2 := 32 - bits.LeadingZeros32(uint32(f.opN)) + // 15 == 16384. we assume that since 16384 is higher than our + // normal maxOpN, it's always a reasonable value. + if pow2 > 15 { + pow2 = 15 + } + // store in inverse order so the lowest slot in the array is the + // highest cardinality + sq.observedOpN[15-pow2]++ return false } @@ -291,10 +319,10 @@ func (sq *prioritySnapshotQueue) needsSnapshot(f *fragment) bool { // indexes/fields/views/fragments, looking for fragments which have OpN // high enough to justify a snapshot but don't seem to have one pending. // It then dumps these in the low priority background queue. -func (sq *prioritySnapshotQueue) ScanHolder(h *Holder) { +func (sq *prioritySnapshotQueue) ScanHolder(h *Holder, done chan struct{}) { sq.mu.Lock() sq.scanWG.Add(1) - go sq.scanHolderWorker(h, sq.background, sq.done) + go sq.scanHolderWorker(h, sq.background, done) sq.mu.Unlock() } @@ -302,6 +330,9 @@ func (sq *prioritySnapshotQueue) ScanHolder(h *Holder) { // fragments which need snapshots taken. It's the cleanup task for snapshots // that would have been requested by Enqueue, but the queue was full. func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan snapshotRequest, done chan struct{}) { + // queueDone is global to this snapshotQueue, done is specific to this holder + // scanner. + queueDone := sq.done defer sq.scanWG.Done() var indexNames, fieldNames, viewNames []string var fragNums []uint64 @@ -379,9 +410,11 @@ func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan sna // background queue. If there's nothing that needs snapshots, // we pause frequently for a second or so at a time. counter++ - if counter == 100 { + if counter == 1000 { select { case <-time.After(1 * time.Second): + case <-queueDone: + return case <-done: return } @@ -400,9 +433,36 @@ func (sq *prioritySnapshotQueue) scanHolderWorker(h *Holder, background chan sna // No reason to be active if we're not finding anything. select { case <-time.After(60 * time.Second): + case <-queueDone: + return case <-done: return } } + sq.logger.Debugf("observedOpN by power of 2: %d\n", sq.observedOpN[:]) + total := 0 + for _, v := range sq.observedOpN { + total += v + } + target := total / 4 + subTotal := 0 + for i, v := range sq.observedOpN { + subTotal += v + if subTotal >= target { + prevMaxOpN := sq.maxOpN + sq.maxOpN = (1 << (15 - uint(i))) / 2 + if sq.maxOpN > 0 { + sq.maxOpN-- + } + if prevMaxOpN != sq.maxOpN { + sq.logger.Printf("background scan: %d/%d fragments considered have opN %d or higher\n", + subTotal, total, sq.maxOpN) + } + break + } + } + for i := range sq.observedOpN { + sq.observedOpN[i] = 0 + } } } diff --git a/test/field.go b/test/field.go index b95e4f88b..f56672834 100644 --- a/test/field.go +++ b/test/field.go @@ -33,7 +33,7 @@ func newField(opts pilosa.FieldOption) *Field { if err != nil { panic(err) } - field, err := pilosa.NewField(path, "i", "f", opts) + field, err := pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i", "f", opts) if err != nil { panic(err) } @@ -63,7 +63,7 @@ func (f *Field) reopen() error { } path, index, name := f.Path(), f.Index(), f.Name() - f.Field, err = pilosa.NewField(path, index, name, pilosa.OptFieldTypeDefault()) + f.Field, err = pilosa.NewField(pilosa.NewHolder(pilosa.DefaultPartitionN), path, index, name, pilosa.OptFieldTypeDefault()) if err != nil { return err } diff --git a/test/index.go b/test/index.go index bf65c6ae2..b90702ae2 100644 --- a/test/index.go +++ b/test/index.go @@ -32,7 +32,7 @@ func newIndex() *Index { if err != nil { panic(err) } - index, err := pilosa.NewIndex(path, "i", pilosa.DefaultPartitionN) + index, err := pilosa.NewIndex(pilosa.NewHolder(pilosa.DefaultPartitionN), path, "i") if err != nil { panic(err) } @@ -62,7 +62,7 @@ func (i *Index) Reopen() error { } path, name := i.Path(), i.Name() - i.Index, err = pilosa.NewIndex(path, name, pilosa.DefaultPartitionN) + i.Index, err = pilosa.NewIndex(pilosa.NewHolder(pilosa.DefaultPartitionN), path, name) if err != nil { return err } diff --git a/view.go b/view.go index 478725e1e..af73b7fe2 100644 --- a/view.go +++ b/view.go @@ -26,7 +26,6 @@ import ( "sync/atomic" "time" - "github.com/pilosa/pilosa/v2/logger" "github.com/pilosa/pilosa/v2/pql" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/stats" @@ -49,6 +48,8 @@ type view struct { field string name string + holder *Holder + fieldType string cacheType string cacheSize uint32 @@ -56,24 +57,24 @@ type view struct { // Fragments by shard. fragments map[uint64]*fragment - broadcaster broadcaster - stats stats.StatsClient - rowAttrStore AttrStore - logger logger.Logger - snapshotQueue snapshotQueue + broadcaster broadcaster + stats stats.StatsClient + rowAttrStore AttrStore knownShards *roaring.Bitmap knownShardsCopied uint32 } // newView returns a new instance of View. -func newView(path, index, field, name string, fieldOptions FieldOptions) *view { +func newView(holder *Holder, path, index, field, name string, fieldOptions FieldOptions) *view { return &view{ path: path, index: index, field: field, name: name, + holder: holder, + fieldType: fieldOptions.Type, cacheType: fieldOptions.CacheType, cacheSize: fieldOptions.CacheSize, @@ -82,7 +83,6 @@ func newView(path, index, field, name string, fieldOptions FieldOptions) *view { broadcaster: NopBroadcaster, stats: stats.NopStatsClient, - logger: logger.NopLogger, knownShards: roaring.NewSliceBitmap(), } } @@ -133,14 +133,14 @@ func (v *view) open() error { if err := func() error { // Ensure the view's path exists. - v.logger.Debugf("ensure view path exists: %s", v.path) + v.holder.Logger.Debugf("ensure view path exists: %s", v.path) if err := os.MkdirAll(v.path, 0777); err != nil { return errors.Wrap(err, "creating view directory") } else if err := os.MkdirAll(filepath.Join(v.path, "fragments"), 0777); err != nil { return errors.Wrap(err, "creating fragments directory") } - v.logger.Debugf("open fragments for index/field/view: %s/%s/%s", v.index, v.field, v.name) + v.holder.Logger.Debugf("open fragments for index/field/view: %s/%s/%s", v.index, v.field, v.name) if err := v.openFragments(); err != nil { return errors.Wrap(err, "opening fragments") } @@ -151,7 +151,7 @@ func (v *view) open() error { return err } - v.logger.Debugf("successfully opened index/field/view: %s/%s/%s", v.index, v.field, v.name) + v.holder.Logger.Debugf("successfully opened index/field/view: %s/%s/%s", v.index, v.field, v.name) return nil } @@ -190,12 +190,12 @@ fileLoop: // Parse filename into integer. shard, err := strconv.ParseUint(filepath.Base(fi.Name()), 10, 64) if err != nil { - v.logger.Debugf("WARNING: couldn't use non-integer file as shard in index/field/view %s/%s/%s: %s", v.index, v.field, v.name, fi.Name()) + v.holder.Logger.Debugf("WARNING: couldn't use non-integer file as shard in index/field/view %s/%s/%s: %s", v.index, v.field, v.name, fi.Name()) continue } workQueue <- struct{}{} - v.logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard) + v.holder.Logger.Debugf("open index/field/view/fragment: %s/%s/%s/%d", v.index, v.field, v.name, shard) eg.Go(func() error { defer func() { <-workQueue @@ -205,7 +205,7 @@ fileLoop: return fmt.Errorf("open fragment: shard=%d, err=%s", frag.shard, err) } frag.RowAttrStore = v.rowAttrStore - v.logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard) + v.holder.Logger.Debugf("add index/field/view/fragment to view.fragments: %s/%s/%s/%d", v.index, v.field, v.name, shard) mu.Lock() v.fragments[frag.shard] = frag v.addKnownShard(frag.shard) @@ -342,7 +342,7 @@ func (v *view) notifyIfNewShard(shard uint64) { // Broadcast a message that a new max shard was just created. err := v.broadcaster.SendSync(msg) if err != nil { - v.logger.Printf("broadcasting create shard: %v", err) + v.holder.Logger.Printf("broadcasting create shard: %v", err) } close(broadcastChan) }() @@ -352,19 +352,15 @@ func (v *view) notifyIfNewShard(shard uint64) { select { case <-broadcastChan: case <-time.After(50 * time.Millisecond): - v.logger.Debugf("broadcasting create shard took >50ms") + v.holder.Logger.Debugf("broadcasting create shard took >50ms") } } func (v *view) newFragment(path string, shard uint64) *fragment { - frag := newFragment(path, v.index, v.field, v.name, shard, v.flags()) + frag := newFragment(v.holder, path, v.index, v.field, v.name, shard, v.flags()) frag.CacheType = v.cacheType frag.CacheSize = v.cacheSize - frag.Logger = v.logger frag.stats = v.stats - if v.snapshotQueue != nil { - frag.snapshotQueue = v.snapshotQueue - } if v.fieldType == FieldTypeMutex { frag.mutexVector = newRowsVector(frag) } else if v.fieldType == FieldTypeBool { @@ -382,7 +378,7 @@ func (v *view) deleteFragment(shard uint64) error { return ErrFragmentNotFound } - v.logger.Printf("delete fragment: (%s/%s/%s) %d", v.index, v.field, v.name, shard) + v.holder.Logger.Printf("delete fragment: (%s/%s/%s) %d", v.index, v.field, v.name, shard) // Close data files before deletion. if err := fragment.Close(); err != nil { @@ -396,7 +392,7 @@ func (v *view) deleteFragment(shard uint64) error { // Delete fragment cache file. if err := os.Remove(fragment.cachePath()); err != nil { - v.logger.Printf("no cache file to delete for shard %d", shard) + v.holder.Logger.Printf("no cache file to delete for shard %d", shard) } delete(v.fragments, shard) diff --git a/view_internal_test.go b/view_internal_test.go index bb50003d2..bbb1a20f4 100644 --- a/view_internal_test.go +++ b/view_internal_test.go @@ -34,7 +34,7 @@ func mustOpenView(index, field, name string) *view { CacheSize: DefaultCacheSize, } - v := newView(path, index, field, name, fo) + v := newView(NewHolder(DefaultPartitionN), path, index, field, name, fo) if err := v.open(); err != nil { panic(err) }