From 6fb79bf6ef4f83c57cd8ab220281db6ac3aadc04 Mon Sep 17 00:00:00 2001 From: chrisnojima Date: Mon, 21 Sep 2026 10:01:39 -0400 Subject: [PATCH 1/3] fix(leveldb): flush by rotating the memtable; open the db atomically Flush used to write a sentinel key and then CompactRange the whole key space just to force the memtable out. That rewrites tables that had nothing to do with the flush. Open a transaction and immediately discard it instead: goleveldb rotates a non-empty memtable and waits for it to land in a table, without compacting anything. Concurrent calls serialize on goleveldb's own write lock, so a call that arrives late still covers every write made before it. The lazily-opened db field was assigned under the *read* lock, racing every other reader. Make it an atomic.Pointer, and have OpenTransaction return LevelDBOpenClosedError instead of dereferencing a nil db after Close. The cleaner no longer runs a separate monitorAppState goroutine with its own cancel channel; clean() samples State()/NextUpdate() inside the same critical section that takes the running flag, so a transition can't slip through between the sample and the batch loop. --- go/libkb/leveldb.go | 100 +++++++++----- go/libkb/leveldb_cleaner.go | 106 +++++++-------- go/libkb/leveldb_cleaner_test.go | 217 +++++++++++++++++++++++++++++++ go/libkb/leveldb_test.go | 193 ++++++++++++++++++++++++++- 4 files changed, 518 insertions(+), 98 deletions(-) create mode 100644 go/libkb/leveldb_cleaner_test.go diff --git a/go/libkb/leveldb.go b/go/libkb/leveldb.go index 7a048acce3ed..b5b777061d95 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,10 +119,17 @@ 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 openedDb(). + db atomic.Pointer[leveldb.DB] dbOpenerOnce *sync.Once cleaner *levelDbCleaner + // flushHook, if set, runs after each memtable rotation. Tests only. + flushHook func() + filename string Contextified } @@ -163,13 +171,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 +188,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,7 +198,7 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) return err } - if l.db == nil { + if l.db.Load() == nil { // This means DB is already closed. We are preventing lazy-opening after // closing, so just return error here. return LevelDBOpenClosedError{} @@ -214,7 +224,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) } @@ -230,42 +240,54 @@ func (l *LevelDb) ForceOpen() error { return l.doWhileOpenAndNukeIfCorrupted(func() 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. +// +// Concurrent calls serialize on goleveldb's own write lock rather than on +// anything of ours: OpenTransaction blocks until any transaction ahead of it +// commits or discards, and by the time it unblocks it rotates whatever +// memtable is current, so a call that lands after another one already +// covers any write made before it arrived. func (l *LevelDb) Flush() (err error) { + return l.flushMemtable() +} + +// openedDb returns the DB without triggering a lazy open, or nil if it isn't +// open. Callers must hold the read lock. +func (l *LevelDb) openedDb() *leveldb.DB { + return l.db.Load() +} + +func (l *LevelDb) flushMemtable() (err error) { defer convertNoSpaceError(&err) l.RLock() defer l.RUnlock() - if l.db == nil { + db := l.openedDb() + if 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 { + // Opening a transaction rotates a non-empty memtable and waits until it + // is written to a table, without compacting any tables. The + // transaction itself is not needed. + tr, err := db.OpenTransaction() + if err != nil { return err } - if err = l.db.CompactRange(util.Range{}); err != nil { - return err + tr.Discard() + if l.flushHook != nil { + l.flushHook() } - return l.db.Delete(levelDbFlushSentinelKey, nil) + return nil } func (l *LevelDb) Stats() (stats string) { if err := l.doWhileOpenAndNukeIfCorrupted(func() (err error) { - stats, err = l.db.GetProperty("leveldb.stats") + stats, err = l.db.Load().GetProperty("leveldb.stats") stats = fmt.Sprintf("%s\n%s", stats, l.cleaner.Status()) return err }); err != nil { @@ -277,7 +299,7 @@ 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) + return l.db.Load().Stats(&dbStats) }); err != nil { return false, false, err } @@ -299,10 +321,10 @@ 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 + 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. @@ -376,13 +398,13 @@ 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 levelDbPut(l.db.Load(), 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) + val, found, err = levelDbGet(l.db.Load(), l.cleaner, id) return err }) return val, found, err @@ -390,7 +412,7 @@ func (l *LevelDb) Get(id DbKey) (val []byte, found bool, err error) { 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) + val, found, err = levelDbLookup(l.db.Load(), l.cleaner, id) return err }) return val, found, err @@ -398,7 +420,7 @@ func (l *LevelDb) Lookup(id DbKey) (val []byte, found bool, err error) { func (l *LevelDb) Delete(id DbKey) error { return l.doWhileOpenAndNukeIfCorrupted(func() error { - return levelDbDelete(l.db, l.cleaner, id) + return levelDbDelete(l.db.Load(), l.cleaner, id) }) } @@ -407,7 +429,13 @@ func (l *LevelDb) OpenTransaction() (LocalDbTransaction, error) { ltr LevelDbTransaction err error ) - if ltr.tr, err = l.db.OpenTransaction(); err != nil { + l.RLock() + db := l.openedDb() + l.RUnlock() + if db == nil { + return LevelDbTransaction{}, LevelDBOpenClosedError{} + } + if ltr.tr, err = db.OpenTransaction(); err != nil { return LevelDbTransaction{}, err } ltr.cleaner = l.cleaner @@ -419,7 +447,7 @@ func (l *LevelDb) KeysWithPrefixes(prefixes ...[]byte) (DBKeySet, error) { err := l.doWhileOpenAndNukeIfCorrupted(func() error { opts := &opt.ReadOptions{DontFillCache: true} for _, prefix := range prefixes { - iter := l.db.NewIterator(util.BytesPrefix(prefix), opts) + iter := l.db.Load().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..c0989313241f 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{} + 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{} isShutdown bool } 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 { @@ -129,27 +119,18 @@ func (c *levelDbCleaner) Stop() { } } -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 +// start attaches the cleaner to a newly opened db, undoing a previous +// Stop/Shutdown from closing it. +func (c *levelDbCleaner) start(db *leveldb.DB) { + c.Lock() + defer c.Unlock() + c.db = db + c.cacheMu.Lock() + defer c.cacheMu.Unlock() + if c.isShutdown { + if cache, err := lru.New(c.config.CacheCapacity); err == nil { + c.cache = cache + c.isShutdown = false } } } @@ -158,12 +139,6 @@ func (c *levelDbCleaner) log(format string, args ...any) { c.M().Debug(fmt.Sprintf("levelDbCleaner(%s): %s", c.dbName, format), args...) } -func (c *levelDbCleaner) setDb(db *leveldb.DB) { - c.Lock() - defer c.Unlock() - c.db = db -} - func (c *levelDbCleaner) cacheKey(key []byte) string { return string(key) } @@ -215,7 +190,16 @@ func (c *levelDbCleaner) clean(force bool) (err error) { c.running = true 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)() @@ -242,12 +226,16 @@ 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: } diff --git a/go/libkb/leveldb_cleaner_test.go b/go/libkb/leveldb_cleaner_test.go new file mode 100644 index 000000000000..1d6b38d3298e --- /dev/null +++ b/go/libkb/leveldb_cleaner_test.go @@ -0,0 +1,217 @@ +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 reset isShutdown so a reopened cleaner's cache +// isn't stuck discarding 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) + }) + } +} diff --git a/go/libkb/leveldb_test.go b/go/libkb/leveldb_test.go index e40565abfa6f..11697f93d90e 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() error { + return db.db.Load().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,181 @@ 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) + }, + }, + { + // A write and a Flush call that a flush's own hook makes reentrantly + // must still be flushed before the outer call returns: the hook runs + // after the transaction is discarded, so goleveldb's write lock is + // already free and the nested call is a plain second flush. + name: "flush-reentrant", testBody: func(t *testing.T) { + tc := SetupTest(t, "LevelDb-flush-reentrant", 0) + defer tc.Cleanup() + db, err := createTempLevelDbForTest(&tc, &td) + require.NoError(t, err) + _, err = testLevelDbPut(db) + require.NoError(t, err) + + rotations := 0 + db.flushHook = func() { + rotations++ + if rotations == 1 { + require.NoError(t, db.db.Load().Put([]byte("kv:late"), []byte{1}, nil)) + require.NoError(t, db.Flush()) + } + } + require.NoError(t, db.Flush()) + require.Equal(t, 2, rotations) + require.Zero(t, levelDbJournalSize(t, db)) + }, + }, + { + // 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) + } + }, + }, + { + 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) From b19fe3185f72000a0098c6c5e3b426b60656545a Mon Sep 17 00:00:00 2001 From: chrisnojima Date: Tue, 22 Sep 2026 15:42:18 -0400 Subject: [PATCH 2/3] fix(leveldb): flush with the fork's FlushMemdb; keep cleans on their own db goleveldb's OpenTransaction keeps its write lock when rotating the memtable or waiting on compaction fails, e.g. a journal create hitting ENOSPC. Flushing through it could leave every later write, and Close, blocked for good. It also held the lock for the whole table write. Flush now calls FlushMemdb from the keybase/goleveldb fork, which releases the lock before waiting and on every error; the fork also fixes OpenTransaction itself. OpenTransaction goes through the lazy open like every other operation instead of reporting a never-opened db as closed. The open db is handed to each action rather than reloaded. A clean keeps the db and stop channel it started with, so a Nuke and reopen mid-clean can't redirect it to the new db. start() resets lastKey and always installs a fresh LRU, which retires isShutdown. Status and clearCache read the cache pointer under its lock. --- go/go.mod | 4 +- go/go.sum | 4 +- go/libkb/leveldb.go | 86 +++++++++++--------------------- go/libkb/leveldb_cleaner.go | 66 ++++++++++++------------ go/libkb/leveldb_cleaner_test.go | 54 ++++++++++++++++++++ go/libkb/leveldb_test.go | 52 ++++++++----------- 6 files changed, 145 insertions(+), 121 deletions(-) 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 b5b777061d95..a82dcc13f5fd 100644 --- a/go/libkb/leveldb.go +++ b/go/libkb/leveldb.go @@ -122,14 +122,11 @@ type LevelDb struct { // 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 openedDb(). + // through dbOpenerOnce, such as Flush. db atomic.Pointer[leveldb.DB] dbOpenerOnce *sync.Once cleaner *levelDbCleaner - // flushHook, if set, runs after each memtable rotation. Tests only. - flushHook func() - filename string Contextified } @@ -161,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() @@ -198,13 +195,14 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func() error) (err error) return err } - if l.db.Load() == 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) }() // If the file is corrupt, just nuke and act like we didn't find anything @@ -237,7 +235,7 @@ 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 }) } // Flush writes the current memtable to disk and starts an empty journal. An @@ -248,46 +246,23 @@ func (l *LevelDb) ForceOpen() error { // opens fast. No-op if the DB is not currently open; does not trigger a lazy // open. // -// Concurrent calls serialize on goleveldb's own write lock rather than on -// anything of ours: OpenTransaction blocks until any transaction ahead of it -// commits or discards, and by the time it unblocks it rotates whatever -// memtable is current, so a call that lands after another one already -// covers any write made before it arrived. +// 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() (err error) { - return l.flushMemtable() -} - -// openedDb returns the DB without triggering a lazy open, or nil if it isn't -// open. Callers must hold the read lock. -func (l *LevelDb) openedDb() *leveldb.DB { - return l.db.Load() -} - -func (l *LevelDb) flushMemtable() (err error) { defer convertNoSpaceError(&err) l.RLock() defer l.RUnlock() - db := l.openedDb() + db := l.db.Load() if db == nil { return nil } - // Opening a transaction rotates a non-empty memtable and waits until it - // is written to a table, without compacting any tables. The - // transaction itself is not needed. - tr, err := db.OpenTransaction() - if err != nil { - return err - } - tr.Discard() - if l.flushHook != nil { - l.flushHook() - } - return nil + return db.FlushMemdb() } func (l *LevelDb) Stats() (stats string) { - if err := l.doWhileOpenAndNukeIfCorrupted(func() (err error) { - stats, err = l.db.Load().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 { @@ -298,8 +273,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.Load().Stats(&dbStats) + if err := l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) (err error) { + return db.Stats(&dbStats) }); err != nil { return false, false, err } @@ -397,30 +372,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.Load(), 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.Load(), 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.Load(), 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.Load(), l.cleaner, id) + return l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) error { + return levelDbDelete(db, l.cleaner, id) }) } @@ -429,13 +404,10 @@ func (l *LevelDb) OpenTransaction() (LocalDbTransaction, error) { ltr LevelDbTransaction err error ) - l.RLock() - db := l.openedDb() - l.RUnlock() - if db == nil { - return LevelDbTransaction{}, LevelDBOpenClosedError{} - } - if ltr.tr, err = db.OpenTransaction(); err != nil { + if err = l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) (err error) { + ltr.tr, err = db.OpenTransaction() + return err + }); err != nil { return LevelDbTransaction{}, err } ltr.cleaner = l.cleaner @@ -444,10 +416,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.Load().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 c0989313241f..313ba1f7da17 100644 --- a/go/libkb/leveldb_cleaner.go +++ b/go/libkb/leveldb_cleaner.go @@ -69,8 +69,6 @@ type levelDbCleaner struct { cacheMu sync.Mutex // protects the pointer to the cache db *leveldb.DB stopCh chan struct{} - - isShutdown bool } func newLevelDbCleaner(mctx MetaContext, dbName string) *levelDbCleaner { @@ -105,8 +103,11 @@ 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) } func (c *levelDbCleaner) Stop() { @@ -120,19 +121,20 @@ func (c *levelDbCleaner) Stop() { } // start attaches the cleaner to a newly opened db, undoing a previous -// Stop/Shutdown from closing it. +// Stop/Shutdown from closing it. A reopened db is a new key space, so cleaning +// starts over from its first key. 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.lastKey = nil c.cacheMu.Lock() defer c.cacheMu.Unlock() - if c.isShutdown { - if cache, err := lru.New(c.config.CacheCapacity); err == nil { - c.cache = cache - c.isShutdown = false - } - } + c.cache = cache } func (c *levelDbCleaner) log(format string, args ...any) { @@ -144,14 +146,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 { @@ -166,14 +167,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 } @@ -183,11 +181,15 @@ 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 + // A clean stays on the db it started with. If that db is closed, Stop + // closes stopCh and the clean exits; a reopen attaches a new db for later + // cleans without this one ever touching it. + db := c.db key := c.lastKey stopCh := c.stopCh // Sample the app state in the same critical section as running=true, so @@ -206,12 +208,14 @@ func (c *levelDbCleaner) clean(force bool) (err error) { 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 }() - dbSize, err := c.getDbSize() + dbSize, err := c.getDbSize(db) if err != nil { return err } @@ -240,7 +244,7 @@ func (c *levelDbCleaner) clean(force bool) (err error) { } start := c.G().GetClock().Now() - numPurged, key, err = c.cleanBatch(key) + numPurged, key, err = c.cleanBatch(db, stopCh, key) if err != nil { return err } @@ -254,7 +258,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 } @@ -268,7 +272,7 @@ func (c *levelDbCleaner) clean(force bool) (err error) { return nil } -func (c *levelDbCleaner) cleanBatch(startKey []byte) (int, []byte, error) { +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)} @@ -278,18 +282,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() + select { + case <-stopCh: + iter.Release() return 0, nil, errors.New("cleanBatch: cancelled due to shutdown") + default: } - cache := c.cache - c.cacheMu.Unlock() + cache := c.getCache() if _, found := cache.Get(c.cacheKey(key)); !found { cp := make([]byte, len(key)) @@ -312,11 +316,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 index 1d6b38d3298e..8b61a88dac88 100644 --- a/go/libkb/leveldb_cleaner_test.go +++ b/go/libkb/leveldb_cleaner_test.go @@ -215,3 +215,57 @@ func TestCleanerCleansAfterReopen(t *testing.T) { }) } } + +// 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 +} diff --git a/go/libkb/leveldb_test.go b/go/libkb/leveldb_test.go index 11697f93d90e..62f6709d260c 100644 --- a/go/libkb/leveldb_test.go +++ b/go/libkb/leveldb_test.go @@ -71,8 +71,8 @@ func doSomeIO() error { } func levelDbStats(t *testing.T, db *LevelDb) (stats leveldb.DBStats) { - require.NoError(t, db.doWhileOpenAndNukeIfCorrupted(func() error { - return db.db.Load().Stats(&stats) + require.NoError(t, db.doWhileOpenAndNukeIfCorrupted(func(ldb *leveldb.DB) error { + return ldb.Stats(&stats) })) return stats } @@ -211,32 +211,6 @@ func TestLevelDb(t *testing.T) { require.Equal(t, bytes.Repeat([]byte{19}, 100), val) }, }, - { - // A write and a Flush call that a flush's own hook makes reentrantly - // must still be flushed before the outer call returns: the hook runs - // after the transaction is discarded, so goleveldb's write lock is - // already free and the nested call is a plain second flush. - name: "flush-reentrant", testBody: func(t *testing.T) { - tc := SetupTest(t, "LevelDb-flush-reentrant", 0) - defer tc.Cleanup() - db, err := createTempLevelDbForTest(&tc, &td) - require.NoError(t, err) - _, err = testLevelDbPut(db) - require.NoError(t, err) - - rotations := 0 - db.flushHook = func() { - rotations++ - if rotations == 1 { - require.NoError(t, db.db.Load().Put([]byte("kv:late"), []byte{1}, nil)) - require.NoError(t, db.Flush()) - } - } - require.NoError(t, db.Flush()) - require.Equal(t, 2, rotations) - require.Zero(t, levelDbJournalSize(t, db)) - }, - }, { // 8 goroutines call Flush with writes interleaved: every call // returns nil, and every writer's last write is durable and @@ -271,6 +245,24 @@ func TestLevelDb(t *testing.T) { } }, }, + { + // 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) + }, + }, { name: "open-transaction-after-close", testBody: func(t *testing.T) { tc := SetupTest(t, "LevelDb-transaction-closed", 0) @@ -381,7 +373,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): @@ -392,7 +384,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): From 3a576131f6511cabc455b711f7ebd137fc8f1413 Mon Sep 17 00:00:00 2001 From: chrisnojima Date: Tue, 22 Sep 2026 15:51:01 -0400 Subject: [PATCH 3/3] fix(leveldb): stop cleans before close; keep OpenTransaction off our lock Close closed goleveldb before stopping the cleaner, so a running clean could hold an iterator across Close, which goleveldb says is unsafe. Stop now runs first, detaches the cleaner from its db, and waits for a running clean to exit; the clean's between-batch sleep wakes on stop. OpenTransaction waited on goleveldb's write lock while holding our read lock. A Close or Nuke queued behind it blocked every new reader, including a Get from the goroutine holding the current transaction, so nobody could proceed. It now opens the db lazily, then waits outside the lock. Flush goes through the same error handling as other operations, so a full disk starts a forced clean and a corrupt db is nuked. Only Nuke resets the cleaner's position; Close and reopen keep the same data, so cleaning resumes where it stopped. --- go/libkb/leveldb.go | 49 ++++++++++++++++------- go/libkb/leveldb_cleaner.go | 54 +++++++++++++++++++------ go/libkb/leveldb_cleaner_test.go | 68 +++++++++++++++++++++++++++++++- go/libkb/leveldb_test.go | 44 +++++++++++++++++++++ 4 files changed, 186 insertions(+), 29 deletions(-) diff --git a/go/libkb/leveldb.go b/go/libkb/leveldb.go index a82dcc13f5fd..48538622b865 100644 --- a/go/libkb/leveldb.go +++ b/go/libkb/leveldb.go @@ -204,7 +204,13 @@ func (l *LevelDb) doWhileOpenAndNukeIfCorrupted(action func(db *leveldb.DB) erro 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 @@ -249,15 +255,18 @@ func (l *LevelDb) ForceOpen() error { // 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() (err error) { - defer convertNoSpaceError(&err) - l.RLock() - defer l.RUnlock() - db := l.db.Load() - if db == nil { - return nil - } - return db.FlushMemdb() +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) { @@ -298,15 +307,16 @@ func (l *LevelDb) closeLocked() error { var err error if db := l.db.Load(); db != nil { l.G().Log.Debug("Closing LevelDB local cache: %s", l.GetFilename()) + // 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 } @@ -353,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) @@ -404,12 +415,20 @@ func (l *LevelDb) OpenTransaction() (LocalDbTransaction, error) { ltr LevelDbTransaction err error ) - if err = l.doWhileOpenAndNukeIfCorrupted(func(db *leveldb.DB) (err error) { - ltr.tr, err = db.OpenTransaction() - return err + // 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 return ltr, nil } diff --git a/go/libkb/leveldb_cleaner.go b/go/libkb/leveldb_cleaner.go index 313ba1f7da17..6184893c0afb 100644 --- a/go/libkb/leveldb_cleaner.go +++ b/go/libkb/leveldb_cleaner.go @@ -69,6 +69,8 @@ type levelDbCleaner struct { 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 { @@ -110,19 +112,34 @@ func (c *levelDbCleaner) Status() string { 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 } + c.Unlock() + if cleanDone != nil { + <-cleanDone + } +} + +// 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 } // start attaches the cleaner to a newly opened db, undoing a previous -// Stop/Shutdown from closing it. A reopened db is a new key space, so cleaning -// starts over from its first key. +// Stop/Shutdown from closing it. func (c *levelDbCleaner) start(db *leveldb.DB) { cache, err := lru.New(c.config.CacheCapacity) if err != nil { @@ -131,7 +148,6 @@ func (c *levelDbCleaner) start(db *leveldb.DB) { c.Lock() defer c.Unlock() c.db = db - c.lastKey = nil c.cacheMu.Lock() defer c.cacheMu.Unlock() c.cache = cache @@ -186,9 +202,11 @@ func (c *levelDbCleaner) clean(force bool) (err error) { return nil } c.running = true - // A clean stays on the db it started with. If that db is closed, Stop - // closes stopCh and the clean exits; a reopen attaches a new db for later - // cleans without this one ever touching it. + 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 @@ -213,6 +231,7 @@ func (c *levelDbCleaner) clean(force bool) (err error) { } c.lastRun = c.G().GetClock().Now() c.running = false + close(cleanDone) }() dbSize, err := c.getDbSize(db) @@ -245,6 +264,10 @@ func (c *levelDbCleaner) clean(force bool) (err error) { start := c.G().GetClock().Now() 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,12 +289,19 @@ 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 } +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. @@ -290,7 +320,7 @@ func (c *levelDbCleaner) cleanBatch(db *leveldb.DB, stopCh chan struct{}, startK select { case <-stopCh: iter.Release() - return 0, nil, errors.New("cleanBatch: cancelled due to shutdown") + return 0, startKey, errCleanStopped default: } cache := c.getCache() diff --git a/go/libkb/leveldb_cleaner_test.go b/go/libkb/leveldb_cleaner_test.go index 8b61a88dac88..cfbe6a69638c 100644 --- a/go/libkb/leveldb_cleaner_test.go +++ b/go/libkb/leveldb_cleaner_test.go @@ -182,8 +182,8 @@ func TestCleanerRearmsIntoBackgroundActive(t *testing.T) { } // A cleaner still cleans after its db is reopened (Nuke, or Close + -// ForceOpen): start() must reset isShutdown so a reopened cleaner's cache -// isn't stuck discarding everything. +// 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) { @@ -269,3 +269,67 @@ func TestCleanerStatusDuringReopen(t *testing.T) { 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 62f6709d260c..e275ba5528f2 100644 --- a/go/libkb/leveldb_test.go +++ b/go/libkb/leveldb_test.go @@ -263,6 +263,50 @@ func TestLevelDb(t *testing.T) { 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)