Skip to content
Merged
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
50 changes: 50 additions & 0 deletions memdb/batch_release_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
15 changes: 10 additions & 5 deletions memdb/block.go
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
9 changes: 6 additions & 3 deletions memdb/db.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
75 changes: 60 additions & 15 deletions memdb/db_internal.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package memdb
import (
"errors"
"io"
"sort"
"sync/atomic"
"time"

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

Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion memdb/recovery.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
22 changes: 15 additions & 7 deletions server/e2e/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading