diff --git a/go/go.mod b/go/go.mod index dbc50e86b6fd..d9310b7d440c 100644 --- a/go/go.mod +++ b/go/go.mod @@ -205,7 +205,9 @@ replace ( // Fork also preserves FetchOptions.PackRefs from the v4 Keybase fork. github.com/go-git/go-git/v5 => github.com/keybase/go-git/v5 v5.19.1-keybase.3 github.com/stellar/go => github.com/keybase/stellar-org v0.0.0-20191010205648-0fc3bfe3dfa7 - github.com/syndtr/goleveldb => github.com/keybase/goleveldb v1.0.1-0.20221007195407-9881c0c26e65 + // Keybase fork of goleveldb: active-compaction stats, DB.FlushMemdb, and + // OpenTransaction releasing the write lock when it fails. + github.com/syndtr/goleveldb => github.com/keybase/goleveldb v1.0.1-0.20260922194017-e81a99618c6c mvdan.cc/xurls/v2 => github.com/keybase/xurls/v2 v2.0.1-0.20190725180013-1e015cacd06c ) diff --git a/go/go.sum b/go/go.sum index cf8ebe21abc5..85fda098f8a6 100644 --- a/go/go.sum +++ b/go/go.sum @@ -304,8 +304,8 @@ github.com/keybase/go-winio v0.4.12-0.20181031203417-0903bf878a72 h1:okCFPNP1a2X github.com/keybase/go-winio v0.4.12-0.20181031203417-0903bf878a72/go.mod h1:Rcswqyeiwun4CF+RpzSNllKs3nO8Es5HZIhl+8YCm94= github.com/keybase/golang-ico v0.0.0-20181117022008-819cbeb217c9 h1:O4kEXd3yzbpHCRqwjbsBNKlanp52zXjKP6nFh9VSGy0= github.com/keybase/golang-ico v0.0.0-20181117022008-819cbeb217c9/go.mod h1:t9Db0X8VAAJTL82huV2G56KTYUeVET8THTLSjyWpE4o= -github.com/keybase/goleveldb v1.0.1-0.20221007195407-9881c0c26e65 h1:+y5yl2o/36at0ANWlFv7orF19GY1TghEsWRYD3JIT0g= -github.com/keybase/goleveldb v1.0.1-0.20221007195407-9881c0c26e65/go.mod h1:RRCYJbIwD5jmqPI9XoAFR0OcDxqUctll6zUj/+B4S48= +github.com/keybase/goleveldb v1.0.1-0.20260922194017-e81a99618c6c h1:R9SY/sdXgYKASxIjuLfeWNXiVBT9ic8nrCzpsT0tx5k= +github.com/keybase/goleveldb v1.0.1-0.20260922194017-e81a99618c6c/go.mod h1:RRCYJbIwD5jmqPI9XoAFR0OcDxqUctll6zUj/+B4S48= github.com/keybase/gomounts v0.0.0-20180302000443-349507f4d353 h1:7+JGi0kc98ugL6PYpihS+nDD+nv/Fj8hGRvlN7595bg= github.com/keybase/gomounts v0.0.0-20180302000443-349507f4d353/go.mod h1:LJNWKEO+J6j4hj9xQWPWnNT8344YyuGKiSCjyJAAmbI= github.com/keybase/keybase-test-vectors v1.0.12-0.20200309162119-ea1e58fecd5d h1:AuV96cODzV3oEHhMcNr+7t2RxfxkAV6WpCWlKLP4mOo= diff --git a/go/libkb/leveldb.go b/go/libkb/leveldb.go index 7a048acce3ed..48538622b865 100644 --- a/go/libkb/leveldb.go +++ b/go/libkb/leveldb.go @@ -11,6 +11,7 @@ import ( "path/filepath" "strings" "sync" + "sync/atomic" "github.com/syndtr/goleveldb/leveldb" errors "github.com/syndtr/goleveldb/leveldb/errors" @@ -118,7 +119,11 @@ type LevelDb struct { // rather than the DB itself. More specifically, close does Lock(), while // other DB operations does RLock(). sync.RWMutex - db *leveldb.DB + // db is an atomic.Pointer rather than a plain field guarded by RLock/Lock + // because the lazy open's assignment runs under the read lock (shared), + // so a plain field would race against other readers that don't go + // through dbOpenerOnce, such as Flush. + db atomic.Pointer[leveldb.DB] dbOpenerOnce *sync.Once cleaner *levelDbCleaner @@ -153,7 +158,7 @@ func (l *LevelDb) Opts() *opt.Options { } } -func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) { +func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func(db *leveldb.DB) error) (err error) { err = func() error { l.RLock() defer l.RUnlock() @@ -163,13 +168,14 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) l.dbOpenerOnce.Do(func() { l.G().Log.Debug("+ LevelDb.open") fn := l.GetFilename() - l.G().Log.Debug("| Opening LevelDB for local cache: %v %s", l, fn) + l.G().Log.Debug("| Opening LevelDB for local cache: %s", fn) l.G().Log.Debug("| Opening LevelDB options: %+v", l.Opts()) - l.db, err = leveldb.OpenFile(fn, l.Opts()) + db, openErr := leveldb.OpenFile(fn, l.Opts()) + err = openErr if _, ok := err.(*errors.ErrCorrupted); ok { l.G().Log.Debug("| LevelDb was corrupted; attempting recovery (%v)", err) var recoveryError error - l.db, recoveryError = leveldb.RecoverFile(fn, nil) + db, recoveryError = leveldb.RecoverFile(fn, nil) if recoveryError != nil { l.G().Log.Debug("| Recovery failed: %v", recoveryError) } else { @@ -179,8 +185,9 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) } } l.G().Log.Debug("- LevelDb.open -> %s", ErrToOk(err)) - if l.db != nil { - l.cleaner.setDb(l.db) + l.db.Store(db) + if db != nil { + l.cleaner.start(db) } }) @@ -188,15 +195,22 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) return err } - if l.db == nil { + db := l.db.Load() + if db == nil { // This means DB is already closed. We are preventing lazy-opening after // closing, so just return error here. return LevelDBOpenClosedError{} } - return action() + return action(db) }() + return l.handleOpError(err) +} +// handleOpError recovers from an error a db operation returned: it nukes a +// corrupt db, starts a forced clean when the disk is full, and lets a failed +// lazy open be retried. +func (l *LevelDb) handleOpError(err error) error { // If the file is corrupt, just nuke and act like we didn't find anything if l.nukeIfCorrupt(err) { err = nil @@ -214,7 +228,7 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) // we should at least try instead of auto returning LevelDBOpenClosederror. if err != nil { l.Lock() - if l.db == nil { + if l.db.Load() == nil { l.G().Log.Debug("LevelDb: doWhileOpenAndNukeIfCorrupted: resetting sync one: %s", err) l.dbOpenerOnce = new(sync.Once) } @@ -227,45 +241,37 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) // where we want to get around the lazy open and make sure we can // use it later. func (l *LevelDb) ForceOpen() error { - return l.doWhileOpenAndNukeIfCorrupted(func() error { return nil }) + return l.doWhileOpenAndNukeIfCorrupted(func(*leveldb.DB) error { return nil }) } -// levelDbFlushSentinelKey lives in the "pm" table so the db cleaner ignores -// it. Written before CompactRange so the memtable contains at least one key -// and isMemOverlaps returns true for the full-range compaction. -var levelDbFlushSentinelKey = []byte(levelDbTablePerm + ":ff:flush-sentinel") - -// Flush writes the current memtable to disk and rotates the journal. An +// Flush writes the current memtable to disk and starts an empty journal. An // unclean process kill (routine on iOS) with a non-empty journal forces a // journal replay on the next open, or worse a whole-DB recovery if the // journal tail is corrupt — both of which block startup. Flushing while -// entering the background leaves a near-empty journal so the next cold start +// entering the background leaves an empty journal so the next cold start // opens fast. No-op if the DB is not currently open; does not trigger a lazy // open. -func (l *LevelDb) Flush() (err error) { - defer convertNoSpaceError(&err) - l.RLock() - defer l.RUnlock() - if l.db == nil { - return nil - } - // Write the sentinel so the memtable is non-empty; then compact the full - // key space (util.Range{} with nil Start/Limit) so isMemOverlaps always - // returns true regardless of what other keys are live. A narrow range - // keyed only on the sentinel could miss the memtable flush if a concurrent - // write rotated the memtable between the Put and CompactRange. - if err = l.db.Put(levelDbFlushSentinelKey, nil, nil); err != nil { - return err - } - if err = l.db.CompactRange(util.Range{}); err != nil { - return err - } - return l.db.Delete(levelDbFlushSentinelKey, nil) +// +// Writers are blocked only while the memtable is swapped out, not while it is +// written. A call returns once every write made before it is in a table, +// including a memtable that a concurrent call or a full write buffer rotated. +func (l *LevelDb) Flush() error { + err := func() (err error) { + defer convertNoSpaceError(&err) + l.RLock() + defer l.RUnlock() + db := l.db.Load() + if db == nil { + return nil + } + return db.FlushMemdb() + }() + return l.handleOpError(err) } func (l *LevelDb) Stats() (stats string) { - if err := l.doWhileOpenAndNukeIfCorrupted(func() (err error) { - stats, err = l.db.GetProperty("leveldb.stats") + if err := l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) (err error) { + stats, err = db.GetProperty("leveldb.stats") stats = fmt.Sprintf("%s\n%s", stats, l.cleaner.Status()) return err }); err != nil { @@ -276,8 +282,8 @@ func (l *LevelDb) Stats() (stats string) { func (l *LevelDb) CompactionStats() (memActive, tableActive bool, err error) { var dbStats leveldb.DBStats - if err := l.doWhileOpenAndNukeIfCorrupted(func() (err error) { - return l.db.Stats(&dbStats) + if err := l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) (err error) { + return db.Stats(&dbStats) }); err != nil { return false, false, err } @@ -299,17 +305,18 @@ func (l *LevelDb) Close() error { func (l *LevelDb) closeLocked() error { var err error - if l.db != nil { + if db := l.db.Load(); db != nil { l.G().Log.Debug("Closing LevelDB local cache: %s", l.GetFilename()) - err = l.db.Close() - l.db = nil + // Stop any active cleaning job first: goleveldb must not be closed + // under an open iterator. + l.cleaner.Stop() + l.cleaner.Shutdown() + err = db.Close() + l.db.Store(nil) // In case we just nuked DB and reset the dbOpenerOnce, this makes sure it // doesn't open the DB again. l.dbOpenerOnce.Do(func() {}) - // stop any active cleaning jobs - l.cleaner.Stop() - l.cleaner.Shutdown() } return err } @@ -356,6 +363,7 @@ func (l *LevelDb) Nuke() (fn string, err error) { if err = os.RemoveAll(fn); err != nil { return fn, err } + l.cleaner.forgetPosition() // reset dbOpenerOnce since this is not a explicit close and there might be // more legitimate DB operations coming in l.dbOpenerOnce = new(sync.Once) @@ -375,30 +383,30 @@ func (l *LevelDb) nukeIfCorrupt(err error) bool { } func (l *LevelDb) Put(id DbKey, aliases []DbKey, value []byte) error { - return l.doWhileOpenAndNukeIfCorrupted(func() error { - return levelDbPut(l.db, l.cleaner, id, aliases, value) + return l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) error { + return levelDbPut(db, l.cleaner, id, aliases, value) }) } func (l *LevelDb) Get(id DbKey) (val []byte, found bool, err error) { - err = l.doWhileOpenAndNukeIfCorrupted(func() error { - val, found, err = levelDbGet(l.db, l.cleaner, id) + err = l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) error { + val, found, err = levelDbGet(db, l.cleaner, id) return err }) return val, found, err } func (l *LevelDb) Lookup(id DbKey) (val []byte, found bool, err error) { - err = l.doWhileOpenAndNukeIfCorrupted(func() error { - val, found, err = levelDbLookup(l.db, l.cleaner, id) + err = l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) error { + val, found, err = levelDbLookup(db, l.cleaner, id) return err }) return val, found, err } func (l *LevelDb) Delete(id DbKey) error { - return l.doWhileOpenAndNukeIfCorrupted(func() error { - return levelDbDelete(l.db, l.cleaner, id) + return l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) error { + return levelDbDelete(db, l.cleaner, id) }) } @@ -407,7 +415,18 @@ func (l *LevelDb) OpenTransaction() (LocalDbTransaction, error) { ltr LevelDbTransaction err error ) - if ltr.tr, err = l.db.OpenTransaction(); err != nil { + // Open the db lazily, but wait on goleveldb's write lock outside our read + // lock: a Close or Nuke queued behind a held read lock would block every + // new reader, including the current transaction's holder. If the db + // closes first, goleveldb returns ErrClosed. + var db *leveldb.DB + if err = l.doWhileOpenAndNukeIfCorrupted(func(opened *leveldb.DB) error { + db = opened + return nil + }); err != nil { + return LevelDbTransaction{}, err + } + if ltr.tr, err = db.OpenTransaction(); err != nil { return LevelDbTransaction{}, err } ltr.cleaner = l.cleaner @@ -416,10 +435,10 @@ func (l *LevelDb) OpenTransaction() (LocalDbTransaction, error) { func (l *LevelDb) KeysWithPrefixes(prefixes ...[]byte) (DBKeySet, error) { m := make(map[DbKey]struct{}) - err := l.doWhileOpenAndNukeIfCorrupted(func() error { + err := l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) error { opts := &opt.ReadOptions{DontFillCache: true} for _, prefix := range prefixes { - iter := l.db.NewIterator(util.BytesPrefix(prefix), opts) + iter := db.NewIterator(util.BytesPrefix(prefix), opts) for iter.Next() { _, dbKey, err := DbKeyParse(string(iter.Key())) if err != nil { diff --git a/go/libkb/leveldb_cleaner.go b/go/libkb/leveldb_cleaner.go index 891329c23049..6184893c0afb 100644 --- a/go/libkb/leveldb_cleaner.go +++ b/go/libkb/leveldb_cleaner.go @@ -60,52 +60,42 @@ type levelDbCleaner struct { MetaContextified sync.Mutex - running bool - lastKey []byte - lastRun time.Time - dbName string - config DbCleanerConfig - cache *lru.Cache - cacheMu sync.Mutex // protects the pointer to the cache - isMobile bool - db *leveldb.DB - stopCh chan struct{} - cancelCh chan struct{} - - isShutdown bool + running bool + lastKey []byte + lastRun time.Time + dbName string + config DbCleanerConfig + cache *lru.Cache + cacheMu sync.Mutex // protects the pointer to the cache + db *leveldb.DB + stopCh chan struct{} + // cleanDone is closed when the running clean exits. + cleanDone chan struct{} } func newLevelDbCleaner(mctx MetaContext, dbName string) *levelDbCleaner { config := DefaultDesktopDbCleanerConfig - isMobile := mctx.G().IsMobileAppType() - if isMobile { + if mctx.G().IsMobileAppType() { config = DefaultMobileDbCleanerConfig } - return newLevelDbCleanerWithConfig(mctx, dbName, config, isMobile) + return newLevelDbCleanerWithConfig(mctx, dbName, config) } -func newLevelDbCleanerWithConfig(mctx MetaContext, dbName string, config DbCleanerConfig, isMobile bool) *levelDbCleaner { +func newLevelDbCleanerWithConfig(mctx MetaContext, dbName string, config DbCleanerConfig) *levelDbCleaner { cache, err := lru.New(config.CacheCapacity) if err != nil { panic(err) } mctx = mctx.WithLogTag("DBCLN") - c := &levelDbCleaner{ + return &levelDbCleaner{ MetaContextified: NewMetaContextified(mctx), // Start the run shortly after starting but not immediately - lastRun: mctx.G().GetClock().Now().Add(-(config.CleanInterval - config.CleanInterval/10)), - dbName: dbName, - config: config, - cache: cache, - isMobile: isMobile, - stopCh: make(chan struct{}), - cancelCh: make(chan struct{}), + lastRun: mctx.G().GetClock().Now().Add(-(config.CleanInterval - config.CleanInterval/10)), + dbName: dbName, + config: config, + cache: cache, + stopCh: make(chan struct{}), } - if isMobile { - stopCh := c.stopCh - go c.monitorAppState(stopCh) - } - return c } func (c *levelDbCleaner) getCache() *lru.Cache { @@ -115,53 +105,56 @@ func (c *levelDbCleaner) getCache() *lru.Cache { } func (c *levelDbCleaner) Status() string { + cacheSize := c.getCache().Len() + c.Lock() + defer c.Unlock() return fmt.Sprintf("levelDbCleaner{cacheSize: %d, lastRun: %v, lastKey: %v, running: %v}\n%v\n", - c.cache.Len(), c.lastRun, c.lastKey, c.running, c.config) + cacheSize, c.lastRun, c.lastKey, c.running, c.config) } +// Stop detaches the cleaner from its db and waits for a running clean to +// exit, so the caller can close the db without an iterator still open on it. func (c *levelDbCleaner) Stop() { c.log("Stop") c.Lock() - defer c.Unlock() - if c.stopCh != nil { - close(c.stopCh) - c.stopCh = make(chan struct{}) + close(c.stopCh) + c.stopCh = make(chan struct{}) + c.db = nil + var cleanDone chan struct{} + if c.running { + cleanDone = c.cleanDone } -} - -func (c *levelDbCleaner) monitorAppState(stopCh chan struct{}) { - c.log("monitorAppState") - state := keybase1.MobileAppState_FOREGROUND - for { - select { - case <-c.G().MobileAppState.NextUpdate(state): - state = c.G().MobileAppState.State() - switch state { - case keybase1.MobileAppState_BACKGROUNDACTIVE: - default: - c.log("monitorAppState: attempting cancel, state: %v", state) - c.Lock() - if c.cancelCh != nil { - close(c.cancelCh) - c.cancelCh = make(chan struct{}) - } - c.Unlock() - } - case <-stopCh: - c.log("monitorAppState: stop") - return - } + c.Unlock() + if cleanDone != nil { + <-cleanDone } } -func (c *levelDbCleaner) log(format string, args ...any) { - c.M().Debug(fmt.Sprintf("levelDbCleaner(%s): %s", c.dbName, format), args...) +// forgetPosition restarts cleaning from the first key, for when the db's +// contents are gone. +func (c *levelDbCleaner) forgetPosition() { + c.Lock() + defer c.Unlock() + c.lastKey = nil } -func (c *levelDbCleaner) setDb(db *leveldb.DB) { +// start attaches the cleaner to a newly opened db, undoing a previous +// Stop/Shutdown from closing it. +func (c *levelDbCleaner) start(db *leveldb.DB) { + cache, err := lru.New(c.config.CacheCapacity) + if err != nil { + panic(err) + } c.Lock() defer c.Unlock() c.db = db + c.cacheMu.Lock() + defer c.cacheMu.Unlock() + c.cache = cache +} + +func (c *levelDbCleaner) log(format string, args ...any) { + c.M().Debug(fmt.Sprintf("levelDbCleaner(%s): %s", c.dbName, format), args...) } func (c *levelDbCleaner) cacheKey(key []byte) string { @@ -169,14 +162,13 @@ func (c *levelDbCleaner) cacheKey(key []byte) string { } func (c *levelDbCleaner) clearCache() { - c.cache.Purge() + c.getCache().Purge() } func (c *levelDbCleaner) Shutdown() { c.cacheMu.Lock() defer c.cacheMu.Unlock() c.cache, _ = lru.New(1) - c.isShutdown = true } func (c *levelDbCleaner) shouldCleanLocked(force bool) bool { @@ -191,14 +183,11 @@ func (c *levelDbCleaner) shouldCleanLocked(force bool) bool { c.G().GetClock().Now().Sub(c.lastRun) >= c.config.CleanInterval } -func (c *levelDbCleaner) getDbSize() (size uint64, err error) { - if c.db == nil { - return 0, nil - } +func (c *levelDbCleaner) getDbSize(db *leveldb.DB) (size uint64, err error) { // get the size from the start of the kv table to the beginning of the perm // table since that is all we can clean dbRange := util.Range{Start: tablePrefix(levelDbTableKv), Limit: tablePrefix(levelDbTablePerm)} - sizes, err := c.db.SizeOf([]util.Range{dbRange}) + sizes, err := db.SizeOf([]util.Range{dbRange}) if err != nil { return 0, err } @@ -208,26 +197,44 @@ func (c *levelDbCleaner) getDbSize() (size uint64, err error) { func (c *levelDbCleaner) clean(force bool) (err error) { c.Lock() // get out without spamming the logs - if !c.shouldCleanLocked(force) { + if c.db == nil || !c.shouldCleanLocked(force) { c.Unlock() return nil } c.running = true + cleanDone := make(chan struct{}) + c.cleanDone = cleanDone + // A clean stays on the db it started with. Before that db is closed, Stop + // closes stopCh and waits for the clean to exit; a reopen attaches a new + // db for later cleans without this one ever touching it. + db := c.db key := c.lastKey stopCh := c.stopCh - cancelCh := c.cancelCh + // Sample the app state in the same critical section as running=true, so + // a transition that lands between here and the batch loop (during + // getDbSize, logging, etc.) is not missed: any caller who observes + // running via c.Lock() only does so after this sample is already taken. + // A clean gives way to the foreground: it keeps running only while the + // app state stays BACKGROUNDACTIVE (or never changes at all, as on + // desktop). NextUpdate collapses intermediate transitions, so a wake + // re-reads the current state rather than assuming what it changed to. + state := c.G().MobileAppState.State() + appCh := c.G().MobileAppState.NextUpdate(state) c.Unlock() defer c.M().Trace(fmt.Sprintf("levelDbCleaner(%s) clean, config: %v", c.dbName, c.config), &err)() defer func() { c.Lock() defer c.Unlock() - c.lastKey = key + if c.db == db { + c.lastKey = key + } c.lastRun = c.G().GetClock().Now() c.running = false + close(cleanDone) }() - dbSize, err := c.getDbSize() + dbSize, err := c.getDbSize(db) if err != nil { return err } @@ -242,17 +249,25 @@ func (c *levelDbCleaner) clean(force bool) (err error) { var totalNumPurged, numPurged int for i := range 100 { select { - case <-cancelCh: - c.log("aborting clean, %d runs, canceled", i) - return nil case <-stopCh: c.log("aborting clean %d runs, stopped", i) return nil + case <-appCh: + state = c.G().MobileAppState.State() + if state != keybase1.MobileAppState_BACKGROUNDACTIVE { + c.log("aborting clean, %d runs, left BACKGROUNDACTIVE for %v", i, state) + return nil + } + appCh = c.G().MobileAppState.NextUpdate(state) default: } start := c.G().GetClock().Now() - numPurged, key, err = c.cleanBatch(key) + numPurged, key, err = c.cleanBatch(db, stopCh, key) + if err == errCleanStopped { + c.log("aborting clean %d runs, stopped", i) + return nil + } if err != nil { return err } @@ -266,7 +281,7 @@ func (c *levelDbCleaner) clean(force bool) (err error) { numPurged, humanize.Bytes(dbSize), key, c.G().GetClock().Now().Sub(start)) } // check if we are within limits - dbSize, err = c.getDbSize() + dbSize, err = c.getDbSize(db) if err != nil { return err } @@ -274,13 +289,20 @@ func (c *levelDbCleaner) clean(force bool) (err error) { if !force && dbSize < c.config.HaltSize { break } - time.Sleep(c.config.SleepInterval) + select { + case <-stopCh: + c.log("aborting clean %d runs, stopped", i) + return nil + case <-time.After(c.config.SleepInterval): + } } c.log("clean complete. purged %d items total, dbSize: %v", totalNumPurged, humanize.Bytes(dbSize)) return nil } -func (c *levelDbCleaner) cleanBatch(startKey []byte) (int, []byte, error) { +var errCleanStopped = errors.New("levelDbCleaner: stopped") + +func (c *levelDbCleaner) cleanBatch(db *leveldb.DB, stopCh chan struct{}, startKey []byte) (int, []byte, error) { // Start our range from wherever we left off last time, and clean up until // the permanent entries table begins. iterRange := &util.Range{Start: startKey, Limit: tablePrefix(levelDbTablePerm)} @@ -290,18 +312,18 @@ func (c *levelDbCleaner) cleanBatch(startKey []byte) (int, []byte, error) { // caching so that the data processed by the bulk read does not end up // displacing most of the cached contents.""" opts := &opt.ReadOptions{DontFillCache: true} - iter := c.db.NewIterator(iterRange, opts) + iter := db.NewIterator(iterRange, opts) batch := new(leveldb.Batch) for batch.Len() < 1000 && iter.Next() { key := iter.Key() - c.cacheMu.Lock() - if c.isShutdown { - c.cacheMu.Unlock() - return 0, nil, errors.New("cleanBatch: cancelled due to shutdown") + select { + case <-stopCh: + iter.Release() + return 0, startKey, errCleanStopped + default: } - cache := c.cache - c.cacheMu.Unlock() + cache := c.getCache() if _, found := cache.Get(c.cacheKey(key)); !found { cp := make([]byte, len(key)) @@ -324,11 +346,11 @@ func (c *levelDbCleaner) cleanBatch(startKey []byte) (int, []byte, error) { if err := iter.Error(); err != nil { return 0, nil, err } - if err := c.db.Write(batch, nil); err != nil { + if err := db.Write(batch, nil); err != nil { return 0, nil, err } // Compact the range we just deleted in so the size changes are reflected - err := c.db.CompactRange(util.Range{Start: startKey, Limit: key}) + err := db.CompactRange(util.Range{Start: startKey, Limit: key}) return batch.Len(), key, err } diff --git a/go/libkb/leveldb_cleaner_test.go b/go/libkb/leveldb_cleaner_test.go new file mode 100644 index 000000000000..cfbe6a69638c --- /dev/null +++ b/go/libkb/leveldb_cleaner_test.go @@ -0,0 +1,335 @@ +package libkb + +import ( + "fmt" + "path/filepath" + "testing" + "time" + + keybase1 "github.com/keybase/client/go/protocol/keybase1" + "github.com/stretchr/testify/require" +) + +// newMobileCleanerDb makes a LevelDb whose cleaner behaves as on mobile. The +// db is not opened. +func newMobileCleanerDb(t *testing.T, tc *TestContext, config DbCleanerConfig) *LevelDb { + dir := t.TempDir() + db := NewLevelDb(tc.G, func() string { return filepath.Join(dir, "test.leveldb") }) + db.cleaner = newLevelDbCleanerWithConfig(NewMetaContextTODO(tc.G), "test", config) + t.Cleanup(func() { _ = db.Close() }) + return db +} + +func testCleanerConfig() DbCleanerConfig { + config := DefaultMobileDbCleanerConfig + config.CacheCapacity = 10 + return config +} + +var cleanerStates = []keybase1.MobileAppState{ + keybase1.MobileAppState_FOREGROUND, + keybase1.MobileAppState_BACKGROUNDACTIVE, + keybase1.MobileAppState_INACTIVE, + keybase1.MobileAppState_BACKGROUND, +} + +// waitCleanerRunning waits until a started clean has taken the running flag. +// clean() samples the app state in the same critical section as setting +// running, so a caller who observes running via this has already lost any +// race against that sample. +func waitCleanerRunning(t *testing.T, c *levelDbCleaner) { + t.Helper() + require.Eventually(t, func() bool { + c.Lock() + defer c.Unlock() + return c.running + }, 10*time.Second, time.Millisecond, "clean did not start running") +} + +// putKeys writes numKeys keys under db and returns the last one. +func putKeys(t *testing.T, db *LevelDb, numKeys int) DbKey { + t.Helper() + var last DbKey + for i := range numKeys { + last = DbKey{Key: fmt.Sprintf("k%05d", i), Typ: 0} + require.NoError(t, db.Put(last, nil, []byte{1})) + } + return last +} + +// A clean in progress stops before finishing when the app leaves +// BACKGROUNDACTIVE for any other state. +func TestCleanerStopsWhenLeavingBackgroundActive(t *testing.T) { + for _, next := range cleanerStates { + if next == keybase1.MobileAppState_BACKGROUNDACTIVE { + continue + } + t.Run(next.String(), func(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-stop", 0) + defer tc.Cleanup() + config := testCleanerConfig() + config.SleepInterval = 100 * time.Millisecond + db := newMobileCleanerDb(t, &tc, config) + tc.G.MobileAppState.Update(keybase1.MobileAppState_BACKGROUNDACTIVE) + + const numKeys = 3500 + lastKey := putKeys(t, db, numKeys) + db.cleaner.clearCache() + + done := make(chan error, 1) + go func() { done <- db.cleaner.clean(true /* force */) }() + waitCleanerRunning(t, db.cleaner) + + tc.G.MobileAppState.Update(next) + require.NoError(t, <-done) + + _, found, err := db.Get(lastKey) + require.NoError(t, err) + require.True(t, found, "a clean canceled by leaving BACKGROUNDACTIVE should not reach the last key") + }) + } +} + +// A clean that starts outside BACKGROUNDACTIVE is also interrupted by a +// transition to a different non-BACKGROUNDACTIVE state: cancellation depends +// on the landing state, not on where the clean started. +func TestCleanerStopsOnTransitionBetweenNonBackgroundActiveStates(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-stop-fg", 0) + defer tc.Cleanup() + config := testCleanerConfig() + config.SleepInterval = 100 * time.Millisecond + db := newMobileCleanerDb(t, &tc, config) + tc.G.MobileAppState.Update(keybase1.MobileAppState_FOREGROUND) + + const numKeys = 3500 + lastKey := putKeys(t, db, numKeys) + db.cleaner.clearCache() + + done := make(chan error, 1) + go func() { done <- db.cleaner.clean(true /* force */) }() + waitCleanerRunning(t, db.cleaner) + + tc.G.MobileAppState.Update(keybase1.MobileAppState_BACKGROUND) + require.NoError(t, <-done) + + _, found, err := db.Get(lastKey) + require.NoError(t, err) + require.True(t, found, "a clean canceled by a transition between non-BACKGROUNDACTIVE states should not reach the last key") +} + +// A clean keeps running, with no early return, for as long as the app state +// stays BACKGROUNDACTIVE, including across an unrelated update that collapses +// to a no-op (NextUpdate only fires on a real change). +func TestCleanerContinuesWhileBackgroundActive(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-continue", 0) + defer tc.Cleanup() + config := testCleanerConfig() + config.SleepInterval = 100 * time.Millisecond + db := newMobileCleanerDb(t, &tc, config) + tc.G.MobileAppState.Update(keybase1.MobileAppState_BACKGROUNDACTIVE) + + const numKeys = 3500 + lastKey := putKeys(t, db, numKeys) + db.cleaner.clearCache() + + done := make(chan error, 1) + go func() { done <- db.cleaner.clean(true /* force */) }() + waitCleanerRunning(t, db.cleaner) + + // A same-value update: no real transition, so it must not interrupt the + // clean. + tc.G.MobileAppState.Update(keybase1.MobileAppState_BACKGROUNDACTIVE) + require.NoError(t, <-done) + + _, found, err := db.Get(lastKey) + require.NoError(t, err) + require.False(t, found, "a clean that never left BACKGROUNDACTIVE should run to completion") +} + +// A clean that starts outside BACKGROUNDACTIVE and then transitions into it +// keeps running: the wake re-arms rather than treating the change itself as +// a cancellation. +func TestCleanerRearmsIntoBackgroundActive(t *testing.T) { + for _, start := range []keybase1.MobileAppState{ + keybase1.MobileAppState_FOREGROUND, + keybase1.MobileAppState_INACTIVE, + keybase1.MobileAppState_BACKGROUND, + } { + t.Run(start.String(), func(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-rearm", 0) + defer tc.Cleanup() + config := testCleanerConfig() + config.SleepInterval = 100 * time.Millisecond + db := newMobileCleanerDb(t, &tc, config) + tc.G.MobileAppState.Update(start) + + const numKeys = 3500 + lastKey := putKeys(t, db, numKeys) + db.cleaner.clearCache() + + done := make(chan error, 1) + go func() { done <- db.cleaner.clean(true /* force */) }() + waitCleanerRunning(t, db.cleaner) + + tc.G.MobileAppState.Update(keybase1.MobileAppState_BACKGROUNDACTIVE) + require.NoError(t, <-done) + + _, found, err := db.Get(lastKey) + require.NoError(t, err) + require.False(t, found, "a clean that transitions into BACKGROUNDACTIVE should run to completion") + }) + } +} + +// A cleaner still cleans after its db is reopened (Nuke, or Close + +// ForceOpen): start() must install a working cache, since Shutdown left a +// one-entry cache that forgets everything. +func TestCleanerCleansAfterReopen(t *testing.T) { + for _, name := range []string{"nuke", "close"} { + t.Run(name, func(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-reopen", 0) + defer tc.Cleanup() + db := newMobileCleanerDb(t, &tc, testCleanerConfig()) + require.NoError(t, db.ForceOpen()) + + switch name { + case "nuke": + _, err := db.Nuke() + require.NoError(t, err) + require.NoError(t, db.ForceOpen()) + case "close": + require.NoError(t, db.Close()) + // The first use after Close fails and rearms the lazy open. + require.Error(t, db.ForceOpen()) + require.NoError(t, db.ForceOpen()) + } + + key := DbKey{Key: "reopen-key", Typ: 0} + require.NoError(t, db.Put(key, nil, []byte{1})) + db.cleaner.clearCache() + require.NoError(t, db.cleaner.clean(true /* force */)) + + _, found, err := db.Get(key) + require.NoError(t, err) + require.False(t, found, "clean after %s left the key", name) + }) + } +} + +// A reopened db is a different key space, so a clean on it starts from the +// beginning rather than where a clean of the old db left off. +func TestCleanerStartsFromBeginningAfterReopen(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-reopen-lastkey", 0) + defer tc.Cleanup() + db := newMobileCleanerDb(t, &tc, testCleanerConfig()) + require.NoError(t, db.ForceOpen()) + + key := DbKey{Key: "aaaa", Typ: 0} + db.cleaner.Lock() + db.cleaner.lastKey = DbKey{Key: "zzzz", Typ: 0}.ToBytes() + db.cleaner.Unlock() + + _, err := db.Nuke() + require.NoError(t, err) + require.NoError(t, db.Put(key, nil, []byte{1})) + db.cleaner.clearCache() + require.NoError(t, db.cleaner.clean(true /* force */)) + + _, found, err := db.Get(key) + require.NoError(t, err) + require.False(t, found, "clean skipped keys below the old db's lastKey") +} + +// Status reads cleaner state that a reopen replaces. +func TestCleanerStatusDuringReopen(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-status-reopen", 0) + defer tc.Cleanup() + db := newMobileCleanerDb(t, &tc, testCleanerConfig()) + require.NoError(t, db.ForceOpen()) + + stop := make(chan struct{}) + done := make(chan struct{}) + go func() { + defer close(done) + for { + select { + case <-stop: + return + default: + _ = db.cleaner.Status() + db.cleaner.clearCache() + } + } + }() + for range 20 { + _, err := db.Nuke() + require.NoError(t, err) + require.NoError(t, db.ForceOpen()) + } + close(stop) + <-done +} + +// Close waits for a running clean to exit before closing the db under it, +// and a clean that sleeps between batches exits promptly when stopped. +func TestCleanerCloseWaitsForRunningClean(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-close-waits", 0) + defer tc.Cleanup() + config := testCleanerConfig() + config.SleepInterval = time.Minute + db := newMobileCleanerDb(t, &tc, config) + + putKeys(t, db, 3500) + db.cleaner.clearCache() + done := make(chan error, 1) + go func() { done <- db.cleaner.clean(true /* force */) }() + // Wait for the first batch to land, so the clean is headed for its sleep. + require.Eventually(t, func() bool { + keys, err := db.KeysWithPrefixes(tablePrefix(levelDbTableKv)) + require.NoError(t, err) + return len(keys) < 3500 + }, 10*time.Second, time.Millisecond, "the first batch never landed") + // The batch's compaction and size check run before the sleep. + time.Sleep(500 * time.Millisecond) + + start := time.Now() + require.NoError(t, db.Close()) + require.Less(t, time.Since(start), 10*time.Second, "Close waited out the clean's sleep") + db.cleaner.Lock() + running := db.cleaner.running + db.cleaner.Unlock() + require.False(t, running, "Close returned with a clean still running on the closed db") + require.NoError(t, <-done) +} + +// A clean that starts after Close has no db to clean. +func TestCleanerAfterCloseIsNoop(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-after-close", 0) + defer tc.Cleanup() + db := newMobileCleanerDb(t, &tc, testCleanerConfig()) + require.NoError(t, db.ForceOpen()) + require.NoError(t, db.Close()) + require.NoError(t, db.cleaner.clean(true /* force */)) +} + +// Close and reopen keeps the same data, so cleaning resumes where it was. +func TestCleanerKeepsPositionAcrossCloseReopen(t *testing.T) { + tc := SetupTest(t, "LevelDb-cleaner-close-reopen-lastkey", 0) + defer tc.Cleanup() + db := newMobileCleanerDb(t, &tc, testCleanerConfig()) + require.NoError(t, db.ForceOpen()) + + lastKey := DbKey{Key: "mmmm", Typ: 0}.ToBytes() + db.cleaner.Lock() + db.cleaner.lastKey = lastKey + db.cleaner.Unlock() + + require.NoError(t, db.Close()) + // The first use after Close fails and rearms the lazy open. + require.Error(t, db.ForceOpen()) + require.NoError(t, db.ForceOpen()) + + db.cleaner.Lock() + defer db.cleaner.Unlock() + require.Equal(t, lastKey, db.cleaner.lastKey) +} diff --git a/go/libkb/leveldb_test.go b/go/libkb/leveldb_test.go index e40565abfa6f..e275ba5528f2 100644 --- a/go/libkb/leveldb_test.go +++ b/go/libkb/leveldb_test.go @@ -12,6 +12,7 @@ import ( "testing" "time" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/syndtr/goleveldb/leveldb" ) @@ -69,6 +70,34 @@ func doSomeIO() error { return os.WriteFile(filepath.Join(dir, "some-io"), []byte("O_O"), 0o600) } +func levelDbStats(t *testing.T, db *LevelDb) (stats leveldb.DBStats) { + require.NoError(t, db.doWhileOpenAndNukeIfCorrupted(func(ldb *leveldb.DB) error { + return ldb.Stats(&stats) + })) + return stats +} + +func levelDbTableCount(t *testing.T, db *LevelDb) (count int) { + for _, n := range levelDbStats(t, db).LevelTablesCounts { + count += n + } + return count +} + +// levelDbJournalSize returns the size of the journal (*.log) files, which +// hold writes not yet flushed to a table. +func levelDbJournalSize(t *testing.T, db *LevelDb) (size int64) { + journals, err := filepath.Glob(filepath.Join(db.GetFilename(), "*.log")) + require.NoError(t, err) + require.NotEmpty(t, journals) + for _, j := range journals { + fi, err := os.Stat(j) + require.NoError(t, err) + size += fi.Size() + } + return size +} + func testLevelDbPut(db *LevelDb) (key DbKey, err error) { key = DbKey{Key: "test-key", Typ: 0} v := []byte{1, 2, 3, 4} @@ -123,23 +152,217 @@ func TestLevelDb(t *testing.T) { key, err := testLevelDbPut(db) require.NoError(t, err) + require.Zero(t, levelDbTableCount(t, db), "the put should still be in the memtable") require.NoError(t, db.Flush()) + require.NotZero(t, levelDbTableCount(t, db), "flush should write the memtable to a table") require.NoError(t, db.Flush()) - // Data survives the flush and the sentinel is cleaned up. + // Data survives the flush. val, found, err := db.Get(key) require.NoError(t, err) require.True(t, found) require.Equal(t, []byte{1, 2, 3, 4}, val) - _, err = db.db.Get(levelDbFlushSentinelKey, nil) - require.Equal(t, leveldb.ErrNotFound, err) // Writes still work after a flush. _, err = testLevelDbPut(db) require.NoError(t, err) }, }, + { + name: "flush-memtable-only", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-flush-memtable-only", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + require.NoError(t, db.ForceOpen()) + + putAcrossPrefixes := func(round int) { + for _, prefix := range []string{"aa", "kv", "lo", "pm", "zz"} { + for i := 0; i < 20; i++ { + key := []byte(fmt.Sprintf("%s:%d:%d", prefix, round, i)) + require.NoError(t, db.db.Load().Put(key, bytes.Repeat([]byte{byte(i)}, 100), nil)) + } + } + } + // Existing tables spanning the whole key space, so a table + // compaction of the flushed memtable would have inputs. + for round := 0; round < 2; round++ { + putAcrossPrefixes(round) + tr, err := db.db.Load().OpenTransaction() + require.NoError(t, err) + tr.Discard() + } + putAcrossPrefixes(2) + require.NotZero(t, levelDbJournalSize(t, db)) + before := levelDbStats(t, db).LevelTablesCounts + beforeTotal := levelDbTableCount(t, db) + + require.NoError(t, db.Flush()) + + after := levelDbStats(t, db).LevelTablesCounts + for level, n := range before { + require.GreaterOrEqual(t, after[level], n, "no table should be compacted away (level %d)", level) + } + require.Equal(t, beforeTotal+1, levelDbTableCount(t, db), "flush should add exactly one table") + require.Zero(t, levelDbJournalSize(t, db), "the flushed memtable's journal should be gone") + val, err := db.db.Load().Get([]byte("zz:2:19"), nil) + require.NoError(t, err) + require.Equal(t, bytes.Repeat([]byte{19}, 100), val) + }, + }, + { + // 8 goroutines call Flush with writes interleaved: every call + // returns nil, and every writer's last write is durable and + // readable once all goroutines finish. + name: "flush-concurrent", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-flush-concurrent", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + require.NoError(t, db.ForceOpen()) + + const writers, iterations = 8, 25 + var wg sync.WaitGroup + for w := 0; w < writers; w++ { + wg.Add(1) + go func(w int) { + defer wg.Done() + for i := 0; i < iterations; i++ { + key := DbKey{Key: fmt.Sprintf("%d-%d", w, i), Typ: 0} + assert.NoError(t, db.Put(key, nil, []byte{byte(i)})) + assert.NoError(t, db.Flush()) + } + }(w) + } + wg.Wait() + + require.Zero(t, levelDbJournalSize(t, db), "the last writes must be flushed") + for w := 0; w < writers; w++ { + _, found, err := db.Get(DbKey{Key: fmt.Sprintf("%d-%d", w, iterations-1), Typ: 0}) + require.NoError(t, err) + require.True(t, found) + } + }, + }, + { + // OpenTransaction opens the db lazily like every other operation. + name: "open-transaction-first", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-transaction-first", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + + tr, err := db.OpenTransaction() + require.NoError(t, err) + key := DbKey{Key: "tr-key", Typ: 0} + require.NoError(t, tr.Put(key, nil, []byte{1})) + require.NoError(t, tr.Commit()) + _, found, err := db.Get(key) + require.NoError(t, err) + require.True(t, found) + }, + }, + { + // An OpenTransaction waiting on goleveldb's write lock must not hold + // our read lock: a Nuke queued behind it would block the Get that + // the transaction holder needs before it can finish. + name: "open-transaction-waiting-vs-nuke", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-transaction-nuke", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + + tr, err := db.OpenTransaction() + require.NoError(t, err) + waiting := make(chan struct{}) + go func() { + close(waiting) + if tr2, err := db.OpenTransaction(); err == nil { + tr2.Discard() + } + }() + <-waiting + time.Sleep(100 * time.Millisecond) + nuked := make(chan struct{}) + go func() { + defer close(nuked) + _, err := db.Nuke() + assert.NoError(t, err) + }() + time.Sleep(100 * time.Millisecond) + + got := make(chan struct{}) + go func() { + defer close(got) + _, _, _ = db.Get(DbKey{Key: "test-key", Typ: 0}) + }() + select { + case <-got: + case <-time.After(5 * time.Second): + tr.Discard() + t.Fatal("Get deadlocked behind a Nuke queued on a waiting OpenTransaction") + } + tr.Discard() + <-nuked + }, + }, + { + name: "open-transaction-after-close", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-transaction-closed", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + + require.NoError(t, db.ForceOpen()) + require.NoError(t, db.Close()) + _, err = db.OpenTransaction() + require.ErrorAs(t, err, &LevelDBOpenClosedError{}) + }, + }, + { + name: "concurrent-open", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-concurrent-open", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + + // Under -race, this catches the lazy open assigning db.db while + // Flush reads it. + var wg sync.WaitGroup + for i := 0; i < 8; i++ { + wg.Add(2) + go func() { + defer wg.Done() + _, _, err := db.Get(DbKey{Key: "test-key", Typ: 0}) + assert.NoError(t, err) + }() + go func() { + defer wg.Done() + assert.NoError(t, db.Flush()) + }() + } + wg.Wait() + + // A lazy open racing a Nuke reopens rather than reporting closed. + for i := 0; i < 8; i++ { + wg.Add(2) + go func() { + defer wg.Done() + key := DbKey{Key: "test-key", Typ: 0} + assert.NoError(t, db.Put(key, nil, []byte{1})) + _, _, err := db.Get(key) + assert.NoError(t, err) + }() + go func() { + defer wg.Done() + _, err := db.Nuke() + assert.NoError(t, err) + }() + } + wg.Wait() + }, + }, { name: "cleaner", testBody: func(t *testing.T) { tc := SetupTest(t, "LevelDb-cleaner", 0) @@ -194,7 +417,7 @@ func TestLevelDb(t *testing.T) { // for sure they can happen concurrently. ch := make(chan struct{}) go func() { - _ = db.doWhileOpenAndNukeIfCorrupted(func() error { + _ = db.doWhileOpenAndNukeIfCorrupted(func(*leveldb.DB) error { defer wg.Done() select { case <-time.After(8 * time.Second): @@ -205,7 +428,7 @@ func TestLevelDb(t *testing.T) { }) }() go func() { - _ = db.doWhileOpenAndNukeIfCorrupted(func() error { + _ = db.doWhileOpenAndNukeIfCorrupted(func(*leveldb.DB) error { defer wg.Done() select { case <-time.After(8 * time.Second):