From 3ed4487ae269be42307ad19fa217a751b38177d4 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 19 Nov 2021 12:31:13 -0600 Subject: [PATCH 01/22] scratch space for FB-992: create benchmark for checkpointing Note also the commented-out debug printf in checkpoint, there as a reference. This is interesting because it turns out that MOST of checkpoint writes is not actually writing new pages in most cases. The actual "pages in WAL : pages in map" ratio is typically around 30:1 apparently. This would likely be different in cases where we were updating existing data, though. This is scratch space to prep for an actual work. The final results will likely be different. --- rbf/db.go | 1 + rbf/db_test.go | 86 ++++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 87 insertions(+) diff --git a/rbf/db.go b/rbf/db.go index 99c4f0e2f..65794c9eb 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -200,6 +200,7 @@ func (db *DB) checkpoint() error { return nil } + // fmt.Printf("checkpoint: walPageN %d, PageMap size %d\n", db.walPageN, db.pageMap.size) for i := 0; i < db.walPageN; i++ { page, err := db.readWALPageAt(i) if err != nil { diff --git a/rbf/db_test.go b/rbf/db_test.go index 56170d6eb..614def465 100644 --- a/rbf/db_test.go +++ b/rbf/db_test.go @@ -13,6 +13,7 @@ import ( _ "net/http/pprof" + "github.com/felixge/fgprof" "github.com/molecula/featurebase/v2/rbf" rbfcfg "github.com/molecula/featurebase/v2/rbf/cfg" "golang.org/x/sync/errgroup" @@ -336,6 +337,91 @@ func TestDB_MultiTx(t *testing.T) { } } +// benchmarkOneCheckpoint +func benchmarkOneCheckpoint(b *testing.B) { + cfg := rbfcfg.NewDefaultConfig() + // extremely low to force checkpointing + cfg.MinWALCheckpointSize = rbf.PageSize * 16 + cfg.MaxWALCheckpointSize = rbf.PageSize * 64 + var _ rbfcfg.Config + db := MustOpenDB(b, cfg) + defer MustCloseDB(b, db) + + // Run multiple readers in separate goroutines. + ctx, cancel := context.WithCancel(context.Background()) + g, ctx := errgroup.WithContext(ctx) + for i := 0; i < 4; i++ { + g.Go(func() error { + for { + if ctx.Err() != nil { + return nil // cancelled, return no error + } else if err := func() error { + tx, err := db.Begin(false) + if err != nil { + return err + } + defer tx.Rollback() + + // time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond)))) + + for i := 0; i < rand.Intn(1000); i++ { + v := rand.Intn(1 << 20) + if _, err := tx.Contains("x", uint64(v)); err != nil { + return err + } + } + return nil + }(); err != nil { + return err + } + + // time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond)))) + } + }) + } + + // Continuously set/clear bits while readers are executing. + for i := 0; i < 1000; i++ { + func() { + tx, err := db.Begin(true) + if err != nil { + b.Fatal(err) + } + defer tx.Rollback() + + for j := 0; j < rand.Intn(100); j++ { + v := rand.Intn(1 << 20) + if _, err := tx.Add("x", uint64(v)); err != nil { + b.Fatal(err) + } + + } + + if err := tx.Commit(); err != nil { + b.Fatal(err) + } + }() + } + + // Stop readers & wait. + cancel() + if err := g.Wait(); err != nil { + b.Fatal(err) + } +} + +func BenchmarkDbCheckpoint(b *testing.B) { + out, err := os.Create("cp.out") + if err != nil { + b.Fatalf("creating log file: %v", err) + } + done := fgprof.Start(out, fgprof.FormatPprof) + for i := 0; i < b.N; i++ { + benchmarkOneCheckpoint(b) + } + done() +} + // better diagnosis of deadlocks/hung situations versus just really slow "Quick" tests. func TestMain(m *testing.M) { l, err := net.Listen("tcp", ":0") From 5c889c72bdab24f960d69470adc99301f558fa96 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 19 Nov 2021 13:04:45 -0600 Subject: [PATCH 02/22] make test hit the lock harder Discovered test was running slightly strange and spending an unreasonable amount of time on rand.Intn(), possibly because we weren't caching the value used as the loop condition. Tweaked that, also made the pool a bit different. Now it takes ~50 seconds for benchtime 100x, and produces a profile with a TON of time spent waiting on sleeps (expected) and the condition variable for waiting on checkpoints (the thing we want to measure, really). --- rbf/db_test.go | 54 ++++++++++++++++++++++++++++++++++++++------------ 1 file changed, 41 insertions(+), 13 deletions(-) diff --git a/rbf/db_test.go b/rbf/db_test.go index 614def465..5c8504b1b 100644 --- a/rbf/db_test.go +++ b/rbf/db_test.go @@ -337,8 +337,11 @@ func TestDB_MultiTx(t *testing.T) { } } +// premake pool of random values +const randPool = (1 << 18) + // benchmarkOneCheckpoint -func benchmarkOneCheckpoint(b *testing.B) { +func benchmarkOneCheckpoint(b *testing.B, randInts []int) { cfg := rbfcfg.NewDefaultConfig() // extremely low to force checkpointing cfg.MinWALCheckpointSize = rbf.PageSize * 16 @@ -350,7 +353,8 @@ func benchmarkOneCheckpoint(b *testing.B) { // Run multiple readers in separate goroutines. ctx, cancel := context.WithCancel(context.Background()) g, ctx := errgroup.WithContext(ctx) - for i := 0; i < 4; i++ { + for i := 0; i < 8; i++ { + i := i g.Go(func() error { for { if ctx.Err() != nil { @@ -362,10 +366,11 @@ func benchmarkOneCheckpoint(b *testing.B) { } defer tx.Rollback() - // time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond)))) + time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond)))) - for i := 0; i < rand.Intn(1000); i++ { - v := rand.Intn(1 << 20) + times := rand.Intn(1000) + 1 + for j := 0; j < times; j++ { + v := randInts[((i<<10)+j)%(randPool-1)] if _, err := tx.Contains("x", uint64(v)); err != nil { return err } @@ -374,13 +379,13 @@ func benchmarkOneCheckpoint(b *testing.B) { }(); err != nil { return err } - // time.Sleep(time.Duration(rand.Intn(int(3 * time.Millisecond)))) } }) } // Continuously set/clear bits while readers are executing. + next := 0 for i := 0; i < 1000; i++ { func() { tx, err := db.Begin(true) @@ -389,14 +394,22 @@ func benchmarkOneCheckpoint(b *testing.B) { } defer tx.Rollback() - for j := 0; j < rand.Intn(100); j++ { - v := rand.Intn(1 << 20) - if _, err := tx.Add("x", uint64(v)); err != nil { - b.Fatal(err) + times := rand.Intn(100) + for j := 0; j < times; j++ { + v := randInts[next] + next = (next + 1) % (randPool - 1) + if j&7 == 0 { + // some removes but they're less frequent + if _, err := tx.Remove("x", uint64(v)); err != nil { + b.Fatal(err) + } + } else { + if _, err := tx.Add("x", uint64(v)); err != nil { + b.Fatal(err) + } } } - if err := tx.Commit(); err != nil { b.Fatal(err) } @@ -416,9 +429,24 @@ func BenchmarkDbCheckpoint(b *testing.B) { b.Fatalf("creating log file: %v", err) } done := fgprof.Start(out, fgprof.FormatPprof) - for i := 0; i < b.N; i++ { - benchmarkOneCheckpoint(b) + b.StopTimer() + // premake these because otherwise it's >5% of CPU in the reads + randInts := make([]int, randPool) + for i := range randInts { + v1, v2 := rand.Intn(1<<24), rand.Intn(1<<24) + // minimum gives us a skewed distribution which makes lower values more + // likely than higher values, so we get a mix of container types + if v1 < v2 { + randInts[i] = v1 + } else { + randInts[i] = v2 + } } + b.StartTimer() + for i := 0; i < b.N; i++ { + benchmarkOneCheckpoint(b, randInts) + } + b.StopTimer() done() } From a631e25dc517e5473ea70863bae8633e0232ce00 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 19 Nov 2021 13:14:42 -0600 Subject: [PATCH 03/22] refactor: removeTx responsible for getting/releasing its own lock We change nothing substantive here, except that there's a window between when a write transaction updates the root pages and when it removes itself from the db tx list and possibly causes a checkpoint where it's not holding the db lock. The issue here is that we want to be able to *keep* the lock but still return, so no one else can start transactions, but the specific Rollback or Commit that removed the last outstanding transaction doesn't block forever. This will, later, allow us to exercise finer-grained control over when we allow transactions. This is a separate commit so we can run the test suite against it, and verify that this part in particular didn't break anything. --- rbf/db.go | 7 ++++++- rbf/tx.go | 14 +++++++++----- 2 files changed, 15 insertions(+), 6 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 65794c9eb..7f27ca418 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -503,8 +503,13 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) { return tx, nil } -// removeTx removes an active transaction from the database. +// removeTx removes an active transaction from the database. it obtains +// the db lock, and currently drops it, but will later possibly be leaving +// it retained by an asynchronous op that wants to happen before we start +// running new tx. func (db *DB) removeTx(tx *Tx) error { + db.mu.Lock() + defer db.mu.Unlock() // Release writer lock if tx is writable. if tx.writable { tx.db.rwmu.Unlock() diff --git a/rbf/tx.go b/rbf/tx.go index 38fed48f6..2cf7420db 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -109,20 +109,24 @@ func (tx *Tx) Commit() error { // future plan: after checkpoint is moved to background // or not every removeTx, then we can move the // tx.db.rootRecords = tx.rootRecords into removeTx(). - + // + // ... or maybe not: let's do that part here, and then removeTx + // may or may not start a checkpoint, possibly asynchronously. + // // avoid race detector firing on a write race here - // vs the read of rootRecords at db.Begin() + // vs the read of rootRecords at db.Begin(), then release + // the lock, because we need removeTx to grab the lock to + // work, but if it wants to checkpoint, it wants to be able to return + // to us here and still be holding the lock. tx.db.mu.Lock() - defer tx.db.mu.Unlock() tx.db.rootRecords = tx.rootRecords tx.db.pageMap = tx.pageMap tx.db.walPageN = tx.walPageN + tx.db.mu.Unlock() return tx.db.removeTx(tx) } // Disconnect transaction from DB. - tx.db.mu.Lock() - defer tx.db.mu.Unlock() return tx.db.removeTx(tx) } From b8f59d922c72e1e1b8d4b1ee36bbd89e6e0d9096 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 19 Nov 2021 14:37:59 -0600 Subject: [PATCH 04/22] checkpoint rework/refactoring: logger, async-ish checkpoint Trying to make the checkpoint be asynchronous-at-all, and also allowing it to log. --- rbf/cfg/cfg.go | 6 ++ rbf/db.go | 162 ++++++++++++++++++++++++++++++++++++++---------- rbf/rbf_test.go | 8 +++ 3 files changed, 144 insertions(+), 32 deletions(-) diff --git a/rbf/cfg/cfg.go b/rbf/cfg/cfg.go index cc2cf7a8a..671c6fe43 100644 --- a/rbf/cfg/cfg.go +++ b/rbf/cfg/cfg.go @@ -2,6 +2,7 @@ package cfg import ( + "github.com/molecula/featurebase/v2/logger" "github.com/spf13/pflag" ) @@ -35,6 +36,11 @@ type Config struct { // CursorCacheSize is the number of copies of Cursor{} to keep in our // readyCursorCh arena to avoid GC pressure. CursorCacheSize int64 `toml:"cursor-cache-size"` + + // Logger specifies a logger for asynchronous errors, such as + // background checkpoints. It cannot be set from toml. The default is + // to use stderr. + Logger logger.Logger `toml:"-"` } func NewDefaultConfig() *Config { diff --git a/rbf/db.go b/rbf/db.go index 7f27ca418..830e31f5a 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -11,6 +11,7 @@ import ( "syscall" "github.com/benbjohnson/immutable" + "github.com/molecula/featurebase/v2/logger" rbfcfg "github.com/molecula/featurebase/v2/rbf/cfg" "github.com/molecula/featurebase/v2/syswrap" ) @@ -38,6 +39,7 @@ type DB struct { pageMap *PageMap // pgno-to-WALID mapping txs map[*Tx]struct{} // active transactions opened bool // true if open + logger logger.Logger // for diagnostics from async things wal []byte // wal mmap walFile *os.File // wal file descriptor @@ -62,6 +64,11 @@ func NewDB(path string, cfg *rbfcfg.Config) *DB { txs: make(map[*Tx]struct{}), pageMap: NewPageMap(), Path: path, + logger: cfg.Logger, + } + if db.logger == nil { + // default to writing to stdout if not told otherwise + db.logger = logger.NewStandardLogger(os.Stderr) } db.haltCond = sync.NewCond(&db.mu) @@ -134,7 +141,7 @@ func (db *DB) Open() (err error) { if err := db.openWAL(); err != nil { return fmt.Errorf("wal open: %w", err) } else if err := db.checkpoint(); err != nil { - return fmt.Errorf("checkpoint: %w", err) + return fmt.Errorf("startup checkpoint: %w", err) } return nil @@ -158,10 +165,12 @@ func (db *DB) openWAL() (err error) { // Determine the number of whole pages in the WAL. var pageN int + var fileSize int64 if fi, err := db.walFile.Stat(); err != nil { return fmt.Errorf("wal stat: %w", err) } else { - pageN = int(fi.Size() / PageSize) + fileSize = fi.Size() + pageN = int(fileSize / PageSize) } // Read backwards through the WAL to find the last valid meta page. @@ -169,14 +178,45 @@ func (db *DB) openWAL() (err error) { if page, err := db.readWALPageAt(pageN - 1); err != nil { return err } else if IsMetaPage(page) { + // We now face a challenge. Probably this is a meta page. + // But consider a sequence of pages written which gets + // interrupted right before the meta page is written. + // If the last page is a bitmap page, it could LOOK LIKE a meta + // page. So we have to check the page before it. If that page + // is a bitmap header, then actually this is a bitmap page, right? + // If that page doesn't exist, of course, we're fine, except + // for the philosophical question of why we wrote a meta page + // when no pages had changed. + if pageN > 1 { + if page, err = db.readWALPageAt(pageN - 2); err != nil { + return err + } + if IsBitmapHeader(page) { + // But wait! + // What if this *is* a meta page, and the page before it is + // actually a *bitmap page* that looks like a bitmap header? And + // so on. + // + // Rather than try to resolve this, in this insanely unlikely + // situation, we read from the beginning which allows us to + // always know what we're seeing, because every bitmap page + // comes *after* a bitmap header page, and thus, we know when + // we might be seeing one. + pageN, err = db.methodicalWALPageN(pageN) + if err != nil { + return err + } + } + } break } } - - // Truncate WAL to the last valid meta page. - if err := db.walFile.Truncate(int64(pageN * PageSize)); err != nil { - return fmt.Errorf("wal truncate: %w", err) - } else if _, err := db.walFile.Seek(int64(pageN*PageSize), io.SeekStart); err != nil { + if fileSize != int64(pageN*PageSize) { + if err := db.walFile.Truncate(int64(pageN * PageSize)); err != nil { + return fmt.Errorf("wal truncate: %w", err) + } + } + if _, err := db.walFile.Seek(int64(pageN*PageSize), io.SeekStart); err != nil { return fmt.Errorf("wal seek: %w", err) } db.walPageN = pageN @@ -184,6 +224,27 @@ func (db *DB) openWAL() (err error) { return nil } +// methodicalWALPageN tries to determine the last meta page in a very reliable +// but slow way. This handles the theoretical but hard to imagine creating +// edge case where we have a bitmap page which happens to look like a meta +// page, and the write got interrupted before the meta page got written. +func (db *DB) methodicalWALPageN(pageN int) (lastMeta int, err error) { + for i := 0; i < pageN; i++ { + var page []byte + if page, err = db.readWALPageAt(i); err != nil { + return -1, err + } + switch { + case IsMetaPage(page): + lastMeta = i + case IsBitmapHeader(page): + // skip the bitmap page, which we can't usefully evaluate + i++ + } + } + return lastMeta, nil +} + // checkpoint moves all WAL pages to the main DB file. // Must be called by a write transaction while under db.mu lock. func (db *DB) checkpoint() error { @@ -199,28 +260,53 @@ func (db *DB) checkpoint() error { if db.walPageN == 0 { return nil } + // We might have either a *PageMap or just the file. If we have the file, + // building the PageMap is fairly expensive because it's fancy and immutable. + // If we have the PageMap *or* some other map, that's two different things + // to iterate. If we have the PageMap, building a map from it is relatively + // cheap, so we'll do it that way. + pages := make(map[uint32]int) + if db.pageMap.size == 0 { + // you'd think we're done, but actually this PROBABLY means that + // this is initial startup, and we haven't read the file yet. We scan + // the file for pages, because it turns out most of them probably + // got overwritten. + for i := 0; i < db.walPageN; i++ { + page, err := db.readWALPageAt(i) + if err != nil { + return err + } + + // Determine page number. Meta pages are always on zero & bitmap + // headers specify the page number of the next page in the WAL. + // All other pages have their page number in the page data. + var pgno uint32 + if IsBitmapHeader(page) { + pgno = readPageNo(page) + if page, err = db.readWALPageAt(i + 1); err != nil { + return err + } + i++ // bitmaps in WAL are two pages + } else if !IsMetaPage(page) { + pgno = readPageNo(page) + } + // record where in the file we have this page + pages[pgno] = i + } + } else { + itr := db.pageMap.Iterator() + itr.First() + for k, v, ok := itr.Next(); ok; k, v, ok = itr.Next() { + pages[k] = int(v) + } + } // fmt.Printf("checkpoint: walPageN %d, PageMap size %d\n", db.walPageN, db.pageMap.size) - for i := 0; i < db.walPageN; i++ { - page, err := db.readWALPageAt(i) + for pgno, walID := range pages { + page, err := db.readWALPageAt(walID) if err != nil { return err } - - // Determine page number. Meta pages are always on zero & bitmap - // headers specify the page number of the next page in the WAL. - // All other pages have their page number in the page data. - var pgno uint32 - if IsBitmapHeader(page) { - pgno = readPageNo(page) - if page, err = db.readWALPageAt(i + 1); err != nil { - return err - } - i++ // bitmaps in WAL are two pages - } else if !IsMetaPage(page) { - pgno = readPageNo(page) - } - // Write data to the data file. if err := db.writeDBPage(pgno, page); err != nil { return err @@ -509,23 +595,35 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) { // running new tx. func (db *DB) removeTx(tx *Tx) error { db.mu.Lock() - defer db.mu.Unlock() - // Release writer lock if tx is writable. + // release the write lock. we have to do this for now. some day we won't, + // and will want to hold it, but right now we can't be sure we can get it. if tx.writable { tx.db.rwmu.Unlock() } - + // remove ourselves from the list of transactions the db is keeping. delete(tx.db.txs, tx) // Disassociate from db. tx.db = nil - // Write pages from WAL to DB. - // TODO(bbj): Move this to an async goroutine. + // Write pages from WAL to DB. As of this instant, we are the ONLY + // transaction, which means that no transaction has an older version + // of the PageMap than we do, and if we're a write, Commit() already + // updated the page map to our page map. So, if we *can* checkpoint, + // the checkpoint gets spawned asynchronously. We could block writes, + // except doing so will deadlock in a weird way. if len(db.txs) == 0 && db.walSize() > db.cfg.MinWALCheckpointSize { - if err := db.checkpoint(); err != nil { - return fmt.Errorf("checkpoint: %w", err) - } + // We are doing this function with the db lock held, which means + // we're *not* releasing the db lock, even though we're returning. + // This is a weird special case, and probably a bad idea. + go func() { + defer db.mu.Unlock() + if err := db.checkpoint(); err != nil { + db.logger.Errorf("async checkpoint: %w", err) + } + }() + } else { + defer db.mu.Unlock() } return nil diff --git a/rbf/rbf_test.go b/rbf/rbf_test.go index 48f16335a..be3461e37 100644 --- a/rbf/rbf_test.go +++ b/rbf/rbf_test.go @@ -11,6 +11,7 @@ import ( "sort" "testing" + "github.com/molecula/featurebase/v2/logger" "github.com/molecula/featurebase/v2/rbf" rbfcfg "github.com/molecula/featurebase/v2/rbf/cfg" "github.com/molecula/featurebase/v2/testhook" @@ -65,6 +66,13 @@ func NewDB(tb testing.TB, cfg ...*rbfcfg.Config) *rbf.DB { // MustOpenDB returns a db opened on a temporary file. On error, fail test. func MustOpenDB(tb testing.TB, cfg ...*rbfcfg.Config) *rbf.DB { tb.Helper() + if len(cfg) == 0 || cfg[0] == nil { + newconf := rbfcfg.NewDefaultConfig() + newconf.Logger = logger.NewLogfLogger(tb) + cfg = []*rbfcfg.Config{newconf} + } else if cfg[0].Logger == nil { + cfg[0].Logger = logger.NewLogfLogger(tb) + } db := NewDB(tb, cfg...) if err := db.Open(); err != nil { tb.Fatal(err) From 806669fa0f5c738a5e1d70ac749c67453540bf1c Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 19 Nov 2021 15:27:16 -0600 Subject: [PATCH 05/22] make db able to fail out if it can't checkpoint, fix silly wrong-units error PageMap uses "WALID", which is a WAL page ID relative to the "base" ID of the WAL, rather than the wal page count you'd get just reading the file. So everything it reports has a fixed offset at any given time. I think this may be left over from a point where there were partial checkpoints. Anyway, the net outcome is that each new transaction was getting different page IDs, but the actual WAL pages did not always reflect that. Each checkpoint increases the offset. This might imply that we can start having problems after 4 billion pages written even if most of them were redundant? Anyway, with that fixed, this seems to work. I think. --- rbf/db.go | 76 ++++++++++++++++++++++++++++++++++++------------------- 1 file changed, 50 insertions(+), 26 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 830e31f5a..40c98e22a 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -49,6 +49,8 @@ type DB struct { rwmu sync.Mutex // mutex for restricting single writer haltCond *sync.Cond // condition for resuming txs after checkpoint + isDead error // this database died in an unrecoverable way, error out opens + // Path represents the path to the database file. Path string } @@ -247,7 +249,7 @@ func (db *DB) methodicalWALPageN(pageN int) (lastMeta int, err error) { // checkpoint moves all WAL pages to the main DB file. // Must be called by a write transaction while under db.mu lock. -func (db *DB) checkpoint() error { +func (db *DB) checkpoint() (err error) { if !db.opened { return nil } else if len(db.txs) > 0 { @@ -260,21 +262,31 @@ func (db *DB) checkpoint() error { if db.walPageN == 0 { return nil } + // wake up things waiting on haltCond when we're done, even if we fail. + // Otherwise, we deadlock with them all stuck waiting on that forever. + defer func() { + if err != nil && db.isDead == nil { + db.isDead = err + } + db.haltCond.Broadcast() + }() + var page []byte // We might have either a *PageMap or just the file. If we have the file, // building the PageMap is fairly expensive because it's fancy and immutable. // If we have the PageMap *or* some other map, that's two different things // to iterate. If we have the PageMap, building a map from it is relatively // cheap, so we'll do it that way. pages := make(map[uint32]int) + if db.pageMap.size == 0 { // you'd think we're done, but actually this PROBABLY means that // this is initial startup, and we haven't read the file yet. We scan // the file for pages, because it turns out most of them probably // got overwritten. for i := 0; i < db.walPageN; i++ { - page, err := db.readWALPageAt(i) + page, err = db.readWALPageAt(i) if err != nil { - return err + return fmt.Errorf("reading WAL page %d: %w", i, err) } // Determine page number. Meta pages are always on zero & bitmap @@ -283,8 +295,12 @@ func (db *DB) checkpoint() error { var pgno uint32 if IsBitmapHeader(page) { pgno = readPageNo(page) - if page, err = db.readWALPageAt(i + 1); err != nil { - return err + if i+1 < db.walPageN { + if page, err = db.readWALPageAt(i + 1); err != nil { + return err + } + } else { + return fmt.Errorf("last page of WAL file (%d) is bitmap header", i) } i++ // bitmaps in WAL are two pages } else if !IsMetaPage(page) { @@ -294,41 +310,37 @@ func (db *DB) checkpoint() error { pages[pgno] = i } } else { + walBase := db.baseWALID() itr := db.pageMap.Iterator() itr.First() for k, v, ok := itr.Next(); ok; k, v, ok = itr.Next() { - pages[k] = int(v) + pages[k] = int(v - walBase - 1) } } // fmt.Printf("checkpoint: walPageN %d, PageMap size %d\n", db.walPageN, db.pageMap.size) for pgno, walID := range pages { - page, err := db.readWALPageAt(walID) + page, err = db.readWALPageAt(walID) if err != nil { - return err + return fmt.Errorf("reading page %d [page number %d]: %v", walID, pgno, err) } // Write data to the data file. - if err := db.writeDBPage(pgno, page); err != nil { - return err + if err = db.writeDBPage(pgno, page); err != nil { + return fmt.Errorf("writing page %d: %v", pgno, err) } } - // Ensure database file is synced and then truncate the WAL file. - if err := db.fsync(db.file); err != nil { + if err = db.fsync(db.file); err != nil { return fmt.Errorf("db file sync: %w", err) - } else if err := db.walFile.Truncate(0); err != nil { + } else if err = db.walFile.Truncate(0); err != nil { return fmt.Errorf("truncate wal file: %w", err) - } else if err := db.fsync(db.walFile); err != nil { + } else if err = db.fsync(db.walFile); err != nil { return fmt.Errorf("wal file sync: %w", err) - } else if _, err := db.walFile.Seek(0, io.SeekStart); err != nil { + } else if _, err = db.walFile.Seek(0, io.SeekStart); err != nil { return fmt.Errorf("seek wal file: %w", err) } db.walPageN = 0 db.pageMap = NewPageMap() - - // Notify halted transactions that the WAL has been checkpointed. - db.haltCond.Broadcast() - return nil } @@ -537,9 +549,21 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) { cleanup() return nil, ErrClosed } + if db.isDead != nil { + err := db.isDead + cleanup() + db.mu.Unlock() + return nil, err + } // Wait for WAL size to be below threshold. for int64(db.walPageN*PageSize) > db.cfg.MaxWALCheckpointSize { + if db.isDead != nil { + err := db.isDead + cleanup() + db.mu.Unlock() + return nil, err + } db.haltCond.Wait() } @@ -616,14 +640,14 @@ func (db *DB) removeTx(tx *Tx) error { // We are doing this function with the db lock held, which means // we're *not* releasing the db lock, even though we're returning. // This is a weird special case, and probably a bad idea. - go func() { - defer db.mu.Unlock() - if err := db.checkpoint(); err != nil { - db.logger.Errorf("async checkpoint: %w", err) - } - }() - } else { + // go func() { defer db.mu.Unlock() + if err := db.checkpoint(); err != nil { + db.logger.Errorf("async checkpoint: %v", err) + } + // }() + } else { + db.mu.Unlock() } return nil From 6d68e719338bb825ec1ff537b35b9538ec1581b4 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 19 Nov 2021 15:32:58 -0600 Subject: [PATCH 06/22] make the checkpoint async --- rbf/db.go | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 40c98e22a..9d4769ff3 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -640,12 +640,12 @@ func (db *DB) removeTx(tx *Tx) error { // We are doing this function with the db lock held, which means // we're *not* releasing the db lock, even though we're returning. // This is a weird special case, and probably a bad idea. - // go func() { - defer db.mu.Unlock() - if err := db.checkpoint(); err != nil { - db.logger.Errorf("async checkpoint: %v", err) - } - // }() + go func() { + defer db.mu.Unlock() + if err := db.checkpoint(); err != nil { + db.logger.Errorf("async checkpoint: %v", err) + } + }() } else { db.mu.Unlock() } From c3c02eabb0587a37a58fcf7e99d3b3016a6b8885 Mon Sep 17 00:00:00 2001 From: Seebs Date: Mon, 22 Nov 2021 14:02:53 -0600 Subject: [PATCH 07/22] almost but not quite support async checkpoint This gets us to being able to run reads during a checkpoint, but now we have to wait for new reads to end before we can release the write lock, etc. This is actually slightly slower, but if we could get ONE more step, we could allow new writes during that phase, to a different WAL, if we had a different WAL to write to. --- rbf/db.go | 157 +++++++++++++++++++++++++++++++++++++++++++----------- 1 file changed, 126 insertions(+), 31 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 9d4769ff3..3db8889ca 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -28,6 +28,17 @@ var cursorSyncPool = &sync.Pool{ }, } +// txWaiter is a representation of "i need to wait for txs to complete". +// it is created with a function, and will run that function, with the db +// lock held, at some point after every Tx that was open when it was created +// has closed. WARNING: A txWaiter may hold db.rwmu. +type txWaiter struct { + mu sync.Mutex + cond *sync.Cond + waitingOn map[*Tx]struct{} + callback func() +} + // DB options like MaxSize, FsyncEnabled, DoAllocZero // can be set before calling DB.Open(). type DB struct { @@ -49,6 +60,8 @@ type DB struct { rwmu sync.Mutex // mutex for restricting single writer haltCond *sync.Cond // condition for resuming txs after checkpoint + txWaiters []*txWaiter // things waiting for Txs to close + isDead error // this database died in an unrecoverable way, error out opens // Path represents the path to the database file. @@ -142,8 +155,12 @@ func (db *DB) Open() (err error) { // Open write-ahead log & checkpoint to the end since no transactions are open. if err := db.openWAL(); err != nil { return fmt.Errorf("wal open: %w", err) - } else if err := db.checkpoint(); err != nil { - return fmt.Errorf("startup checkpoint: %w", err) + } else { + // checkpoint wants to hold the rwmu lock. + db.rwmu.Lock() + if err := db.checkpoint(); err != nil { + return fmt.Errorf("startup checkpoint: %w", err) + } } return nil @@ -247,9 +264,18 @@ func (db *DB) methodicalWALPageN(pageN int) (lastMeta int, err error) { return lastMeta, nil } -// checkpoint moves all WAL pages to the main DB file. -// Must be called by a write transaction while under db.mu lock. +// checkpoint moves all WAL pages to the main DB file. Must be called +// while holding both db.mu and db.rwmu. Should release db.rwmu, but not +// db.mu. func (db *DB) checkpoint() (err error) { + // if we don't spin off a possible async waiter, we should release the + // write lock, if we do, that will release it. + releaseLock := true + defer func() { + if releaseLock { + db.rwmu.Unlock() + } + }() if !db.opened { return nil } else if len(db.txs) > 0 { @@ -332,15 +358,26 @@ func (db *DB) checkpoint() (err error) { // Ensure database file is synced and then truncate the WAL file. if err = db.fsync(db.file); err != nil { return fmt.Errorf("db file sync: %w", err) - } else if err = db.walFile.Truncate(0); err != nil { - return fmt.Errorf("truncate wal file: %w", err) - } else if err = db.fsync(db.walFile); err != nil { - return fmt.Errorf("wal file sync: %w", err) - } else if _, err = db.walFile.Seek(0, io.SeekStart); err != nil { - return fmt.Errorf("seek wal file: %w", err) } - db.walPageN = 0 - db.pageMap = NewPageMap() + // now we've updated the file. There are existing transactions that are still + // using the WAL, though. So we wait for them to terminate before we unlock + // the rwmu and update the metadata about the WAL. + releaseLock = false + // fmt.Printf("checkpoint mostly done, waiting for Tx cleanup...\n") + db.afterCurrentTx(func() { + // fmt.Printf("truncating WAL\n") + defer db.rwmu.Unlock() + if err = db.walFile.Truncate(0); err != nil { + db.logger.Errorf("truncate wal file: %w", err) + } else if err = db.fsync(db.walFile); err != nil { + db.logger.Errorf("wal file sync: %w", err) + } else if _, err = db.walFile.Seek(0, io.SeekStart); err != nil { + db.logger.Errorf("seek wal file: %w", err) + } + db.walPageN = 0 + db.pageMap = NewPageMap() + // fmt.Printf("checkpoint actually done\n") + }) return nil } @@ -613,43 +650,101 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) { return tx, nil } +// afterCurrentTx produces runs the provided callback, with the db lock +// held, after all current Tx terminate. It should be called with the db +// lock held. +func (db *DB) afterCurrentTx(callback func()) { + if len(db.txs) == 0 { + callback() + return + } + txw := &txWaiter{} + txw.cond = sync.NewCond(&txw.mu) + txw.callback = callback + txw.waitingOn = make(map[*Tx]struct{}, len(db.txs)) + for k := range db.txs { + txw.waitingOn[k] = struct{}{} + } + db.txWaiters = append(db.txWaiters, txw) + txw.mu.Lock() + go func() { + for len(txw.waitingOn) > 0 { + // fmt.Printf("afterCurrentTx: %d left\n", len(txw.waitingOn)) + txw.cond.Wait() + } + // fmt.Printf("afterCurrentTx: locking db\n") + db.mu.Lock() + defer db.mu.Unlock() + // remove us from the db's list + for i, v := range db.txWaiters { + if v == txw { + // remove us from the list + copy(db.txWaiters[i:], db.txWaiters[i+1:]) + db.txWaiters = db.txWaiters[:len(db.txWaiters)-1] + break + } + } + // fmt.Printf("afterCurrentTx: running callback\n") + txw.callback() + }() + return +} + // removeTx removes an active transaction from the database. it obtains // the db lock, and currently drops it, but will later possibly be leaving // it retained by an asynchronous op that wants to happen before we start // running new tx. func (db *DB) removeTx(tx *Tx) error { db.mu.Lock() - // release the write lock. we have to do this for now. some day we won't, - // and will want to hold it, but right now we can't be sure we can get it. + defer db.mu.Unlock() + // We might want to trigger a checkpoint. Only for writable + // transactions, and only when either there's nothing else open or we + // really need to. + checkpoint := false if tx.writable { - tx.db.rwmu.Unlock() + walSize := db.walSize() + if walSize > db.cfg.MinWALCheckpointSize { + // Might be a good time for a checkpoint. We'll do a checkpoint + // if we're the only transaction, or if we have to. + if len(db.txs) == 1 || walSize > db.cfg.MaxWALCheckpointSize { + checkpoint = true + } + } + // During checkpointing, we'll be preventing writes, but allowing reads. + if !checkpoint { + tx.db.rwmu.Unlock() + } } // remove ourselves from the list of transactions the db is keeping. delete(tx.db.txs, tx) + for _, txw := range tx.db.txWaiters { + delete(txw.waitingOn, tx) + // let it know we're done. we've still got db.mu.lock, so it won't + // happen just yet, but it'll be able to continue. + if len(txw.waitingOn) == 0 { + txw.cond.Broadcast() + } + } // Disassociate from db. tx.db = nil - // Write pages from WAL to DB. As of this instant, we are the ONLY - // transaction, which means that no transaction has an older version - // of the PageMap than we do, and if we're a write, Commit() already - // updated the page map to our page map. So, if we *can* checkpoint, - // the checkpoint gets spawned asynchronously. We could block writes, - // except doing so will deadlock in a weird way. - if len(db.txs) == 0 && db.walSize() > db.cfg.MinWALCheckpointSize { - // We are doing this function with the db lock held, which means - // we're *not* releasing the db lock, even though we're returning. - // This is a weird special case, and probably a bad idea. - go func() { - defer db.mu.Unlock() + if checkpoint { + // We need to run a checkpoint. This can be semi-asynchronous. + // It needs to wait until every existing transaction has finished, + // because every existing transaction could want to look up pages + // which are in the database before our operations, but which should + // now be in the WAL. We want them to use the WAL instead. + // fmt.Printf("possibly-async checkpoint...\n") + db.afterCurrentTx(func() { + // We still hold db.rwmu here. checkpoint unlocks it when it's + // ready. + // fmt.Printf("checkpoint starting\n") if err := db.checkpoint(); err != nil { db.logger.Errorf("async checkpoint: %v", err) } - }() - } else { - db.mu.Unlock() + }) } - return nil } From 4279e2cb2d8010639de375021a99e40a64538e9f Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Thu, 16 Dec 2021 09:17:13 -0700 Subject: [PATCH 08/22] rebase fixes --- rbf/cursor_test.go | 4 ++-- rbf/db.go | 2 -- 2 files changed, 2 insertions(+), 4 deletions(-) diff --git a/rbf/cursor_test.go b/rbf/cursor_test.go index 4f7464bd1..940798c04 100644 --- a/rbf/cursor_test.go +++ b/rbf/cursor_test.go @@ -973,8 +973,8 @@ func TestCursor_SplitBranchCells(t *testing.T) { } // c, _ := tx.Cursor("x") //added just for dot code coverage - c.Dump("ignore for coverage") - + c.Dump("test.dump") + os.Remove("test.dump") } func TestCursor_RemoveCells(t *testing.T) { diff --git a/rbf/db.go b/rbf/db.go index 3db8889ca..21596d144 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -695,8 +695,6 @@ func (db *DB) afterCurrentTx(callback func()) { // it retained by an asynchronous op that wants to happen before we start // running new tx. func (db *DB) removeTx(tx *Tx) error { - db.mu.Lock() - defer db.mu.Unlock() // We might want to trigger a checkpoint. Only for writable // transactions, and only when either there's nothing else open or we // really need to. From 29f5f6d7c2dc84aebfae9d4a67bf9f7e65bcb3c2 Mon Sep 17 00:00:00 2001 From: Seebs Date: Wed, 15 Dec 2021 12:09:40 -0600 Subject: [PATCH 09/22] copy things rows after getting them and before their finishers during writes When a qcx is a write, every Tx under it closes immediately, thus invalidating all returned data. Thus, if you do a Not() inside a Store(), you're doing a difference on an existence row and some other row call... and both of those rows were run, individually, as separate transactions that got invalidated the moment they were fetched. Oops. --- executor.go | 22 ++++++++++++++++++++-- 1 file changed, 20 insertions(+), 2 deletions(-) diff --git a/executor.go b/executor.go index 226759b55..8b0409eb7 100644 --- a/executor.go +++ b/executor.go @@ -4493,7 +4493,11 @@ func (e *executor) executeRowShard(ctx context.Context, qcx *Qcx, index string, return nil, err } defer finisher(&err0) - return frag.row(tx, rowID) + row, err := frag.row(tx, rowID) + if qcx.write && err == nil { + row = row.Clone() + } + return row, err } // If no quantum exists then return an empty bitmap. @@ -4532,15 +4536,21 @@ func (e *executor) executeRowShard(ctx context.Context, qcx *Qcx, index string, if len(rows) == 0 { return &Row{}, nil } else if len(rows) == 1 { + if qcx.write { + return rows[0].Clone(), nil + } return rows[0], nil } row := rows[0].Union(rows[1:]...) + if qcx.write { + row = row.Clone() + } return row, nil } // executeRowBSIGroupShard executes a range(bsiGroup) call for a local shard. -func (e *executor) executeRowBSIGroupShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (_ *Row, err0 error) { +func (e *executor) executeRowBSIGroupShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (cloneable *Row, err0 error) { span, _ := tracing.StartSpanFromContext(ctx, "Executor.executeRowBSIGroupShard") defer span.Finish() @@ -4572,6 +4582,11 @@ func (e *executor) executeRowBSIGroupShard(ctx context.Context, qcx *Qcx, index return nil, err } defer finisher(&err0) + defer func() { + if qcx.write && cloneable != nil { + cloneable = cloneable.Clone() + } + }() // EQ null _exists - frag.NotNull() // NEQ null frag.NotNull() @@ -4822,6 +4837,9 @@ func (e *executor) executeNotShard(ctx context.Context, qcx *Qcx, index string, if existenceRow, err = existenceFrag.row(tx, 0); err != nil { return nil, err } + if qcx.write { + existenceRow = existenceRow.Clone() + } } // the finishers returned by a write tx, which we might be in if there's // a higher-level write in this call OR ANY OTHER CALL, are safe to From 5764d98f6d0198b2b61b32f5e336de5d59b110d2 Mon Sep 17 00:00:00 2001 From: Seebs Date: Thu, 16 Dec 2021 11:06:17 -0600 Subject: [PATCH 10/22] test fixes and order of operations on changing db.PageMap We need to update db.PageMap after we write the db, but before we truncate the WAL, so new transactions don't pick up the old PageMap and then get a truncated WAL. Also, checkpoint should not abort if there's txs -- that's okay now. --- rbf/db.go | 30 +++++++++++++++++------------- rbf/db_test.go | 6 ++++-- rbf/tx.go | 5 ++--- 3 files changed, 23 insertions(+), 18 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 21596d144..cfa5ede43 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -278,8 +278,6 @@ func (db *DB) checkpoint() (err error) { }() if !db.opened { return nil - } else if len(db.txs) > 0 { - return nil // skip if transactions open } // Check if there are any WAL pages, if not do nothing as @@ -364,6 +362,8 @@ func (db *DB) checkpoint() (err error) { // the rwmu and update the metadata about the WAL. releaseLock = false // fmt.Printf("checkpoint mostly done, waiting for Tx cleanup...\n") + db.walPageN = 0 + db.pageMap = NewPageMap() db.afterCurrentTx(func() { // fmt.Printf("truncating WAL\n") defer db.rwmu.Unlock() @@ -374,8 +374,6 @@ func (db *DB) checkpoint() (err error) { } else if _, err = db.walFile.Seek(0, io.SeekStart); err != nil { db.logger.Errorf("seek wal file: %w", err) } - db.walPageN = 0 - db.pageMap = NewPageMap() // fmt.Printf("checkpoint actually done\n") }) return nil @@ -589,19 +587,22 @@ func (db *DB) Begin(writable bool) (_ *Tx, err error) { if db.isDead != nil { err := db.isDead cleanup() - db.mu.Unlock() return nil, err } - // Wait for WAL size to be below threshold. - for int64(db.walPageN*PageSize) > db.cfg.MaxWALCheckpointSize { - if db.isDead != nil { - err := db.isDead - cleanup() - db.mu.Unlock() - return nil, err + // Wait for WAL size to be below threshold, if we're going to write. + // Reads don't care. + if writable { + for int64(db.walPageN*PageSize) > db.cfg.MaxWALCheckpointSize { + if db.isDead != nil { + err := db.isDead + cleanup() + return nil, err + } + // This implicitly releases db.mu.Lock and comes back with it + // held again. + db.haltCond.Wait() } - db.haltCond.Wait() } tx := &Tx{ @@ -775,6 +776,9 @@ func (db *DB) baseWALID() int64 { // readWALPageByID reads a WAL page by WAL ID. func (db *DB) readWALPageByID(id int64) ([]byte, error) { + if id == db.baseWALID() { + fmt.Printf("id %d oops\n", id) + } return db.readWALPageAt(int(id - db.baseWALID() - 1)) } diff --git a/rbf/db_test.go b/rbf/db_test.go index 5c8504b1b..c0eb3b3c7 100644 --- a/rbf/db_test.go +++ b/rbf/db_test.go @@ -291,7 +291,8 @@ func TestDB_MultiTx(t *testing.T) { time.Sleep(time.Duration(rand.Intn(100)) * time.Millisecond) - for i := 0; i < rand.Intn(1000); i++ { + n := rand.Intn(500) + 500 + for i := 0; i < n; i++ { v := rand.Intn(1 << 20) if _, err := tx.Contains("x", uint64(v)); err != nil { return err @@ -316,7 +317,8 @@ func TestDB_MultiTx(t *testing.T) { } defer tx.Rollback() - for j := 0; j < rand.Intn(100); j++ { + n := rand.Intn(90) + 10 + for j := 0; j < n; j++ { v := rand.Intn(1 << 20) if _, err := tx.Add("x", uint64(v)); err != nil { t.Fatal(err) diff --git a/rbf/tx.go b/rbf/tx.go index 2cf7420db..bfb8fc00c 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -102,6 +102,8 @@ func (tx *Tx) Commit() error { // If any pages have been written, ensure we write a new meta page with // the commit flag to mark the end of the transaction. + tx.db.mu.Lock() + defer tx.db.mu.Unlock() if tx.dirty() { if err := tx.flush(); err != nil { return err @@ -118,12 +120,9 @@ func (tx *Tx) Commit() error { // the lock, because we need removeTx to grab the lock to // work, but if it wants to checkpoint, it wants to be able to return // to us here and still be holding the lock. - tx.db.mu.Lock() tx.db.rootRecords = tx.rootRecords tx.db.pageMap = tx.pageMap tx.db.walPageN = tx.walPageN - tx.db.mu.Unlock() - return tx.db.removeTx(tx) } // Disconnect transaction from DB. From 994cc03e88717642f3b06614eef6d7afa859329b Mon Sep 17 00:00:00 2001 From: Seebs Date: Thu, 16 Dec 2021 14:01:38 -0600 Subject: [PATCH 11/22] fix locking and list management for afterCurrentTx Two issues: First, there was a race condition because we were never using the mutex for anything but the condvar broadcast, second, there was no reason for the afterCurrentTx to need to maintain the list since we already know where in the list we are when we are waking it up. afterCurrentTx still wants to run with the db lock held, because the degenerate case (no outstanding Tx) means that it will be running with it held already. That's for another commit. --- rbf/db.go | 22 ++++++++++++---------- 1 file changed, 12 insertions(+), 10 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index cfa5ede43..0662f59d6 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -676,15 +676,6 @@ func (db *DB) afterCurrentTx(callback func()) { // fmt.Printf("afterCurrentTx: locking db\n") db.mu.Lock() defer db.mu.Unlock() - // remove us from the db's list - for i, v := range db.txWaiters { - if v == txw { - // remove us from the list - copy(db.txWaiters[i:], db.txWaiters[i+1:]) - db.txWaiters = db.txWaiters[:len(db.txWaiters)-1] - break - } - } // fmt.Printf("afterCurrentTx: running callback\n") txw.callback() }() @@ -716,12 +707,23 @@ func (db *DB) removeTx(tx *Tx) error { } // remove ourselves from the list of transactions the db is keeping. delete(tx.db.txs, tx) - for _, txw := range tx.db.txWaiters { + for i := 0; i < len(tx.db.txWaiters); i++ { + txw := tx.db.txWaiters[i] + // in practice this probably never matters, but theoretically the + // goroutine that's waiting on the condition variable may + // not have performed its first test on len(txw.waitingOn) yet. + txw.mu.Lock() delete(txw.waitingOn, tx) + txw.mu.Unlock() // let it know we're done. we've still got db.mu.lock, so it won't // happen just yet, but it'll be able to continue. if len(txw.waitingOn) == 0 { + // remove us from the db's list + copy(db.txWaiters[i:], db.txWaiters[i+1:]) + db.txWaiters = db.txWaiters[:len(db.txWaiters)-1] txw.cond.Broadcast() + // decrement i so we don't skip an entry we just copied in to [i] + i-- } } From 47e098c3b1a84c4e0108b10ba050ddf37381448e Mon Sep 17 00:00:00 2001 From: Seebs Date: Thu, 16 Dec 2021 15:54:06 -0600 Subject: [PATCH 12/22] simplify txWaiter We don't need a condition variable for a thing with a single waiter which waits only once, and a data structure which only one side ever modifies. That's a closable channel. --- rbf/db.go | 15 ++++----------- 1 file changed, 4 insertions(+), 11 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 0662f59d6..4129ac1fe 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -33,8 +33,7 @@ var cursorSyncPool = &sync.Pool{ // lock held, at some point after every Tx that was open when it was created // has closed. WARNING: A txWaiter may hold db.rwmu. type txWaiter struct { - mu sync.Mutex - cond *sync.Cond + ready chan struct{} waitingOn map[*Tx]struct{} callback func() } @@ -660,19 +659,15 @@ func (db *DB) afterCurrentTx(callback func()) { return } txw := &txWaiter{} - txw.cond = sync.NewCond(&txw.mu) + txw.ready = make(chan struct{}) txw.callback = callback txw.waitingOn = make(map[*Tx]struct{}, len(db.txs)) for k := range db.txs { txw.waitingOn[k] = struct{}{} } db.txWaiters = append(db.txWaiters, txw) - txw.mu.Lock() go func() { - for len(txw.waitingOn) > 0 { - // fmt.Printf("afterCurrentTx: %d left\n", len(txw.waitingOn)) - txw.cond.Wait() - } + <-txw.ready // fmt.Printf("afterCurrentTx: locking db\n") db.mu.Lock() defer db.mu.Unlock() @@ -712,16 +707,14 @@ func (db *DB) removeTx(tx *Tx) error { // in practice this probably never matters, but theoretically the // goroutine that's waiting on the condition variable may // not have performed its first test on len(txw.waitingOn) yet. - txw.mu.Lock() delete(txw.waitingOn, tx) - txw.mu.Unlock() // let it know we're done. we've still got db.mu.lock, so it won't // happen just yet, but it'll be able to continue. if len(txw.waitingOn) == 0 { // remove us from the db's list copy(db.txWaiters[i:], db.txWaiters[i+1:]) db.txWaiters = db.txWaiters[:len(db.txWaiters)-1] - txw.cond.Broadcast() + close(txw.ready) // decrement i so we don't skip an entry we just copied in to [i] i-- } From 57ca5591a264740daddb7e819f745f71c9163951 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Fri, 17 Dec 2021 10:59:09 -0700 Subject: [PATCH 13/22] Unlock rbf.DB during WAL copy & fsync() --- rbf/db.go | 149 +++++++++++++++++++++++++++++------------------------- 1 file changed, 79 insertions(+), 70 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index 4129ac1fe..e5a81690e 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -51,9 +51,10 @@ type DB struct { opened bool // true if open logger logger.Logger // for diagnostics from async things - wal []byte // wal mmap - walFile *os.File // wal file descriptor - walPageN int // wal page count + wal []byte // wal mmap + walFile *os.File // wal file descriptor + walPageN int // wal page count + baseWALID int64 // WAL ID of first page mu sync.RWMutex // general mutex rwmu sync.Mutex // mutex for restricting single writer @@ -238,6 +239,7 @@ func (db *DB) openWAL() (err error) { return fmt.Errorf("wal seek: %w", err) } db.walPageN = pageN + db.baseWALID = readMetaWALID(db.data) return nil } @@ -293,79 +295,92 @@ func (db *DB) checkpoint() (err error) { } db.haltCond.Broadcast() }() - var page []byte - // We might have either a *PageMap or just the file. If we have the file, - // building the PageMap is fairly expensive because it's fancy and immutable. - // If we have the PageMap *or* some other map, that's two different things - // to iterate. If we have the PageMap, building a map from it is relatively - // cheap, so we'll do it that way. - pages := make(map[uint32]int) - if db.pageMap.size == 0 { - // you'd think we're done, but actually this PROBABLY means that - // this is initial startup, and we haven't read the file yet. We scan - // the file for pages, because it turns out most of them probably - // got overwritten. - for i := 0; i < db.walPageN; i++ { - page, err = db.readWALPageAt(i) - if err != nil { - return fmt.Errorf("reading WAL page %d: %w", i, err) - } + // Copy the pages from the WAL back to the database outside of the lock. + if err := func() error { + db.mu.Unlock() // This is intentionally reversed so run w/o lock + defer db.mu.Lock() - // Determine page number. Meta pages are always on zero & bitmap - // headers specify the page number of the next page in the WAL. - // All other pages have their page number in the page data. - var pgno uint32 - if IsBitmapHeader(page) { - pgno = readPageNo(page) - if i+1 < db.walPageN { - if page, err = db.readWALPageAt(i + 1); err != nil { - return err - } - } else { - return fmt.Errorf("last page of WAL file (%d) is bitmap header", i) + var page []byte + // We might have either a *PageMap or just the file. If we have the file, + // building the PageMap is fairly expensive because it's fancy and immutable. + // If we have the PageMap *or* some other map, that's two different things + // to iterate. If we have the PageMap, building a map from it is relatively + // cheap, so we'll do it that way. + pages := make(map[uint32]int) + + if db.pageMap.size == 0 { + // you'd think we're done, but actually this PROBABLY means that + // this is initial startup, and we haven't read the file yet. We scan + // the file for pages, because it turns out most of them probably + // got overwritten. + for i := 0; i < db.walPageN; i++ { + page, err = db.readWALPageAt(i) + if err != nil { + return fmt.Errorf("reading WAL page %d: %w", i, err) } - i++ // bitmaps in WAL are two pages - } else if !IsMetaPage(page) { - pgno = readPageNo(page) + + // Determine page number. Meta pages are always on zero & bitmap + // headers specify the page number of the next page in the WAL. + // All other pages have their page number in the page data. + var pgno uint32 + if IsBitmapHeader(page) { + pgno = readPageNo(page) + if i+1 < db.walPageN { + if page, err = db.readWALPageAt(i + 1); err != nil { + return err + } + } else { + return fmt.Errorf("last page of WAL file (%d) is bitmap header", i) + } + i++ // bitmaps in WAL are two pages + } else if !IsMetaPage(page) { + pgno = readPageNo(page) + } + // record where in the file we have this page + pages[pgno] = i + } + } else { + itr := db.pageMap.Iterator() + itr.First() + for k, v, ok := itr.Next(); ok; k, v, ok = itr.Next() { + pages[k] = int(v - db.baseWALID - 1) } - // record where in the file we have this page - pages[pgno] = i } - } else { - walBase := db.baseWALID() - itr := db.pageMap.Iterator() - itr.First() - for k, v, ok := itr.Next(); ok; k, v, ok = itr.Next() { - pages[k] = int(v - walBase - 1) + + // fmt.Printf("checkpoint: walPageN %d, PageMap size %d\n", db.walPageN, db.pageMap.size) + for pgno, walID := range pages { + page, err = db.readWALPageAt(walID) + if err != nil { + return fmt.Errorf("reading page %d [page number %d]: %v", walID, pgno, err) + } + + // Write data to the data file. + if err = db.writeDBPage(pgno, page); err != nil { + return fmt.Errorf("writing page %d: %v", pgno, err) + } } + + // Ensure database file is synced and then truncate the WAL file. + if err = db.fsync(db.file); err != nil { + return fmt.Errorf("db file sync: %w", err) + } + + return nil + }(); err != nil { + return err } - // fmt.Printf("checkpoint: walPageN %d, PageMap size %d\n", db.walPageN, db.pageMap.size) - for pgno, walID := range pages { - page, err = db.readWALPageAt(walID) - if err != nil { - return fmt.Errorf("reading page %d [page number %d]: %v", walID, pgno, err) - } - // Write data to the data file. - if err = db.writeDBPage(pgno, page); err != nil { - return fmt.Errorf("writing page %d: %v", pgno, err) - } - } - // Ensure database file is synced and then truncate the WAL file. - if err = db.fsync(db.file); err != nil { - return fmt.Errorf("db file sync: %w", err) - } // now we've updated the file. There are existing transactions that are still // using the WAL, though. So we wait for them to terminate before we unlock // the rwmu and update the metadata about the WAL. releaseLock = false - // fmt.Printf("checkpoint mostly done, waiting for Tx cleanup...\n") db.walPageN = 0 db.pageMap = NewPageMap() + db.afterCurrentTx(func() { - // fmt.Printf("truncating WAL\n") defer db.rwmu.Unlock() + if err = db.walFile.Truncate(0); err != nil { db.logger.Errorf("truncate wal file: %w", err) } else if err = db.fsync(db.walFile); err != nil { @@ -373,8 +388,10 @@ func (db *DB) checkpoint() (err error) { } else if _, err = db.walFile.Seek(0, io.SeekStart); err != nil { db.logger.Errorf("seek wal file: %w", err) } - // fmt.Printf("checkpoint actually done\n") + + db.baseWALID = readMetaWALID(db.data) }) + return nil } @@ -764,17 +781,9 @@ func (db *DB) readDBPage(pgno uint32) ([]byte, error) { return db.data[offset : offset+PageSize], nil } -// baseWALID returns the WAL ID stored in the database file meta page. -func (db *DB) baseWALID() int64 { - return readMetaWALID(db.data) -} - // readWALPageByID reads a WAL page by WAL ID. func (db *DB) readWALPageByID(id int64) ([]byte, error) { - if id == db.baseWALID() { - fmt.Printf("id %d oops\n", id) - } - return db.readWALPageAt(int(id - db.baseWALID() - 1)) + return db.readWALPageAt(int(id - db.baseWALID - 1)) } // readWALPageAt reads the i-th page in the WAL file. From 1fd872b126356c07cac635888590067bf7f0682f Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Fri, 17 Dec 2021 12:55:32 -0600 Subject: [PATCH 14/22] less write locks in fragment.importRoaring/row --- ctl/server.go | 2 +- fragment.go | 75 ++++++++++++++++++++++++------------------------ server.go | 1 - server/config.go | 6 ++-- server/server.go | 2 +- 5 files changed, 42 insertions(+), 44 deletions(-) diff --git a/ctl/server.go b/ctl/server.go index 83edb5456..c5da42847 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -85,7 +85,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.BoolVar(&srv.Config.Storage.FsyncEnabled, "storage.fsync", true, "enable fsync fully safe flush-to-disk") // RowcacheOn - flags.BoolVar((&srv.Config.RowcacheOn), "rowcache-on", srv.Config.RowcacheOn, "turn on the rowcache for all backends (may speed some queries)") + flags.BoolVar((&srv.Config.RowcacheOn), "rowcache-on", srv.Config.RowcacheOn, "Do not use, permanently disabled. Flag exists for backwards compatibility and will be removed.") // RBF specific flags. See pilosa/rbf/cfg/cfg.go for definitions. srv.Config.RBFConfig.DefineFlags(flags) diff --git a/fragment.go b/fragment.go index 83514edad..9b77f90f7 100644 --- a/fragment.go +++ b/fragment.go @@ -593,8 +593,8 @@ func (f *fragment) mutexCheck(tx Tx, details bool, limit int) (map[uint64][]uint // row returns a row by ID. func (f *fragment) row(tx Tx, rowID uint64) (*Row, error) { - f.mu.Lock() - defer f.mu.Unlock() + f.mu.RLock() + defer f.mu.RUnlock() return f.unprotectedRow(tx, rowID) } @@ -937,9 +937,12 @@ func (f *fragment) unprotectedClearRow(tx Tx, rowID uint64) (changed bool, err e return changed, nil } -// unprotectedClearBlock clears all rows for a given block. +// clearBlock clears all rows for a given block. // This updates both the on-disk storage and the in-cache bitmap. -func (f *fragment) unprotectedClearBlock(tx Tx, block int) (changed bool, err error) { +func (f *fragment) clearBlock(tx Tx, block int) (changed bool, err error) { + f.mu.Lock() + defer f.mu.Unlock() + firstRow := uint64(block * HashBlockSize) var wp *io.Writer if f.storage != nil { @@ -2708,20 +2711,24 @@ func (f *fragment) importValue(tx Tx, columnIDs []uint64, values []int64, bitDep func (f *fragment) importRoaring(ctx context.Context, tx Tx, data []byte, clear bool) error { span, ctx := tracing.StartSpanFromContext(ctx, "fragment.importRoaring") defer span.Finish() - span, ctx = tracing.StartSpanFromContext(ctx, "importRoaring.AcquireFragmentLock") - f.mu.Lock() - defer f.mu.Unlock() - span.Finish() - return f.unprotectedImportRoaring(ctx, tx, data, clear) + rowSet, updateCache, err := f.doImportRoaring(ctx, tx, data, clear) + if err != nil { + return errors.Wrap(err, "doImportRoaring") + } + if updateCache { + return f.updateCachePostImport(ctx, rowSet) + } + return nil } -func (f *fragment) unprotectedImportRoaring(ctx context.Context, tx Tx, data []byte, clear bool) error { +func (f *fragment) doImportRoaring(ctx context.Context, tx Tx, data []byte, clear bool) (map[uint64]int, bool, error) { + f.mu.RLock() + defer f.mu.RUnlock() rowSize := uint64(1 << shardVsContainerExponent) span, ctx := tracing.StartSpanFromContext(ctx, "importRoaring.ImportRoaringBits") + defer span.Finish() - useRowCache := storage.RowCacheEnabled() - var changed int var rowSet map[uint64]int var wp *io.Writer if f.storage != nil { @@ -2734,37 +2741,37 @@ func (f *fragment) unprotectedImportRoaring(ctx context.Context, tx Tx, data []b return err } - changed, rowSet, err = tx.ImportRoaringBits(f.index(), f.field(), f.view(), f.shard, rit, clear, true, rowSize) + _, rowSet, err = tx.ImportRoaringBits(f.index(), f.field(), f.view(), f.shard, rit, clear, true, rowSize) return err }) - span.Finish() if err != nil { - return err + return nil, false, err } updateCache := f.CacheType != CacheTypeNone + return rowSet, updateCache, err +} + +func (f *fragment) updateCachePostImport(ctx context.Context, rowSet map[uint64]int) error { + f.mu.Lock() + defer f.mu.Unlock() anyChanged := false for rowID, changes := range rowSet { if changes == 0 { continue } - if useRowCache && f.rowCache != nil { - f.rowCache.Add(rowID, nil) - } - if updateCache { - anyChanged = true - if changes < 0 { - absChanges := uint64(-1 * changes) - if absChanges <= f.cache.Get(rowID) { - f.cache.BulkAdd(rowID, f.cache.Get(rowID)-absChanges) - } else { - f.cache.BulkAdd(rowID, 0) - } + anyChanged = true + if changes < 0 { + absChanges := uint64(-1 * changes) + if absChanges <= f.cache.Get(rowID) { + f.cache.BulkAdd(rowID, f.cache.Get(rowID)-absChanges) } else { - f.cache.BulkAdd(rowID, f.cache.Get(rowID)+uint64(changes)) + f.cache.BulkAdd(rowID, 0) } + } else { + f.cache.BulkAdd(rowID, f.cache.Get(rowID)+uint64(changes)) } } // we only set this if we need to update the cache @@ -2772,26 +2779,18 @@ func (f *fragment) unprotectedImportRoaring(ctx context.Context, tx Tx, data []b f.cache.Invalidate() } - span, _ = tracing.StartSpanFromContext(ctx, "importRoaring.incrementOpN") - - f.incrementOpN(changed) - - span.Finish() return nil } // importRoaringOverwrite overwrites the specified block with the provided data. func (f *fragment) importRoaringOverwrite(ctx context.Context, tx Tx, data []byte, block int) error { - f.mu.Lock() - defer f.mu.Unlock() - // Clear the existing data from fragment block. - if _, err := f.unprotectedClearBlock(tx, block); err != nil { + if _, err := f.clearBlock(tx, block); err != nil { return errors.Wrapf(err, "clearing block: %d", block) } // Union the new block data with the fragment data. - return f.unprotectedImportRoaring(ctx, tx, data, false) + return f.importRoaring(ctx, tx, data, false) } // incrementOpN increase the operation count by one. diff --git a/server.go b/server.go index 8858ab122..0c74a2fa5 100644 --- a/server.go +++ b/server.go @@ -476,7 +476,6 @@ func NewServer(opts ...ServerOption) (*Server, error) { } s.holder = NewHolder(path, s.holderConfig) s.holder.Stats.SetLogger(s.logger) - s.holder.Logger.Infof("RowCacheOn: %v", s.holderConfig.RowcacheOn) cwd, err := os.Getwd() if err != nil { return nil, err diff --git a/server/config.go b/server/config.go index c215d1596..09d82bfdb 100644 --- a/server/config.go +++ b/server/config.go @@ -200,9 +200,9 @@ type Config struct { // "rbf". Storage *storage.Config `toml:"storage"` - // RowcacheOn, if true, turns on the row cache for all storage backends. - // The default is now off because it makes rbf queries faster and uses - // much less memory. + // RowcacheOn permanently disabled. No longer useful w/ RBF. Left + // for backward compatibility but will be removed in a future + // version. RowcacheOn bool `toml:"rowcache-on"` // RBFConfig defines all externally configurable RBF flags. diff --git a/server/server.go b/server/server.go index a6d0049ae..b373d9e35 100644 --- a/server/server.go +++ b/server/server.go @@ -482,7 +482,7 @@ func (m *Command) SetupServer() error { pilosa.OptServerClusterName(m.Config.Cluster.Name), pilosa.OptServerSerializer(proto.Serializer{}), pilosa.OptServerStorageConfig(m.Config.Storage), - pilosa.OptServerRowcacheOn(m.Config.RowcacheOn), + pilosa.OptServerRowcacheOn(false), pilosa.OptServerRBFConfig(m.Config.RBFConfig), pilosa.OptServerMaxQueryMemory(m.Config.MaxQueryMemory), pilosa.OptServerQueryHistoryLength(m.Config.QueryHistoryLength), From 07998622667818811715d993c1f220820a751435 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Fri, 17 Dec 2021 14:50:08 -0600 Subject: [PATCH 15/22] move some locks, nbd --- rbf/db.go | 4 +++- rbf/tx.go | 6 ++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/rbf/db.go b/rbf/db.go index e5a81690e..f6b9ce5bc 100644 --- a/rbf/db.go +++ b/rbf/db.go @@ -380,6 +380,9 @@ func (db *DB) checkpoint() (err error) { db.afterCurrentTx(func() { defer db.rwmu.Unlock() + db.baseWALID = readMetaWALID(db.data) + db.mu.Unlock() + defer db.mu.Lock() if err = db.walFile.Truncate(0); err != nil { db.logger.Errorf("truncate wal file: %w", err) @@ -389,7 +392,6 @@ func (db *DB) checkpoint() (err error) { db.logger.Errorf("seek wal file: %w", err) } - db.baseWALID = readMetaWALID(db.data) }) return nil diff --git a/rbf/tx.go b/rbf/tx.go index bfb8fc00c..bceabd8f1 100644 --- a/rbf/tx.go +++ b/rbf/tx.go @@ -102,8 +102,6 @@ func (tx *Tx) Commit() error { // If any pages have been written, ensure we write a new meta page with // the commit flag to mark the end of the transaction. - tx.db.mu.Lock() - defer tx.db.mu.Unlock() if tx.dirty() { if err := tx.flush(); err != nil { return err @@ -120,11 +118,15 @@ func (tx *Tx) Commit() error { // the lock, because we need removeTx to grab the lock to // work, but if it wants to checkpoint, it wants to be able to return // to us here and still be holding the lock. + tx.db.mu.Lock() tx.db.rootRecords = tx.rootRecords tx.db.pageMap = tx.pageMap tx.db.walPageN = tx.walPageN + tx.db.mu.Unlock() } + tx.db.mu.Lock() + defer tx.db.mu.Unlock() // Disconnect transaction from DB. return tx.db.removeTx(tx) } From 41f6156bda7e7e21e37e2ca65f1675c03560e688 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 17 Dec 2021 17:57:33 -0600 Subject: [PATCH 16/22] don't use write Tx even when we're using the expensive logic for write Tx --- txfactory.go | 15 ++++++++++----- 1 file changed, 10 insertions(+), 5 deletions(-) diff --git a/txfactory.go b/txfactory.go index 62f527653..12fa12981 100644 --- a/txfactory.go +++ b/txfactory.go @@ -242,10 +242,15 @@ func (qcx *Qcx) GetTx(o Txo) (tx Tx, finisher func(perr *error), err error) { } // qcx.write reflects the top executor determination - // if a write will be done at the end, so we upgrade - // the "local" read Tx to be writes, so that they - // don't deadlock against themselves. - o.Write = o.Write || qcx.write + // if a write will be happen at some point, in which case, to avoid + // locking problems with multi-shard things, we (probably incorrectly) + // treat every Tx as its own individual separate Tx. + // + // But we still want to open non-write transactions individually, we + // just can't recycle them (because write operations will come in and + // we want them to work and commit right away so we're not holding a write + // lock for long). + writeLogic := o.Write || qcx.write // In general, we make ALL write transactions local, and never reuse them // below. Previously this was to help lmdb. @@ -273,7 +278,7 @@ func (qcx *Qcx) GetTx(o Txo) (tx Tx, finisher func(perr *error), err error) { return *qcx.RequiredForAtomicWriteTx, NoopFinisher, nil } - if !o.Write && qcx.Grp != nil { + if !writeLogic && qcx.Grp != nil { // read, with a group in place. finisher = func(perr *error) {} // finisher is a returned value From 9945575bf111f3193ed02212dd4ddc109999a0c0 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 17 Dec 2021 21:21:03 -0600 Subject: [PATCH 17/22] create a new worker every so often if progress isn't happening this is very approximate and may be a mess and may be unbounded, but in practice i think it should be okay. if it's not we'll have an adventure. --- executor.go | 48 ++++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 42 insertions(+), 6 deletions(-) diff --git a/executor.go b/executor.go index 8b0409eb7..899ac0e19 100644 --- a/executor.go +++ b/executor.go @@ -51,6 +51,9 @@ type executor struct { Node *topology.Node Cluster *cluster + // how many jobs the work queue has seen + workCounter uint64 + // Client used for remote requests. client InternalQueryClient @@ -128,15 +131,47 @@ func newExecutor(opts ...executorOption) *executor { e.work = make(chan job, e.workerPoolSize) _ = testhook.Opened(NewAuditor(), e, nil) for i := 0; i < e.workerPoolSize; i++ { - e.workersWG.Add(1) - go func() { - defer e.workersWG.Done() - worker(e.work) - }() + e.addWorker() } + go func() { + // background task: every so often, check to see whether we have + // work in the queue but none has been taken for a while. if so, we + // need more workers. + prev := atomic.LoadUint64(&e.workCounter) + periodic := time.NewTicker(50 * time.Millisecond) + defer periodic.Stop() + running := true + for running { + <-periodic.C + func() { + e.workMu.Lock() + defer e.workMu.Unlock() + if e.shutdown { + running = false + return + } + if len(e.work) == 0 { + return + } + next := atomic.LoadUint64(&e.workCounter) + if next == prev { + e.addWorker() + prev = next + } + }() + } + }() return e } +func (e *executor) addWorker() { + e.workersWG.Add(1) + go func() { + defer e.workersWG.Done() + e.worker(e.work) + }() +} + func (e *executor) Close() error { e.workMu.Lock() defer e.workMu.Unlock() @@ -5935,8 +5970,9 @@ type job struct { resultChan chan mapResponse } -func worker(work chan job) { +func (e *executor) worker(work chan job) { for j := range work { + atomic.AddUint64(&e.workCounter, 1) // Skip out early if the context is done, but still send // an ack so mapperLocal can be sure we aren't about to // work on something it sent us. From 8f217ab099f4a6a3952b15ba190c6311133c6f39 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 17 Dec 2021 22:05:41 -0600 Subject: [PATCH 18/22] scale down worker pool when it's large if we have more than twice our starting worker pool, and have had no tasks when checking the queue for multiple rounds, send a job telling the system to retire a worker. eventually we'll get down to about 2x the starting pool size if we stay idle. --- executor.go | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/executor.go b/executor.go index 899ac0e19..be4def11a 100644 --- a/executor.go +++ b/executor.go @@ -64,6 +64,7 @@ type executor struct { workMu sync.RWMutex workersWG sync.WaitGroup workerPoolSize int + currentWorkers int work chan job // Maximum per-request memory usage (Extract() only) @@ -141,6 +142,7 @@ func newExecutor(opts ...executorOption) *executor { periodic := time.NewTicker(50 * time.Millisecond) defer periodic.Stop() running := true + idle := 0 for running { <-periodic.C func() { @@ -151,6 +153,17 @@ func newExecutor(opts ...executorOption) *executor { return } if len(e.work) == 0 { + idle++ + if idle > 10 && e.currentWorkers > (e.workerPoolSize*2) { + select { + case e.work <- job{idleHands: true}: + // we closed an excess worker + default: + // somehow between our test above and now the work + // queue FILLED UP and we stoically accept this + } + } + idle = 0 return } next := atomic.LoadUint64(&e.workCounter) @@ -166,9 +179,11 @@ func newExecutor(opts ...executorOption) *executor { func (e *executor) addWorker() { e.workersWG.Add(1) + e.currentWorkers++ go func() { defer e.workersWG.Done() e.worker(e.work) + e.currentWorkers-- }() } @@ -5968,11 +5983,15 @@ type job struct { ctx context.Context memoryAvailable *int64 // shared, atomic value resultChan chan mapResponse + idleHands bool } func (e *executor) worker(work chan job) { for j := range work { atomic.AddUint64(&e.workCounter, 1) + if j.idleHands { + return + } // Skip out early if the context is done, but still send // an ack so mapperLocal can be sure we aren't about to // work on something it sent us. From 9a2a8f964c99f177697661e9c98dad59eb968ff5 Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 17 Dec 2021 22:22:27 -0600 Subject: [PATCH 19/22] fix silly typo in worker pool downscaling --- executor.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/executor.go b/executor.go index be4def11a..69bdbf15d 100644 --- a/executor.go +++ b/executor.go @@ -162,8 +162,8 @@ func newExecutor(opts ...executorOption) *executor { // somehow between our test above and now the work // queue FILLED UP and we stoically accept this } + idle = 0 } - idle = 0 return } next := atomic.LoadUint64(&e.workCounter) From 1439c316d348ac43d084e88d414e5899ecdfe9dd Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 17 Dec 2021 22:25:01 -0600 Subject: [PATCH 20/22] read-only lock for check of shutdown --- executor.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/executor.go b/executor.go index 69bdbf15d..45d5169c5 100644 --- a/executor.go +++ b/executor.go @@ -146,8 +146,8 @@ func newExecutor(opts ...executorOption) *executor { for running { <-periodic.C func() { - e.workMu.Lock() - defer e.workMu.Unlock() + e.workMu.RLock() + defer e.workMu.RUnlock() if e.shutdown { running = false return From 8974014d5783dc61c6a850c03f00119194e10f1a Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 17 Dec 2021 22:48:40 -0600 Subject: [PATCH 21/22] too tired to be writing code --- executor.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/executor.go b/executor.go index 45d5169c5..fb5ef68ce 100644 --- a/executor.go +++ b/executor.go @@ -169,8 +169,8 @@ func newExecutor(opts ...executorOption) *executor { next := atomic.LoadUint64(&e.workCounter) if next == prev { e.addWorker() - prev = next } + prev = next }() } }() From 7826c06eee6014cc1ec5b9079f94ed86c53e1dee Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Sat, 18 Dec 2021 08:58:19 -0600 Subject: [PATCH 22/22] use atomics for currentWorker to avoid race --- executor.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/executor.go b/executor.go index fb5ef68ce..5e03e1a72 100644 --- a/executor.go +++ b/executor.go @@ -64,7 +64,7 @@ type executor struct { workMu sync.RWMutex workersWG sync.WaitGroup workerPoolSize int - currentWorkers int + currentWorkers int64 work chan job // Maximum per-request memory usage (Extract() only) @@ -154,7 +154,7 @@ func newExecutor(opts ...executorOption) *executor { } if len(e.work) == 0 { idle++ - if idle > 10 && e.currentWorkers > (e.workerPoolSize*2) { + if idle > 10 && atomic.LoadInt64(&e.currentWorkers) > int64(e.workerPoolSize*2) { select { case e.work <- job{idleHands: true}: // we closed an excess worker @@ -179,11 +179,11 @@ func newExecutor(opts ...executorOption) *executor { func (e *executor) addWorker() { e.workersWG.Add(1) - e.currentWorkers++ + atomic.AddInt64(&e.currentWorkers, 1) go func() { defer e.workersWG.Done() e.worker(e.work) - e.currentWorkers-- + atomic.AddInt64(&e.currentWorkers, -1) }() }