diff --git a/memdb/batch_release_test.go b/memdb/batch_release_test.go index 96edf3c..72895de 100644 --- a/memdb/batch_release_test.go +++ b/memdb/batch_release_test.go @@ -178,3 +178,53 @@ func TestRewritesFreeBlocks(t *testing.T) { } t.Logf("%d time blocks, %d logs, %d keys", timeBlockCount(db), len(logs), db.Size()) } + +// TestBlocksDeletingFromEachOtherGo puts a key in the current block, then a +// batch of it and another, then the other in the current block again: the +// batch's block deletes from the current block, which deletes from the +// batch's. Each waited for the other's logs to go, and both kept theirs for +// good; a store recovered from logs of several versions of a key had +// thousands of such blocks. +func TestBlocksDeletingFromEachOtherGo(t *testing.T) { + const d = 300 * time.Millisecond + dir := t.TempDir() + db, err := Open(WithLogFilePath(dir), WithTimeBlockInterval(d), WithLogInterval(2*time.Millisecond)) + if err != nil { + t.Fatal(err) + } + defer db.Close() + // Start at the beginning of a block, so the writes share one. + time.Sleep(time.Until(time.Now().Truncate(d).Add(d))) + if _, err := db.Put(1, []byte("current")); err != nil { + t.Fatal(err) + } + if err := db.Batch(func(b *Batch, _ <-chan struct{}) error { + if err := b.Put(1, []byte("batch")); err != nil { + return err + } + return b.Put(2, []byte("batch")) + }); err != nil { + t.Fatal(err) + } + if _, err := db.Put(2, []byte("current")); err != nil { + t.Fatal(err) + } + for _, k := range []uint64{1, 2} { + if err := db.Delete(k); err != nil { + t.Fatal(err) + } + } + if err := db.Flush(); err != nil { + t.Fatal(err) + } + for deadline := time.Now().Add(3 * time.Second); ; { + logs, _ := filepath.Glob(filepath.Join(dir, logDir, "*.log")) + if timeBlockCount(db) <= 1 && len(logs) <= 1 { + break + } + if time.Now().After(deadline) { + t.Fatalf("%d time blocks and %d logs left", timeBlockCount(db), len(logs)) + } + time.Sleep(d / 3) + } +} diff --git a/memdb/block.go b/memdb/block.go index 6b580c8..9cb8978 100644 --- a/memdb/block.go +++ b/memdb/block.go @@ -46,15 +46,20 @@ type ( // The block's logs in the WAL, and what keeps them there; guarded // by the DB's logMu, not the block's lock. A log holding the delete // of a version must stay in the WAL as long as the log holding its - // put, or the version comes back on the next recovery: a block's - // logs go once the block is released and the logs of every block - // it deletes versions from have gone. - timeRefs []_TimeID + // put, or the version comes back on the next recovery: see + // applyLogs. + timeRefs []_LogRef state blockState deletes map[*_Block]bool // blocks it deletes versions from - waitFor int // blocks in deletes whose logs are in the WAL waiters []*_Block // blocks deleting versions from this one } + + // _LogRef is a log of a block, and its place in the order logs were + // written to the WAL, which recovery replays them in. + _LogRef struct { + id _TimeID + seq uint64 + } ) // blockState is where a block is in its life, which goes one way: diff --git a/memdb/db.go b/memdb/db.go index ed45412..a01b88c 100644 --- a/memdb/db.go +++ b/memdb/db.go @@ -257,10 +257,13 @@ func (db *DB) replace(loc _Loc, key uint64) (bool, error) { } // releaseEmptied releases a block a delete emptied, unless writes may still -// go to it; one with more to write is released once it is (releaseEmpty). -// The caller holds no index shard: releasing takes them. +// go to it, as they do to the current block; one with more to write is +// released once it is (releaseEmpty). A batch's block, written, takes no +// more writes, though its time ID may be past the current block's: it was +// left, emptied, for good. The caller holds no index shard: releasing takes +// them. func (db *DB) releaseEmptied(emptied bool, timeID _TimeID) error { - if !emptied || timeID >= db.timeID() { + if !emptied || timeID == db.timeID() { return nil } if err := db.releaseLog(timeID); err != nil && err != errEntryDoesNotExist { diff --git a/memdb/db_internal.go b/memdb/db_internal.go index af06207..00b994a 100644 --- a/memdb/db_internal.go +++ b/memdb/db_internal.go @@ -19,6 +19,7 @@ package memdb import ( "errors" "io" + "sort" "sync/atomic" "time" @@ -73,6 +74,9 @@ type _DB struct { // logMu guards the blocks' logs and what keeps them in the WAL. It is // taken after a block's lock and before db.mu. logMu mutex[logRank] + // logSeq counts the logs written to the WAL, or replayed from it. + // Guarded by logMu. + logSeq uint64 // close closed uint32 @@ -181,7 +185,7 @@ func (db *DB) blocks() []*_Block { } // deleteFrom records that b holds deletes of versions in from: b's logs -// stay in the WAL as long as from's. The caller holds logMu. +// stay in the WAL as long as from's (applyLogs). The caller holds logMu. func (b *_Block) deleteFrom(from *_Block) { if from == b || from.state == blockGone || b.deletes[from] { return @@ -190,26 +194,60 @@ func (b *_Block) deleteFrom(from *_Block) { b.deletes = make(map[*_Block]bool) } b.deletes[from] = true - b.waitFor++ from.waiters = append(from.waiters, b) } -// applyLogs marks the logs of a released block applied once the logs of -// the blocks it deletes versions from are, and then those of the blocks -// that waited for it. It returns the logs, in the order they may go. The -// caller holds logMu. +// applyLogs marks the logs of a released block applied, with those of the +// blocks it deletes versions from, and from which they delete, and so on, +// once every one of them is released; then those of the blocks that waited +// for them. It returns the logs, in the order they may go. The caller holds +// logMu. +// +// A block's logs went once the logs of the blocks it deletes from had gone. +// Blocks can delete from each other, a batch's block and the block current +// as it was written, or blocks a WAL of several versions of a key recovers: +// each waited for the other, and kept its logs for good. Logs can't: a +// delete is written after the put it deletes. The blocks go together, and +// their logs in the order they were written, so that a put's log goes +// before its delete's, and a crash between leaves no version back. func (b *_Block) applyLogs() []_TimeID { - if b.state != blockReleased || b.waitFor > 0 { + if b.state != blockReleased { return nil } - b.setState(blockGone) - logs := b.timeRefs - for _, w := range b.waiters { - w.waitFor-- - logs = append(logs, w.applyLogs()...) + var group []*_Block + seen := map[*_Block]bool{b: true} + for stack := []*_Block{b}; len(stack) > 0; { + x := stack[len(stack)-1] + stack = stack[:len(stack)-1] + if x.state == blockLive { + // A version it deletes is in a block in use. + return nil + } + group = append(group, x) + for d := range x.deletes { + if d.state != blockGone && !seen[d] { + seen[d] = true + stack = append(stack, d) + } + } + } + var refs []_LogRef + for _, x := range group { + x.setState(blockGone) + refs = append(refs, x.timeRefs...) + } + sort.Slice(refs, func(i, j int) bool { return refs[i].seq < refs[j].seq }) + logs := make([]_TimeID, len(refs)) + for i, r := range refs { + logs[i] = r.id + } + for _, x := range group { + for _, w := range x.waiters { + logs = append(logs, w.applyLogs()...) + } + x.waiters = nil + x.deletes = nil } - b.waiters = nil - b.deletes = nil return logs } @@ -263,7 +301,7 @@ func (db *DB) tinyWrite(tinyLog *_TinyLog) error { if block.state == blockGone { return db.internal.wal.SignalLogApplied(int64(tinyLog.ID())) } - block.timeRefs = append(block.timeRefs, tinyLog.ID()) + block.timeRefs = append(block.timeRefs, db.nextLogRef(tinyLog.ID())) return nil } @@ -359,6 +397,13 @@ func (db *DB) releaseLog(timeID _TimeID) error { return nil } +// nextLogRef returns the ref of a log written to the WAL now, the next in +// the order written. The caller holds logMu. +func (db *DB) nextLogRef(id _TimeID) _LogRef { + db.internal.logSeq++ + return _LogRef{id: id, seq: db.internal.logSeq} +} + // signalApplied marks logs applied, in order. func (db *DB) signalApplied(logs []_TimeID) error { for _, timeRef := range logs { diff --git a/memdb/recovery.go b/memdb/recovery.go index bddefef..5e1bdc9 100644 --- a/memdb/recovery.go +++ b/memdb/recovery.go @@ -125,7 +125,7 @@ func (db *DB) startRecovery() error { block.Lock() block.lastOffset = block.size() block.Unlock() - block.timeRefs = append(block.timeRefs, _TimeID(ID)) + block.timeRefs = append(block.timeRefs, db.nextLogRef(_TimeID(ID))) db.internal.timeMark.release(timeID) db.internal.meter.Recovers.Inc(puts) return false, nil diff --git a/server/e2e/client_test.go b/server/e2e/client_test.go index 688825e..71f11cd 100644 --- a/server/e2e/client_test.go +++ b/server/e2e/client_test.go @@ -371,19 +371,27 @@ func (c *client) connectWith(o connectOpts) (*utp.ConnectAcknowledge, error) { if err := c.writeRaw(buf.Bytes()); err != nil { return nil, err } + var p *utp.ControlMessage select { case <-c.closed: - return nil, fmt.Errorf("connection closed during connect: %v", c.readErr) - case p := <-c.connack: - ack := &utp.ConnectAcknowledge{} - ack.FromBinary(utp.FixedHeader{}, p.Message) - if ack.ReturnCode != utp.Accepted { - return ack, fmt.Errorf("connect refused, return code %d", ack.ReturnCode) + // A refusal is written and the connection closed at once: the + // CONNACK can be waiting when the close is seen, and select took + // either. + select { + case p = <-c.connack: + default: + return nil, fmt.Errorf("connection closed during connect: %v", c.readErr) } - return ack, nil + case p = <-c.connack: case <-time.After(3 * time.Second): return nil, fmt.Errorf("connect timeout") } + ack := &utp.ConnectAcknowledge{} + ack.FromBinary(utp.FixedHeader{}, p.Message) + if ack.ReturnCode != utp.Accepted { + return ack, fmt.Errorf("connect refused, return code %d", ack.ReturnCode) + } + return ack, nil } func (c *client) publish(mode uint8, topic string, payload []byte, ttl string) (uint16, error) {