From 82c75851dfa19da218a3d1b51e18652dce2a1e78 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Fri, 21 Jan 2022 17:00:45 -0600 Subject: [PATCH 1/3] rip out generation stuff it was somewhat difficult to avoid ripping this out without also touching some of the stuff that supports roaring backend. That's going soon too, so no worries :) --- fragment.go | 242 ++---------------- fragment_internal_test.go | 6 +- gendebug_test.go | 36 --- generation.go | 450 ---------------------------------- generation_debug.go | 152 ------------ generation_nodebug.go | 25 -- generation_test.go | 59 ----- holder.go | 2 +- mmap_test.go | 28 --- roaring/filter.go | 2 +- roaring/generation_debug.go | 7 - roaring/generation_nodebug.go | 7 - roaring/roaring.go | 3 +- 13 files changed, 20 insertions(+), 999 deletions(-) delete mode 100644 gendebug_test.go delete mode 100644 generation.go delete mode 100644 generation_debug.go delete mode 100644 generation_nodebug.go delete mode 100644 generation_test.go delete mode 100644 roaring/generation_debug.go delete mode 100644 roaring/generation_nodebug.go diff --git a/fragment.go b/fragment.go index 1eceb076f..41b1f6a79 100644 --- a/fragment.go +++ b/fragment.go @@ -140,7 +140,6 @@ type fragment struct { // File-backed storage flags byte // user-defined flags passed to roaring - gen generation storage *roaring.Bitmap opN int // number of ops since snapshot (may be approximate for imports) ops int // number of higher-level operations, as opposed to bit changes @@ -319,159 +318,15 @@ func (f *fragment) emptyStorage(file *os.File) (bool, error) { return false, nil } -// importStorage attempts to import data from storage -- for instance, -// reading in a roaring bitmap from media. -func (f *fragment) importStorage(data []byte, file *os.File, newGen generation, mapped bool) (bool, error) { - f.storage.PreferMapping(mapped) - if len(data) == 0 { - return f.emptyStorage(file) - } - - // UnmarshalBinary will have remapped the storage to newGen if it - // succeeded, or if it fails but the error is advisory-only. So we - // optimistically set the source here, but if there's a non-advisory - // error, we'll unmap it and then set the source to nil. - f.storage.SetSource(newGen) - if err := f.storage.UnmarshalBinary(data); err != nil { - // roaring can report advisory-only errors... - cause := errors.Cause(err) - _, ok := cause.(roaring.AdvisoryError) - if !ok { - _, e2 := f.storage.RemapRoaringStorage(nil) - f.storage.SetSource(nil) - if e2 != nil { - return false, fmt.Errorf("unmarshal storage: file=%s, err=%s, clearing old mapping also failed: %v", file.Name(), err, e2) - } - return false, fmt.Errorf("unmarshal storage: file=%s, err=%s", file.Name(), err) - } - f.holder.Logger.Warnf("unmarshal storage, file=%s, err=%v", file.Name(), err) - trunc, ok := cause.(roaring.FileShouldBeTruncatedError) - if ok && !f.holder.Opts.ReadOnly { - // if the holder is ReadOnly, we silently ignore the "advisory" - // error. This may be a bad idea. - - // generation code looks for a FileShouldBeTruncatedError - return false, trunc - } - } - f.ops, f.opN = f.storage.Ops() - // For now, we assume that UnmarshalBinary will have mapped at least - // one container if we told it the storage was mapped and it didn't - // error out. This might be wrong in occasional trivial cases, but - // it should be harmless. - return mapped, nil -} - -// applyStorage applies storage to a fragment that may already have -// usable data. For instance, this would try to remap existing containers -// to use a new storage as backing store. -func (f *fragment) applyStorage(data []byte, file *os.File, newGen generation, mapped bool) (bool, error) { - if len(data) == 0 { - // This shouldn't be used anyway in this path, but just in - // case, we'll be explicit about it. - f.storage.PreferMapping(false) - if file != nil { - fi, err := file.Stat() - if err != nil { - f.holder.Logger.Errorf("trying to apply new storage to existing bitmap, stat failed: %v", err) - } - if err == nil && fi != nil && fi.Size() == 0 { - return f.emptyStorage(file) - } - } - // if we can't be sure of that, we assume data is 0 because - // we couldn't mmap it, and since all we'd be doing is remapping - // our containers to use that storage *to take advantage of - // mmap*, we'll just make sure our containers aren't pointing to - // old storage and say "nope". - _, _ = f.storage.RemapRoaringStorage(nil) - f.storage.SetSource(nil) - return false, nil - } - // Tell storage to prefer mapping if and only if we think the data - // is mmapped and valid. - f.storage.PreferMapping(mapped) - // RemapRoaringStorage will fix any mapped containers to point either - // to the provided data (if PreferMapping was called with true and - // data is provided and there's a corresponding container) or to - // allocated storage, so when it's done, there's nothing in it that - // is mapped to anything *other than* the provided data. - mapped, err := f.storage.RemapRoaringStorage(data) - if err != nil { - // OOPS! something went wrong, we don't know why, we can't - // sanely recover from that. - _, _ = f.storage.RemapRoaringStorage(nil) - mapped = false - f.storage.SetSource(nil) - } else { - f.storage.SetSource(newGen) - } - return mapped, err -} - -func (f *fragment) inspectStorage(data []byte, file *os.File, newGen generation, mapped bool) (didMap bool, err error) { - f.bitmapInfo = &roaring.BitmapInfo{} - f.storage, didMap, err = roaring.InspectBinary(data, mapped, f.bitmapInfo) - return didMap, err -} - -// openStorage opens the storage bitmap. -// -// This has been massively reworked recently, and now hands a lot of -// file management off to the generation object and the Done method -// of that object. Similarly, the bitmap mapping/remapping -// logic is now mostly in importStorage (reading in a bitmap) and applyStorage -// (remapping an existing bitmap to match a new backing store). +// openStorage opens the storage bitmap. Does nothing in RBF-world and will be removed soon. func (f *fragment) openStorage(unmarshalData bool) error { if !f.idx.NeedsSnapshot() { - f.gen = &NopGeneration{} f.currdata = struct{ from, to uintptr }{} f.prevdata = f.currdata return nil // openStorage becomes a noop under RBF, Badger, etc. } - // Create a roaring bitmap to serve as storage for the shard. - if f.storage == nil { - f.storage = roaring.NewFileBitmap() - f.storage.Flags = f.flags - // if we didn't actually have storage, we *do* need to - // unmarshal this data in order to have any. - unmarshalData = true - } - - var storageOp func([]byte, *os.File, generation, bool) (bool, error) - if f.holder.Opts.Inspect { - // note that this will unmarshal even if we already have - // storage; when Inspect is on for a holder, we actually want - // to be able to report this. - storageOp = f.inspectStorage - } else { - if unmarshalData { - storageOp = f.importStorage - } else { - storageOp = f.applyStorage - } - } - var err error - 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 - var scratchAddrs struct{ from, to uintptr } - if scratchData != nil { - scratchAddrs.from = uintptr(unsafe.Pointer(&scratchData[0])) - scratchAddrs.to = scratchAddrs.from + uintptr(len(scratchData)) - } - f.currdata = scratchAddrs - } - if generationDebug { - // We might have already done this anyway, if we think we - // mapped stuff, but when debugging we want to do it - // unconditionally, because the test cases otherwise won't - // exercise this code well. - f.storage.SetSource(f.gen) - } - return err + return nil } // openCache initializes the cache from row ids persisted to disk. @@ -557,17 +412,13 @@ func (f *fragment) close() error { return nil } -// closeStorage marks the current generation as done. It is not necessary -// to call this before openStorage. +// closeStorage is essentially a no-op and will go away soon. func (f *fragment) closeStorage() error { // opN is determined by how many bit set/clear operations are in the storage // write log, so once the storage is closed it should be 0. Opening new // storage will set opN appropriately. f.opN = 0 - if f.gen != nil { - f.gen.Done() - } return nil } @@ -640,10 +491,7 @@ func (f *fragment) rowFromStorage(tx Tx, rowID uint64) (*Row, error) { func (f *fragment) setBit(tx Tx, rowID, columnID uint64) (changed bool, err error) { f.mu.Lock() // controls access to the file. defer f.mu.Unlock() - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } + doSetFunc := func() error { // handle mutux field type if f.mutexVector != nil { @@ -654,16 +502,7 @@ func (f *fragment) setBit(tx Tx, rowID, columnID uint64) (changed bool, err erro changed, err = f.unprotectedSetBit(tx, rowID, columnID) return err } - // avoid crashing when f.gen is nil - if f.gen != nil { - err = f.gen.Transaction(wp, doSetFunc) - } else { - if tx.Type() == RoaringTxn { - return changed, errors.New("internal error: f.gen was nil and tx.Type is RoaringTxn - should never happen under roaring b/c storage should be open") - } - // else transactional backend. Just do it. - err = doSetFunc() - } + err = doSetFunc() return changed, err } @@ -728,15 +567,7 @@ func (f *fragment) unprotectedSetBit(tx Tx, rowID, columnID uint64) (changed boo func (f *fragment) clearBit(tx Tx, rowID, columnID uint64) (changed bool, err error) { f.mu.Lock() defer f.mu.Unlock() - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err = f.gen.Transaction(wp, func() error { - changed, err = f.unprotectedClearBit(tx, rowID, columnID) - return err - }) - return changed, err + return f.unprotectedClearBit(tx, rowID, columnID) } // unprotectedClearBit TODO should be replaced by an invocation of @@ -788,15 +619,7 @@ func (f *fragment) unprotectedClearBit(tx Tx, rowID, columnID uint64) (changed b func (f *fragment) setRow(tx Tx, row *Row, rowID uint64) (changed bool, err error) { f.mu.Lock() defer f.mu.Unlock() - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err = f.gen.Transaction(wp, func() error { - changed, err = f.unprotectedSetRow(tx, row, rowID) - return err - }) - return changed, err + return f.unprotectedSetRow(tx, row, rowID) } func (f *fragment) unprotectedSetRow(tx Tx, row *Row, rowID uint64) (changed bool, err error) { @@ -854,15 +677,7 @@ func (f *fragment) unprotectedSetRow(tx Tx, row *Row, rowID uint64) (changed boo func (f *fragment) clearRow(tx Tx, rowID uint64) (changed bool, err error) { f.mu.Lock() defer f.mu.Unlock() - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err = f.gen.Transaction(wp, func() error { - changed, err = f.unprotectedClearRow(tx, rowID) - return err - }) - return changed, err + return f.unprotectedClearRow(tx, rowID) } func (f *fragment) unprotectedClearRow(tx Tx, rowID uint64) (changed bool, err error) { @@ -903,11 +718,7 @@ func (f *fragment) clearBlock(tx Tx, block int) (changed bool, err error) { defer f.mu.Unlock() firstRow := uint64(block * HashBlockSize) - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err = f.gen.Transaction(wp, func() error { + err = func() error { var rowChanged bool for rowID := uint64(firstRow); rowID < firstRow+HashBlockSize; rowID++ { if chang, err := f.unprotectedClearRow(tx, rowID); err != nil { @@ -918,7 +729,7 @@ func (f *fragment) clearBlock(tx Tx, block int) (changed bool, err error) { } changed = rowChanged return nil - }) + }() return changed, err } @@ -1028,11 +839,7 @@ func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint64, val }() } - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err = f.gen.Transaction(wp, func() error { + err = func() error { // Convert value to an unsigned representation. uvalue := uint64(value) if value < 0 { @@ -1086,7 +893,7 @@ func (f *fragment) setValueBase(txOrig Tx, columnID uint64, bitDepth uint64, val } return nil - }) + }() return changed, err } @@ -2369,11 +2176,6 @@ func (p parallelSlices) Swap(i, j int) { // operations to the op log. func (f *fragment) importPositions(tx Tx, set, clear []uint64, rowSet map[uint64]struct{}) error { //tx.AddN() - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - doFunc := func() error { if len(set) > 0 { f.stats.Count(MetricImportingN, int64(len(set)), 1) @@ -2420,16 +2222,7 @@ func (f *fragment) importPositions(tx Tx, set, clear []uint64, rowSet map[uint64 } return nil } - var err error - if f.gen != nil { - err = f.gen.Transaction(wp, doFunc) - } else { - if tx.Type() == RoaringTxn { - return errors.New("internal error: f.gen was nil and tx.Type is RoaringTxn - should never happen under roaring b/c storage should be open") - } - err = doFunc() - } - + err := doFunc() if err != nil && f.storage != nil { // 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) @@ -2685,11 +2478,7 @@ func (f *fragment) doImportRoaring(ctx context.Context, tx Tx, data []byte, clea defer span.Finish() var rowSet map[uint64]int - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err := f.gen.Transaction(wp, func() (err error) { + err := func() (err error) { var rit roaring.RoaringIterator rit, err = roaring.NewRoaringIterator(data) if err != nil { @@ -2698,8 +2487,7 @@ func (f *fragment) doImportRoaring(ctx context.Context, tx Tx, data []byte, clea _, rowSet, err = tx.ImportRoaringBits(f.index(), f.field(), f.view(), f.shard, rit, clear, true, rowSize) return err - }) - + }() if err != nil { return nil, false, err } diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 0ecdfb787..e079c05ef 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -3010,7 +3010,7 @@ func BenchmarkImportRoaringUpdate(b *testing.B) { b.StopTimer() var stat os.FileInfo var statTarget io.Writer - err = f.gen.Transaction(&statTarget, func() error { + err = func() error { targetFile, ok := statTarget.(*os.File) if ok { stat, _ = targetFile.Stat() @@ -3018,7 +3018,7 @@ func BenchmarkImportRoaringUpdate(b *testing.B) { b.Errorf("couldn't stat file") } return nil - }) + }() if err != nil { b.Errorf("transaction error: %v", err) } @@ -3486,8 +3486,6 @@ func (f *fragment) Clean(t testing.TB) { } }() errc := f.Close() - // prevent double-closes of generation during testing. - f.gen = nil if errc != nil { t.Fatalf("error closing fragment: %v", errc) } diff --git a/gendebug_test.go b/gendebug_test.go deleted file mode 100644 index 05094bd22..000000000 --- a/gendebug_test.go +++ /dev/null @@ -1,36 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -// -//go:build generationdebug -// +build generationdebug - -package pilosa - -import ( - "errors" - "fmt" - "runtime" - - "github.com/molecula/featurebase/v3/testhook" -) - -func examineResults() error { - runtime.GC() - stats, results := reportGenerations() - if len(stats) > 0 { - fmt.Printf("generation stats: %s\n", stats) - } - if len(results) == 0 { - return nil - } - if len(results) > 0 { - fmt.Printf("generations:\n") - for _, res := range results { - fmt.Printf(" %s\n", res) - } - } - return errors.New("outstanding generations detected") -} - -func init() { - testhook.RegisterPostTestHook(examineResults) -} diff --git a/generation.go b/generation.go deleted file mode 100644 index dfb8a3277..000000000 --- a/generation.go +++ /dev/null @@ -1,450 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -package pilosa - -import ( - "fmt" - "io" - "io/ioutil" - "os" - "runtime" - // "runtime/debug" - "sync" - "syscall" - "time" - - "github.com/molecula/featurebase/v3/logger" - "github.com/molecula/featurebase/v3/roaring" - "github.com/molecula/featurebase/v3/syswrap" - "github.com/pkg/errors" -) - -// generation represents one "generation" of opening a data file. -// This is what determines when it's safe to unmap a data file, if it -// got mapped, and handles closing/reopening files if we need to -// manage file handle availability. It's an interface because this -// lets us write simpler code for specific cases, rather than handling -// the whole matrix of mapped/unmapped, staying open/being reopened, -// etcetera. -// -// You create a generation by calling newGeneration with a file -// path. If it succeeds in opening that path, it calls a provided -// setup function with the data from the generation, and a flag -// indicating whether the data is mmapped. If the setup function -// fails, newGeneration cleans things up and closes. Otherwise, -// it returns a generation. -// -// The generation itself uses runtime.SetFinalizer to clean up when -// the last reference to it goes away. You should store a pointer -// to the generation in any object which is reliant on the generation. -// -// When you anticipate a generation should be done (for instance, -// opening a new generation), the old one gets marked done, which -// stashes a timestamp in it. Later operations can check whether -// the timestamp is a while back, and if so, complain that something -// might be wrong. -// -// In some cases, we don't have enough open file limit to keep every -// file actually open. To address this, use the `Transaction` function, -// which ensures that the file is open, stores a reference to it in -// a provided `*io.Writer`, and then restores the previous value of -// the io.Writer when it's done. For instance, for a bitmap, this might -// be used with `&b.OpWriter`. -// -// newGeneration takes an optional previous generation; it calls -// that generation's Done function after running the provided setup, -// and bumps the generation count. -type generation interface { - // Transaction runs the given transaction with the generation's - // file open. If the **os.File parameter is - // non-nil, the generation's file will be open, and stored - // into that pointer, during the execution of func, after - // which the previous contents are restored. Otherwise - // the file may or may not be open during the operation. - Transaction(*io.Writer, func() error) error - // Done() should be called exactly once, to indicate that a - // generation is expected not to be in use for long -- for instance, - // when a new generation replaces it. - Done() - // Generation count. - Generation() int64 - // ID indicates the source -- path and generation number -- that - // this generation represents. - ID() string - // Dead indicates whether this generation is Done. - Dead() bool - // Bytes reports the storage associated with this generation, if any. - // DO NOT USE THIS. Except if you're debugging mmap segfaults. - Bytes() []byte -} - -type mmapGeneration struct { - mu sync.Mutex // mutex guards modifiers of generation, not of data - transMu sync.Mutex // guards transactions, specifically - path string - id string - file *os.File - data []byte - generation int64 // generation counter - dead bool // we think this generation is dead - deadSince time.Time // when this generation was marked dead - retries int // for cases where we're retrying - logger logger.Logger -} - -func (m *mmapGeneration) Dead() bool { - m.mu.Lock() - defer m.mu.Unlock() - return m.dead -} - -func (m *mmapGeneration) ID() string { - return m.id -} - -func (m *mmapGeneration) Generation() int64 { - return m.generation -} - -// Transaction runs an exclusive call, ensuring that the file is open if -// the *io.Writer parameter is present. -func (m *mmapGeneration) Transaction(fileP *io.Writer, fn func() error) (transactionErr error) { - m.transMu.Lock() - defer m.transMu.Unlock() - // HEY LOOK CAREFULLY AT THIS BIT: - // We can't just defer this unlock. We specifically want to be - // sure to unlock the regular mutex *before* this function is over, - // and if we error out trying to open the file, we want to do it - // even sooner. If we deferred this, the transaction would block - // *everything*, including things like sanity checks against the - // generation being Dead(), but also including the deferred - // re-close-the-file. - m.mu.Lock() - // if we've been asked for a file pointer, we need to ensure that - // our file is open, and that the file pointer to it is stored in - // the requested location, then revert that when we're done. - // if we aren't asked for a file pointer, nothing needs the file - // open. - if m.dead { - elapsed := time.Since(m.deadSince) - m.logger.Warnf("transaction against %s, which has been dead for %v\n", m.id, elapsed) - } - if fileP != nil { - if m.file == nil { - // we ignore the shouldClose response here; if this - // fragment was previously not being kept open, we're - // going to stick with that. - _, err := m.openFile() - if err != nil { - m.mu.Unlock() - return err - } - defer func() { - // report a close error if we have no other error to report - m.mu.Lock() - defer m.mu.Unlock() - err := m.closeFile() - if transactionErr == nil { - transactionErr = err - } - }() - } - var fileStash io.Writer - fileStash, *fileP = *fileP, m.file - defer func() { - *fileP = fileStash - }() - } - // We are done locking the generation itself for now. - m.mu.Unlock() - // wouldPanic := debug.SetPanicOnFault(true) - // defer func() { - // debug.SetPanicOnFault(wouldPanic) - // if r := recover(); r != nil { - // if err, ok := r.(error); ok { - // // special case: if we caught a page fault, we diagnose that directly. sadly, - // // we can't see the actual values that were used to generate this, probably. - // if err.Error() == "runtime error: invalid memory address or nil pointer dereference" { - // if transactionErr == nil { - // transactionErr = errors.New("invalid memory access during transaction") - // } else { - // transactionErr = fmt.Errorf("invalid memory access during transaction, previous error %v", transactionErr) - // } - // return - // } - // } - // if transactionErr == nil { - // transactionErr = fmt.Errorf("panic during transaction: %v", r) - // } else { - // transactionErr = fmt.Errorf("panic during erroring transaction: panic %v, previous error %v", r, transactionErr) - // } - // } - // }() - return fn() -} - -func (m *mmapGeneration) Bytes() []byte { - return m.data -} - -// Done marks the generation done, and closes its file, but may not unmap it. -// It's still conceptually possible to end up doing a Transaction against a -// done generation, but it's a red flag. -func (m *mmapGeneration) Done() { - if m == nil { - return - } - m.mu.Lock() - defer m.mu.Unlock() - if m.dead { - oops := fmt.Sprintf("generation %s, marked done again at %v, previously marked dead at %v", - m.id, time.Now(), m.deadSince) - panic(oops) - } - m.dead = true - m.deadSince = time.Now() - err := m.closeFile() - if err != nil { - m.logger.Errorf("error closing generation %s: %v", m.id, err) - } - // If we're not debugging, the finalizer won't have been enabled - // previously. Finalizers have non-zero cost, so having them not be - // created until they're needed seems rewarding? - if !generationDebug { - runtime.SetFinalizer(m, generationFinalizer) - } - endGeneration(m.id) - // note, Done() doesn't close the file; only the finalizer actually - // does the shutdown. -} - -// Try to close the file if it's currently open. -func (m *mmapGeneration) closeFile() error { - var lastErr error - // report the most serious error encountered, but still close - // file even if something else failed. - if m.file != nil { - if err := m.file.Sync(); err != nil { - lastErr = fmt.Errorf("sync: %s", err) - } - if err := syscall.Flock(int(m.file.Fd()), syscall.LOCK_UN); err != nil { - lastErr = fmt.Errorf("unlock: %s", err) - } - if err := syswrap.CloseFile(m.file); err != nil { - lastErr = fmt.Errorf("close file: %s", err) - } - m.file = nil - } - return lastErr -} - -// openFile ensures the file is open and locked, or fails. If it does -// open the file, it will also report the "you need to close this file -// when you're done" flag from syswrap. -func (m *mmapGeneration) openFile() (shouldClose bool, err error) { - if m.file != nil { - return false, nil - } - m.file, shouldClose, err = syswrap.OpenFile(m.path, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) - if err != nil { - return false, err - } - - // do we actually want this in every openFile? I don't know. - if err := syscall.Flock(int(m.file.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil { - _ = syswrap.CloseFile(m.file) - m.file = nil - return false, fmt.Errorf("flock: %s", err) - } - return shouldClose, nil -} - -func generationFinalizer(m *mmapGeneration) { - m.mu.Lock() - if !m.dead { - m.logger.Infof("finalizing generation %s which isn't dead yet\n", - m.id) - } - m.mu.Unlock() - err := m.closeFile() - if err != nil { - m.logger.Errorf("finalizing generation, closing file: %v\n", err) - } - if m.data != nil { - err := syswrap.Munmap(m.data) - if err != nil { - m.logger.Errorf("finalizing generation, munmap: %v\n", err) - } - m.data = nil - } - finalizeGeneration(m.id) -} - -// Cancel closes a generation out entirely. It cancels any finalizer, -// unmaps any data, ends generation tracking, and closes any files. -// It does each of these separately whether or not the others need to be done, -// or succeed. It's used to handle failures from newGeneration; it makes sure -// the generation isn't holding any resources and doesn't need to be cleaned -// up otherwise. -// -// Mostly a helper function because there's several cases where newGeneration -// might fail. -func (m *mmapGeneration) Cancel() { - if m.data != nil { - _ = syswrap.Munmap(m.data) - m.data = nil - } - err := m.closeFile() - if err != nil { - m.logger.Errorf("error cancelling generation %s: %v", m.id, err) - } - runtime.SetFinalizer(m, nil) - m.dead = true - m.deadSince = time.Now() - cancelGeneration(m.id) -} - -// newGeneration creates a new generation using the given file path. It -// then calls the provided setup function with the allocated storage, a -// file handle, the new generation, and a flag indicatting whether the storage -// is memory-mapped. If the setup function returns a non-nil error, the -// generation is cleaned up, and newGeneration fails. The setup function -// also returns a boolean indicating whether it used the mapping; if it -// didn't, newGeneration discards the mapping and returns a nil generation. -// -// If generationDebug is enabled, we track the generation even if no mapping -// is actually in use, so we can verify that the tracking is working. -// -// On failure, newGeneration returns nil values for generation and func, -// and an error. On success, the func returned is the close func to use -// when the generation is no longer needed by the caller. -func newGeneration(existing generation, path string, readData bool, setup func([]byte, *os.File, generation, bool) (bool, error), logger logger.Logger) (generation, error) { - m := mmapGeneration{path: path, logger: logger} - if existing != nil { - m.generation = existing.Generation() + 1 - m.retries = existing.(*mmapGeneration).retries - // we might keep a previous generation around just for its generation count. - if !existing.Dead() { - defer existing.Done() - } - } - shouldClose, err := m.openFile() - if err != nil { - return nil, err - } - m.id = fmt.Sprintf("%s:%d", m.path, m.generation) - // possibly assign new generation ID if this one's been used, which can - // happen with reopens, especially during testing. - m.id = registerGeneration(m.id) - // if debugging, we always want the finalizer on so we notice if a - // generation is finalized without being closed. for non-debugging - // use, we only need it when the generation is closed. - if generationDebug { - runtime.SetFinalizer(&m, generationFinalizer) - } - // Mmap the underlying file so it can be zero copied. - var mapped bool - var data []byte - fi, err := m.file.Stat() - if err == nil && fi.Size() > 0 { - data, err = syswrap.Mmap(int(m.file.Fd()), 0, int(fi.Size()), syscall.PROT_READ, syscall.MAP_SHARED) - if err == syswrap.ErrMaxMapCountReached { - // I have no idea where/how to display this message. - m.logger.Warnf("maximum number of maps reached, reading file '%s' instead", m.path) - } else if err != nil { - m.Cancel() - return nil, errors.Wrap(err, "mmap failed") - } else { - mapped = true - } - } - if data == nil && readData { - data, err = ioutil.ReadAll(m.file) - if err != nil { - m.Cancel() - return nil, errors.Wrap(err, "failure file readall") - } - } - // if we got here, data's the expected data, so let's try to use it - mappedAny, err := setup(data, m.file, &m, mapped) - - // if the setup failed, we unmap data if we previously mapped it, - // and exit. Note that having no data, or having only trivial - // data (like a zero-container Roaring file) isn't "failed". - if err != nil { - m.Cancel() - // Unless, that is, we think the file probably ought to - // be truncated: For instance, if a bitmap has a corrupted - // ops log, we could truncate that part of it and retry. - if err, ok := err.(roaring.FileShouldBeTruncatedError); ok && m.retries < 1 { - m.logger.Infof("file %s read partially, but should-be-truncated at %d bytes\n", m.path, err.SuggestedLength()) - // close this generation, then try again. once. - m.retries++ - err := os.Truncate(m.path, err.SuggestedLength()) - if err != nil { - m.logger.Errorf("truncating file failed [but retrying anyway]: %v\n", err) - } - return newGeneration(&m, path, readData, setup, logger) - } - return nil, err - } - - if mapped { - // when generationDebug is on, we want to track this even - // if it's not being used. - if generationDebug || mappedAny { - // Advise the kernel that the mmap is accessed randomly. - // We don't care much about errors with this. - _ = madvise(data, syscall.MADV_RANDOM) - // store the data, so we can unmap it when this generation - // gets finalized. - m.data = data - } else { - // unmap the data and don't stash the pointer in this - // generation. It's not being used. This generation - // doesn't need to exist, yay. - unmapErr := syswrap.Munmap(data) - if unmapErr != nil { - m.logger.Errorf("error unmapping (probably harmless): %v", unmapErr) - } - } - } - // shouldClose comes from underlying syswrap.OpenFile, which checks - // a count of open files to hint at us when we need to start closing - // files to preserve open file descriptor limit. - if shouldClose { - err := m.closeFile() - if err != nil { - m.logger.Errorf("closing file to preserve open files failed: %v\n", err) - } - } - // It's possible that the generation has no actual data to track, - // because nothing's mapped, in which case there won't be any bitmap - // sources following this, just the fragment source. (Bitmaps won't - // be attached to the source unless they're actually mapped to it, - // or generationDebug is true). That's okay. We pay a tiny cost - // for the finalizer, but we also get higher confidence that it really - // does get cleaned up. - return &m, nil -} - -// NopGeneration is used in fragment.openStorage() to short-circuit -// generation stuff that only applies to RoaringTx; doesn't apply to RBFTx/BadgerTx/etc. -type NopGeneration struct { -} - -func (g *NopGeneration) Transaction(w *io.Writer, f func() error) error { - return f() -} -func (g *NopGeneration) Done() {} -func (g *NopGeneration) Generation() int64 { - return 0 -} -func (g *NopGeneration) ID() string { - return "NOP" -} -func (g *NopGeneration) Dead() bool { - return true -} -func (g *NopGeneration) Bytes() (ret []byte) { - return -} diff --git a/generation_debug.go b/generation_debug.go deleted file mode 100644 index 03c384b28..000000000 --- a/generation_debug.go +++ /dev/null @@ -1,152 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -//go:build generationdebug -// +build generationdebug - -package pilosa - -import ( - "fmt" - "math/rand" - "runtime" - "runtime/debug" - "sort" - "sync" - "time" -) - -const generationDebug = true - -type lifespan struct { - from, to, finalized time.Time - stack []byte -} - -var knownGenerations map[string]lifespan -var knownGenerationLock sync.Mutex - -var timeZero time.Time - -var generationDebugVerbose bool - -// History reports the finalized/dead/created status of a span which we think -// is in some way in error. It's shared between a couple of places. -func (span *lifespan) History() string { - dead := "not dead" - finalized := "not finalized" - if span.finalized != timeZero { - finalized = fmt.Sprintf("finalized at %v", span.finalized) - } - if span.to != timeZero { - dead = fmt.Sprintf("dead at %v", span.to) - } - return fmt.Sprintf("%s, %s, created at %v at %s", dead, finalized, span.from, span.stack) -} - -func (span *lifespan) reportHistory(reason string, id string) string { - return fmt.Sprintf("%s %s: %s", id, reason, span.History()) -} - -func registerGeneration(id string) string { - knownGenerationLock.Lock() - defer knownGenerationLock.Unlock() - if knownGenerations == nil { - knownGenerations = make(map[string]lifespan) - } - newSpan := lifespan{from: time.Now(), stack: debug.Stack()} - origId := id - - // if you have more than 65k of the same file open, maybe you have bigger - // problems than this. - for span, exists := knownGenerations[id]; exists; span, exists = knownGenerations[id] { - suffix := fmt.Sprintf("::%04x", rand.Int63n(65536)) - if generationDebugVerbose { - history := span.History() - fmt.Printf("new generation: adding suffix %s, previous %s\n", - suffix, history) - } - id = origId + suffix - } - if generationDebugVerbose { - fmt.Printf("new generation %s\n", id) - } - knownGenerations[id] = newSpan - return id -} - -func endGeneration(id string) { - knownGenerationLock.Lock() - defer knownGenerationLock.Unlock() - span, exists := knownGenerations[id] - if !exists { - oops := fmt.Sprintf("ending generation %s: unknown", id) - panic(oops) - } - if span.finalized != timeZero || span.to != timeZero { - panic(span.reportHistory("ending generation", id)) - } - span.to = time.Now() - knownGenerations[id] = span -} - -// cancelGeneration marks the generation as finalized. In principle it's -// only used in cases where we just started a generation but something -// went wrong. it's not fancier than this because of the weird cases -// where the same generation shows up again, such as when closing and -// reopening an index so we don't know about previous instances of the -// same files. -func cancelGeneration(id string) { - knownGenerationLock.Lock() - defer knownGenerationLock.Unlock() - span, exists := knownGenerations[id] - if exists { - span.finalized = time.Now() - span.to = span.finalized - knownGenerations[id] = span - } -} - -func finalizeGeneration(id string) { - knownGenerationLock.Lock() - defer knownGenerationLock.Unlock() - span, exists := knownGenerations[id] - if !exists { - oops := fmt.Sprintf("finalizing generation %s: unknown", id) - panic(oops) - } - if span.finalized != timeZero { - panic(span.reportHistory("finalizing", id)) - } - span.finalized = time.Now() - knownGenerations[id] = span -} - -func reportGenerations() (stats string, surviving []string) { - runtime.GC() - knownGenerationLock.Lock() - defer knownGenerationLock.Unlock() - times := make([]int64, 0, len(knownGenerations)) - for id, span := range knownGenerations { - if span.to == timeZero || span.finalized == timeZero { - surviving = append(surviving, span.reportHistory("surviving", id)) - } else { - times = append(times, int64(span.finalized.Sub(span.to))) - } - } - stats = "no recorded finalized spans" - if len(times) > 0 { - sort.Slice(times, func(i, j int) bool { return times[i] < times[j] }) - var total int64 - for _, d := range times { - total += d - } - var mean, median, p90, p99, worst int64 - mean = total / int64(len(times)) - median = times[len(times)/2] - p90 = times[(len(times)*9)/10] - p99 = times[(len(times)*99)/100] - worst = times[len(times)-1] - stats = fmt.Sprintf("%d finalized spans. lag: mean %v, median %v, p90 %v, p99 %v, worst %v", - len(times), time.Duration(mean), time.Duration(median), time.Duration(p90), time.Duration(p99), time.Duration(worst)) - } - return stats, surviving -} diff --git a/generation_nodebug.go b/generation_nodebug.go deleted file mode 100644 index c49d3f219..000000000 --- a/generation_nodebug.go +++ /dev/null @@ -1,25 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -//go:build !generationdebug -// +build !generationdebug - -package pilosa - -const generationDebug = false - -func registerGeneration(id string) string { - return id -} - -func endGeneration(id string) { -} - -func cancelGeneration(id string) { -} - -func finalizeGeneration(id string) { -} - -//lint:ignore U1000 this is conditional on a build flag, see generation_test.go. -func reportGenerations() []string { //nolint:unused,deadcode - return nil -} diff --git a/generation_test.go b/generation_test.go deleted file mode 100644 index e30653301..000000000 --- a/generation_test.go +++ /dev/null @@ -1,59 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -// -//go:build generationparanoia -// +build generationparanoia - -package pilosa - -import ( - "runtime" - "testing" - "unsafe" -) - -func TestGenerationPanic(t *testing.T) { - f := mustOpenFragment("i", "f", viewStandard, 0, "none") - defer f.Clean(t) - - for i := 0; i < f.MaxOpN; i++ { - _, _ = f.setBit(0, uint64(i*32)) - } - // force snapshot so we get a mmapped row... - _ = f.Snapshot() - _ = f.row(0) - var prevData []byte - - if f.gen.(*mmapGeneration).data == nil { - t.Fatalf("generation code didn't create a mapping, apparently?") - } - prevData = f.gen.(*mmapGeneration).data - f.mu.Lock() - _ = 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() - defaultSnapshotQueue.Await(f) - f.mu.Unlock() - runtime.GC() - newData := f.gen.(*mmapGeneration).data - if unsafe.Pointer(&prevData[0]) == unsafe.Pointer(&newData[0]) { - t.Fatalf("test can't run usefully, didn't get new data pointer") - } - var wp *io.Writer - if f.storage != nil { - wp = &f.storage.OpWriter - } - err := f.gen.Transaction(wp, func() error { - prevData[0] = 0x3c - return nil - }) - if err == nil { - t.Fatalf("expected a panic to get caught, but nothing happened") - } - if err.Error() != "invalid memory access during transaction" { - t.Fatalf("expected \"invalid memory access during transaction\", got %q", err.Error()) - } -} diff --git a/holder.go b/holder.go index cf35bf998..88f7ff877 100644 --- a/holder.go +++ b/holder.go @@ -653,7 +653,7 @@ func (h *Holder) Open() error { return errors.Wrap(err, "opening index") } - // Since we don't have createAt stored on disk within the data + // Since we don't have createdAt stored on disk within the data // directory, we need to populate it from the etcd schema data. // TODO: we may no longer need the createdAt value stored in memory on // the index struct; it may only be needed in the schema return value diff --git a/mmap_test.go b/mmap_test.go index ef7513b4a..d1905eb9e 100644 --- a/mmap_test.go +++ b/mmap_test.go @@ -2,13 +2,11 @@ package pilosa import ( - "fmt" "math/rand" "runtime" "testing" "github.com/molecula/featurebase/v3/logger" - "github.com/molecula/featurebase/v3/syswrap" ) type cv struct { @@ -68,29 +66,3 @@ func forceSnapshotsCheckMapping(t *testing.T) { } } } - -// This test should basically never fail, but it might if you were running -// out of available mmaps. Which you can fake up by adding '&& false' to the test -// in newGeneration in generation.go. So this is probably useless but it's -// a failure mode we've been bitten by once... -func TestMmapBehavior(t *testing.T) { - // rbf and lmdb not happy with this test. - roaringOnlyTest(t) - - var changed bool - var original uint64 - defer func() { - syswrap.SetMaxMapCount(original) - }() - - for _, mmapMaxVal := range []uint64{0, 3} { - prev := syswrap.SetMaxMapCount(mmapMaxVal) - if !changed { - original = prev - changed = true - } - t.Run(fmt.Sprintf("maps%d", mmapMaxVal), func(t *testing.T) { - forceSnapshotsCheckMapping(t) - }) - } -} diff --git a/roaring/filter.go b/roaring/filter.go index 15a4a2f90..0f080cbc3 100644 --- a/roaring/filter.go +++ b/roaring/filter.go @@ -579,7 +579,7 @@ func (b *BitmapRowFilterMultiFilter) ConsiderData(key FilterKey, data *Container // offsets the input bitmap's containers have, it matches them against // corresponding keys. type BitmapBitmapFilter struct { - filter *Bitmap // We don't use this while iterating, but in ludicrous edge cases it might be holding a generation we need. + filter *Bitmap // We don't use this while iterating, but in ludicrous edge cases it might be holding a generation we need. TODO @seebs I don't understand why this mentions generations containers []*Container nextOffsets []uint64 callback func(uint64) error diff --git a/roaring/generation_debug.go b/roaring/generation_debug.go deleted file mode 100644 index ecbcae70a..000000000 --- a/roaring/generation_debug.go +++ /dev/null @@ -1,7 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -//go:build generationdebug -// +build generationdebug - -package roaring - -const generationDebug = true diff --git a/roaring/generation_nodebug.go b/roaring/generation_nodebug.go deleted file mode 100644 index 05fb122f7..000000000 --- a/roaring/generation_nodebug.go +++ /dev/null @@ -1,7 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -//go:build !generationdebug -// +build !generationdebug - -package roaring - -const generationDebug = false diff --git a/roaring/roaring.go b/roaring/roaring.go index fe416b50b..cdcf25b89 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -620,8 +620,7 @@ func (b *Bitmap) OffsetRange(offset, start, end uint64) *Bitmap { } other.Containers.Put(off+(k-hi0), c.Freeze()) } - // if b.Source != nil && mappedAny { - if b.Source != nil && (generationDebug || mappedAny) { + if b.Source != nil && mappedAny { other.Source = b.Source } return other From 2d44c23ac073f6012be08169fd900c7bccd4f8f7 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Fri, 21 Jan 2022 17:11:21 -0600 Subject: [PATCH 2/3] remove problematic roaring-only test --- rrtx_internal_test.go | 49 ------------------------------------------- 1 file changed, 49 deletions(-) delete mode 100644 rrtx_internal_test.go diff --git a/rrtx_internal_test.go b/rrtx_internal_test.go deleted file mode 100644 index 64c996bc4..000000000 --- a/rrtx_internal_test.go +++ /dev/null @@ -1,49 +0,0 @@ -// Copyright 2021 Molecula Corp. All rights reserved. -package pilosa - -import ( - "testing" - - . "github.com/molecula/featurebase/v3/vprint" // nolint:staticcheck -) - -func TestRoaring_HasData(t *testing.T) { - holder := newHolderWithTempPath(t, "roaring") - - idx, err := holder.CreateIndex("i", IndexOptions{}) - PanicOn(err) - defer idx.Close() - - db, err := globalRoaringReg.OpenDBWrapper(idx.path, false, nil) - PanicOn(err) - db.SetHolder(idx.holder) - - // HasData should start out false. - hasAnything, err := db.HasData() - PanicOn(err) - - if hasAnything { - t.Fatalf("HasData reported existing data on an empty database") - } - - // check that HasData sees a committed record. - - field, shard := "f", uint64(123) - - tx := idx.holder.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard}) - defer tx.Rollback() - - f, err := idx.CreateField(field) - PanicOn(err) - _, err = f.SetBit(tx, 1, 1, nil) - PanicOn(err) - PanicOn(tx.Commit()) - - hasAnything, err = db.HasData() - if err != nil { - t.Fatal(err) - } - if !hasAnything { - t.Fatalf("HasData() reported no data on a database that has 'x' written to it") - } -} From a2e109a07c4790039448a6a9dd4309bbb6d7a9c7 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Fri, 21 Jan 2022 17:27:21 -0600 Subject: [PATCH 3/3] disable roaring backend in test --- executor_test.go | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/executor_test.go b/executor_test.go index a074cf90f..15d792987 100644 --- a/executor_test.go +++ b/executor_test.go @@ -6915,11 +6915,9 @@ func TestTimelessClearRegression(t *testing.T) { } func TestMissingKeyRegression(t *testing.T) { - c := test.MustRunCluster(t, 1, []server.CommandOption{server.OptCommandServerOptions( - pilosa.OptServerStorageConfig(&storage.Config{ - Backend: "roaring", - FsyncEnabled: false, - }))}) + // this used to be explicitly roaring backend... I'm not sure + // whether it is a useful test in post-roaring world. + c := test.MustRunCluster(t, 1) defer c.Close() c.CreateField(t, "i", pilosa.IndexOptions{Keys: true, TrackExistence: true}, "f", pilosa.OptFieldKeys())