Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion go/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
)

Expand Down
4 changes: 2 additions & 2 deletions go/go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
135 changes: 77 additions & 58 deletions go/libkb/leveldb.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"path/filepath"
"strings"
"sync"
"sync/atomic"

"github.com/syndtr/goleveldb/leveldb"
errors "github.com/syndtr/goleveldb/leveldb/errors"
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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()
Expand All @@ -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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

?

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 {
Expand All @@ -179,24 +185,32 @@ 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)
}
})

if err != nil {
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
Expand All @@ -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)
}
Expand All @@ -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 {
Expand All @@ -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
}
Expand All @@ -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
}
Expand Down Expand Up @@ -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)
Expand All @@ -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)
})
}

Expand All @@ -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
Expand All @@ -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 {
Expand Down
Loading