From a32147ea6a0b18440cc6c8776fd84ebc08985716 Mon Sep 17 00:00:00 2001 From: unit-adm <314923187+mumtaz6@users.noreply.github.com> Date: Tue, 6 Oct 2026 13:58:08 +0530 Subject: [PATCH 1/2] memdb: compact the values left in mostly dead blocks A block's logs stay in the WAL while it holds a value, and so do the logs of every block that deletes from it, and of those that delete from them: a delete must outlive the put it deletes. A value kept for long, a session's row, a message waiting on a subscriber away, then keeps the logs of most blocks written after it, in a store that rewrites and deletes its keys. A server's message store on v0.7.0 kept every log since it started, about 12 a minute idle; another, 1,379 logs for 1,644 keys, each block held by one value or chained to one. Compact moves the values left in blocks whose values take up at most half their data to the current block, so that those blocks hold none and go, with the logs that waited on them. A value moves under its key's index lock, as a put does: a put or delete of the key waits, other writes go on, and a key put again or deleted since it was listed doesn't move. The engine's own memdb doesn't compact: it finds an entry by its block. The server's message store compacts at open and every minute. A test store that kept 152 logs for one value keeps 1 after compacting, moving 2. Co-Authored-By: Claude Opus 5.5 (1M context) --- memdb/compact.go | 162 +++++++++++++++ memdb/compact_test.go | 236 ++++++++++++++++++++++ memdb/meter.go | 35 +++- memdb/model_test.go | 10 +- server/internal/db/unitdb/adapter.go | 48 +++++ server/internal/db/unitdb/adapter_test.go | 48 +++++ 6 files changed, 527 insertions(+), 12 deletions(-) create mode 100644 memdb/compact.go create mode 100644 memdb/compact_test.go diff --git a/memdb/compact.go b/memdb/compact.go new file mode 100644 index 0000000..f6cc2b9 --- /dev/null +++ b/memdb/compact.go @@ -0,0 +1,162 @@ +/* + * Copyright 2020 Saffat Technologies, Ltd. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package memdb + +import "encoding/binary" + +// A block's logs stay in the WAL while it holds a value, and so do the logs +// of every block that deletes from it, and of those that delete from them +// (applyLogs): a delete must outlive the put it deletes. A value kept for +// long, a session's row, a message waiting on a subscriber away, then keeps +// the logs of most blocks written after it, in a store that rewrites and +// deletes its keys: a server's message store kept 85,000. +// +// Compact moves the values left in mostly dead blocks to the current block, +// so that those blocks hold none, and go, with the logs that waited on them. +// A block is mostly dead when its values take up at most half its data: one +// mostly live isn't worth moving, and goes as its values are replaced. A +// value moves under its key's index lock, as a put does: a put or delete of +// the key waits, and other writes go on. +// +// A store that keeps a value per entry and frees blocks itself, as the +// engine's does, doesn't compact: it finds an entry by its block. + +// Compact moves the values left in mostly dead blocks to the current block, +// and returns how many it moved. +func (db *DB) Compact() (int, error) { + if err := db.ok(); err != nil { + return 0, err + } + type candidate struct { + block *_Block + keys []uint64 + } + current := db.timeID() + var candidates []candidate + for timeID, b := range db.blocksByID() { + if timeID == current { + continue + } + b.RLock() + if b.data == nil || b.count == 0 { + b.RUnlock() + continue + } + var live int64 + for ik, off := range b.records { + if ik.delFlag == 0 { + live += b.entryLen(off) + } + } + mostlyDead := live*2 <= b.size() + var keys []uint64 + if mostlyDead { + keys = b.liveKeys() + } + b.RUnlock() + if mostlyDead { + candidates = append(candidates, candidate{block: b, keys: keys}) + } + } + + moved := 0 + for _, c := range candidates { + for _, key := range c.keys { + ok, err := db.move(key, c.block) + if err != nil { + return moved, err + } + if ok { + moved++ + } + } + } + db.internal.meter.Compactions.Inc(1) + db.internal.meter.Moves.Inc(int64(moved)) + return moved, nil +} + +// move moves key's value from the block from to the current block, if it is +// there still, as Put would put it again, and reports whether it did. +func (db *DB) move(key uint64, from *_Block) (bool, error) { + if err := db.ok(); err != nil { + return false, err + } + db.internal.logManager.rotateMu.RLock() + defer db.internal.logManager.rotateMu.RUnlock() + timeID := db.timeID() + block, ok := db.timeBlock(timeID) + if !ok { + return false, errForbidden + } + sh := db.index.shard(key) + sh.Lock() + loc, had := sh.keys[key] + if !had || loc.block != from || from == block { + // Put again, deleted, or moved since it was listed. + sh.Unlock() + return false, nil + } + from.RLock() + off, ok := from.records[iKey(false, key)] + var val []byte + var err error + if ok { + val, err = from.get(off) + } + from.RUnlock() + if !ok || err != nil { + sh.Unlock() + return false, err + } + emptied, err := db.replace(loc, key) + if err != nil { + sh.Unlock() + return false, err + } + delete(sh.keys, key) + err = db.putEntry(block, key, val) + if err == nil { + sh.keys[key] = _Loc{timeID: timeID, block: block} + } + sh.Unlock() + if err != nil { + return false, err + } + return true, db.releaseEmptied(emptied, loc.timeID) +} + +// blocksByID returns the blocks in use, by time ID. +func (db *DB) blocksByID() map[_TimeID]*_Block { + db.mu.RLock() + defer db.mu.RUnlock() + blocks := make(map[_TimeID]*_Block, len(db.timeBlocks)) + for id, b := range db.timeBlocks { + blocks[id] = b + } + return blocks +} + +// entryLen returns the length of the entry at off, as put wrote it. The +// caller holds the block's lock. +func (b *_Block) entryLen(off int64) int64 { + head, err := b.data.Slice(off, off+4) + if err != nil || len(head) < 4 { + return 0 + } + return int64(binary.LittleEndian.Uint32(head)) +} diff --git a/memdb/compact_test.go b/memdb/compact_test.go new file mode 100644 index 0000000..842d366 --- /dev/null +++ b/memdb/compact_test.go @@ -0,0 +1,236 @@ +package memdb + +import ( + "fmt" + "math/rand" + "path/filepath" + "sync" + "testing" + "time" +) + +func walLogs(t *testing.T, dir string) int { + t.Helper() + logs, err := filepath.Glob(filepath.Join(dir, logDir, "*.log")) + if err != nil { + t.Fatal(err) + } + return len(logs) +} + +// TestCompactFreesPinnedBlocks keeps a key put first, and rewrites and +// deletes others over many blocks, as a message store does: the key holds +// its block, and every block chaining deletes back to it keeps its logs. +// Compact moves the key, and the blocks go with their logs. +func TestCompactFreesPinnedBlocks(t *testing.T) { + const d = 10 * time.Millisecond + dir := t.TempDir() + db, err := Open(WithLogFilePath(dir), WithTimeBlockInterval(d), WithLogInterval(2*time.Millisecond)) + if err != nil { + t.Fatal(err) + } + const pinned = 7 + if _, err := db.Put(pinned, []byte("kept")); err != nil { + t.Fatal(err) + } + // A key put in the first block too, and rewritten in each later one: + // each rewrite deletes from the block before, which chains back. + for i := 0; i < 150; i++ { + if _, err := db.Put(8, []byte(fmt.Sprintf("row-%d", i))); err != nil { + t.Fatal(err) + } + k := uint64(1000 + i) + if _, err := db.Put(k, []byte("message")); err != nil { + t.Fatal(err) + } + if err := db.Delete(k); err != nil { + t.Fatal(err) + } + time.Sleep(d) + } + if err := db.Flush(); err != nil { + t.Fatal(err) + } + time.Sleep(5 * d) + before := walLogs(t, dir) + if before < 50 { + t.Fatalf("%d logs before compacting: the key didn't hold the blocks", before) + } + + moved, err := db.Compact() + if err != nil { + t.Fatal(err) + } + if moved == 0 { + t.Fatal("Compact moved no value") + } + if err := db.Flush(); err != nil { + t.Fatal(err) + } + var after int + for deadline := time.Now().Add(2 * time.Second); ; time.Sleep(5 * d) { + if after = walLogs(t, dir); after <= 10 || time.Now().After(deadline) { + break + } + } + if after > 10 { + t.Fatalf("%d logs after compacting, %d before", after, before) + } + t.Logf("moved %d values; %d logs before, %d after", moved, before, after) + if err := db.Verify(); err != nil { + t.Fatal(err) + } + check := func(db *DB, when string) { + t.Helper() + if v, err := db.Get(pinned); err != nil || string(v) != "kept" { + t.Errorf("%s: the kept key reads %q, %v", when, v, err) + } + if v, err := db.Get(8); err != nil || string(v) != "row-149" { + t.Errorf("%s: the rewritten key reads %q, %v", when, v, err) + } + if v, err := db.Get(1000); err == nil { + t.Errorf("%s: a deleted key reads %q", when, v) + } + if n := db.Size(); n != 2 { + t.Errorf("%s: %d keys; want 2", when, n) + } + } + check(db, "compacted") + if err := db.Close(); err != nil { + t.Fatal(err) + } + db, err = Open(WithLogFilePath(dir), WithTimeBlockInterval(d)) + if err != nil { + t.Fatal(err) + } + defer db.Close() + check(db, "reopened") + if err := db.Verify(); err != nil { + t.Fatal(err) + } +} + +// TestCompactLeavesMostlyLiveBlocks compacts a store whose blocks are mostly +// live: nothing moves. +func TestCompactLeavesMostlyLiveBlocks(t *testing.T) { + const d = 10 * time.Millisecond + db, err := Open(WithLogFilePath(t.TempDir()), WithTimeBlockInterval(d), WithLogInterval(2*time.Millisecond)) + if err != nil { + t.Fatal(err) + } + defer db.Close() + for i := 0; i < 5; i++ { + for k := 0; k < 10; k++ { + if _, err := db.Put(uint64(i*10+k), []byte("value")); err != nil { + t.Fatal(err) + } + } + // One of ten deleted: the block stays mostly live. + if err := db.Delete(uint64(i * 10)); err != nil { + t.Fatal(err) + } + time.Sleep(2 * d) + } + if moved, err := db.Compact(); err != nil || moved != 0 { + t.Fatalf("Compact moved %d, %v; want none", moved, err) + } +} + +// TestCompactWhileWriting compacts over and over while writers put and +// delete keys of their own: each key keeps the last value its writer put, +// or none after its delete. +func TestCompactWhileWriting(t *testing.T) { + const d = 5 * time.Millisecond + dir := t.TempDir() + db, err := Open(WithLogFilePath(dir), WithTimeBlockInterval(d), WithLogInterval(time.Millisecond)) + if err != nil { + t.Fatal(err) + } + const writers, keys = 4, 20 + want := make([]map[uint64]string, writers) + stop := make(chan struct{}) + var wg sync.WaitGroup + for w := 0; w < writers; w++ { + want[w] = make(map[uint64]string) + wg.Add(1) + go func(w int) { + defer wg.Done() + rnd := rand.New(rand.NewSource(int64(w))) + for i := 0; i < 1500; i++ { + k := uint64(w*keys + rnd.Intn(keys)) + if rnd.Intn(3) == 0 { + if err := db.Delete(k); err == nil { + delete(want[w], k) + } + } else { + v := fmt.Sprintf("w%d-%d", w, i) + if _, err := db.Put(k, []byte(v)); err != nil { + t.Error(err) + return + } + want[w][k] = v + } + if i%15 == 0 { + time.Sleep(d) + } + } + }(w) + } + var compactions int + var cwg sync.WaitGroup + cwg.Add(1) + go func() { + defer cwg.Done() + for { + select { + case <-stop: + return + default: + } + if _, err := db.Compact(); err != nil { + t.Error(err) + return + } + compactions++ + time.Sleep(d) + } + }() + wg.Wait() + close(stop) + cwg.Wait() + + check := func(db *DB, when string) { + t.Helper() + if err := db.Verify(); err != nil { + t.Fatalf("%s: %v", when, err) + } + n := 0 + for w := 0; w < writers; w++ { + for k := uint64(w * keys); k < uint64((w+1)*keys); k++ { + v, err := db.Get(k) + if s, ok := want[w][k]; ok { + n++ + if err != nil || string(v) != s { + t.Errorf("%s: key %d reads %q, %v; want %q", when, k, v, err, s) + } + } else if err == nil { + t.Errorf("%s: key %d, deleted, reads %q", when, k, v) + } + } + } + if got := db.Size(); got != int64(n) { + t.Errorf("%s: %d keys; want %d", when, got, n) + } + } + check(db, "after the writes") + t.Logf("%d compactions", compactions) + if err := db.Close(); err != nil { + t.Fatal(err) + } + db, err = Open(WithLogFilePath(dir), WithTimeBlockInterval(d)) + if err != nil { + t.Fatal(err) + } + defer db.Close() + check(db, "reopened") +} diff --git a/memdb/meter.go b/memdb/meter.go index 818bf6e..80eb997 100644 --- a/memdb/meter.go +++ b/memdb/meter.go @@ -34,6 +34,9 @@ type Meter struct { Syncs metrics.Counter Recovers metrics.Counter Dels metrics.Counter + // Compactions counts Compact's runs, and Moves the values they moved. + Compactions metrics.Counter + Moves metrics.Counter } // NewMeter provide meter to capture statistics. @@ -47,6 +50,9 @@ func NewMeter() *Meter { Syncs: metrics.NewCounter(), Recovers: metrics.NewCounter(), Dels: metrics.NewCounter(), + + Compactions: metrics.NewCounter(), + Moves: metrics.NewCounter(), } c.TimeSeries.Time(func() {}) @@ -55,6 +61,8 @@ func NewMeter() *Meter { Metrics.GetOrRegister("Syncs", c.Syncs) Metrics.GetOrRegister("Recovers", c.Recovers) Metrics.GetOrRegister("Dels", c.Dels) + Metrics.GetOrRegister("Compactions", c.Compactions) + Metrics.GetOrRegister("Moves", c.Moves) return c } @@ -75,17 +83,20 @@ type Varz struct { Syncs int64 `json:"syncs"` Recovers int64 `json:"recovers"` Dels int64 `json:"Dels"` - HMean float64 `json:"hmean"` // Event duration harmonic mean. - P50 float64 `json:"p50"` // Event duration nth percentiles. - P75 float64 `json:"p75"` - P95 float64 `json:"p95"` - P99 float64 `json:"p99"` - P999 float64 `json:"p999"` - Long5p float64 `json:"long_5p"` // Average of the longest 5% event durations. - Short5p float64 `json:"short_5p"` // Average of the shortest 5% event durations. - Max float64 `json:"max"` // Highest event duration. - Min float64 `json:"min"` // Lowest event duration. - StdDev float64 `json:"stddev"` // Standard deviation. + // Compactions counts Compact's runs, and Moves the values they moved. + Compactions int64 `json:"compactions"` + Moves int64 `json:"moves"` + HMean float64 `json:"hmean"` // Event duration harmonic mean. + P50 float64 `json:"p50"` // Event duration nth percentiles. + P75 float64 `json:"p75"` + P95 float64 `json:"p95"` + P99 float64 `json:"p99"` + P999 float64 `json:"p999"` + Long5p float64 `json:"long_5p"` // Average of the longest 5% event durations. + Short5p float64 `json:"short_5p"` // Average of the shortest 5% event durations. + Max float64 `json:"max"` // Highest event duration. + Min float64 `json:"min"` // Lowest event duration. + StdDev float64 `json:"stddev"` // Standard deviation. } func uptime(d time.Duration) string { @@ -122,6 +133,8 @@ func (db *DB) Varz() (*Varz, error) { v.Syncs = db.internal.meter.Syncs.Count() v.Recovers = db.internal.meter.Recovers.Count() v.Dels = db.internal.meter.Dels.Count() + v.Compactions = db.internal.meter.Compactions.Count() + v.Moves = db.internal.meter.Moves.Count() ts := db.internal.meter.TimeSeries.Snapshot() v.HMean = float64(ts.HMean()) v.P50 = float64(ts.P50()) diff --git a/memdb/model_test.go b/memdb/model_test.go index 761a292..53007c2 100644 --- a/memdb/model_test.go +++ b/memdb/model_test.go @@ -188,10 +188,18 @@ func runModel(t *testing.T, seed int64, ops int) error { } op = fmt.Sprintf("Batch(%v, kept %v) in %d", keys, keep, timeID) opKeys = keys - case p < 88: + case p < 84: time.Sleep(time.Duration(rnd.Intn(int(2 * modelBlockDuration)))) op = "sleep" opKeys = nil + case p < 88: + // Compacting moves values; it changes none. + moved, err := db.Compact() + if err != nil { + return fail(fmt.Errorf("op %d: Compact: %v", i, err)) + } + op = fmt.Sprintf("Compact (moved %d)", moved) + opKeys = nil case p < 94: if err := db.Flush(); err != nil { return fail(fmt.Errorf("op %d: Flush: %v", i, err)) diff --git a/server/internal/db/unitdb/adapter.go b/server/internal/db/unitdb/adapter.go index 165cd6c..9496752 100644 --- a/server/internal/db/unitdb/adapter.go +++ b/server/internal/db/unitdb/adapter.go @@ -19,6 +19,7 @@ package adapter import ( "encoding/json" "errors" + "fmt" "io" "os" "sync" @@ -66,6 +67,8 @@ type adapter struct { wmu sync.RWMutex // lastWrite is when the DB was last written to, in unix nanoseconds. lastWrite atomic.Int64 + // stop stops the compactor (startCompactor). + stop chan struct{} // close closer io.Closer @@ -107,14 +110,59 @@ func (a *adapter) Open(path, jsonconfig string, reset bool) error { a.config = &config a.path = path + a.compact() + a.startCompactor(compactEvery) return nil } +// compactEvery is how often a running store compacts its message log. +var compactEvery = time.Minute + +// startCompactor compacts the message log every interval, until Close: a +// message kept for long, a session's row or a publish waiting on a +// subscriber away, holds its block, and the WAL keeps every log chaining +// deletes back to it (memdb's Compact). +func (a *adapter) startCompactor(interval time.Duration) { + stop := make(chan struct{}) + a.stop = stop + go func() { + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-stop: + return + case <-ticker.C: + a.compact() + } + } + }() +} + +// compact moves the messages left in mostly dead blocks, so that the blocks, +// and the WAL's logs of them, go. It is a write: a checkpoint waits for it. +func (a *adapter) compact() { + a.wmu.RLock() + defer a.wmu.RUnlock() + if a.mem == nil { + return + } + if n, err := a.mem.Compact(); err != nil { + log.Error("adapter.compact", err.Error()) + } else if n > 0 { + log.Info("adapter.compact", fmt.Sprintf("moved %d messages", n)) + } +} + // Close closes the underlying database connection func (a *adapter) Close() error { a.wmu.Lock() defer a.wmu.Unlock() + if a.stop != nil { + close(a.stop) + a.stop = nil + } var err error if a.db != nil { err = a.db.Close() diff --git a/server/internal/db/unitdb/adapter_test.go b/server/internal/db/unitdb/adapter_test.go index f0f592e..ec91c17 100644 --- a/server/internal/db/unitdb/adapter_test.go +++ b/server/internal/db/unitdb/adapter_test.go @@ -2,7 +2,9 @@ package adapter import ( "fmt" + "path/filepath" "testing" + "time" "github.com/unit-io/unitdb/memdb" ) @@ -113,3 +115,49 @@ func TestDeleteClosed(t *testing.T) { t.Errorf("Count on a closed store: %d", n) } } + +// TestCompactorFreesLogs keeps a message while others are put and deleted +// over several seconds: the kept message holds its block, and the WAL kept +// every log chaining back to it, until the compactor moved it. +func TestCompactorFreesLogs(t *testing.T) { + saved := compactEvery + compactEvery = 100 * time.Millisecond + defer func() { compactEvery = saved }() + dir := t.TempDir() + a := &adapter{} + if err := a.Open(dir, `{}`, false); err != nil { + t.Fatal(err) + } + defer a.Close() + if err := a.PutMessage(1, []byte("kept")); err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(4 * time.Second) + for i := 0; time.Now().Before(deadline); i++ { + key := uint64(i%64)<<32 | 0x40dd7a16 + if err := a.PutMessage(key, []byte("publish")); err != nil { + t.Fatal(err) + } + if err := a.DeleteMessage(key); err != nil { + t.Fatal(err) + } + if err := a.PutMessage(2, []byte(fmt.Sprintf("row-%d", i))); err != nil { + t.Fatal(err) + } + time.Sleep(5 * time.Millisecond) + } + var logs []string + for wait := time.Now().Add(3 * time.Second); time.Now().Before(wait); time.Sleep(100 * time.Millisecond) { + logs, _ = filepath.Glob(filepath.Join(dir, "logs", "*.log")) + if len(logs) <= 20 { + break + } + } + if len(logs) > 20 { + t.Fatalf("%d logs after the churn, with the compactor running", len(logs)) + } + if got, err := a.GetMessage(1); err != nil || string(got) != "kept" { + t.Fatalf("the kept message reads %q, %v", got, err) + } + t.Logf("%d logs", len(logs)) +} From 55683b7d4264bad3e0e285f956b1c6ad5a076ca2 Mon Sep 17 00:00:00 2001 From: unit-adm <314923187+mumtaz6@users.noreply.github.com> Date: Tue, 6 Oct 2026 14:53:49 +0530 Subject: [PATCH 2/2] memdb: keep each round of TestCompactLeavesMostlyLiveBlocks in one block The test failed now and then (1 in 30 runs, and in CI): "Compact moved 1; want none". A write goes to the block of the log's last rotation, every log interval, not to the block of the time it is made; a round of writes made as a block began could put its first writes in the block before. Split between two blocks, a round left one mostly dead, which Compact rightly moved. Each round now starts a few log intervals into a new block, with blocks of 50 ms, far longer than a round. 200 runs pass. Co-Authored-By: Claude Opus 5.5 (1M context) --- memdb/compact_test.go | 14 +++++++++++--- 1 file changed, 11 insertions(+), 3 deletions(-) diff --git a/memdb/compact_test.go b/memdb/compact_test.go index 842d366..ed3804f 100644 --- a/memdb/compact_test.go +++ b/memdb/compact_test.go @@ -111,15 +111,22 @@ func TestCompactFreesPinnedBlocks(t *testing.T) { } // TestCompactLeavesMostlyLiveBlocks compacts a store whose blocks are mostly -// live: nothing moves. +// live: nothing moves. Each round of writes goes in one block: it starts once +// the log has rotated into a new block, and takes far less than a block. A +// write goes to the block of the log's last rotation, every log interval, so +// a round that started as a block began could put its first writes in the +// block before: split, it could leave one mostly dead, and rightly compacted. func TestCompactLeavesMostlyLiveBlocks(t *testing.T) { - const d = 10 * time.Millisecond + const d = 50 * time.Millisecond db, err := Open(WithLogFilePath(t.TempDir()), WithTimeBlockInterval(d), WithLogInterval(2*time.Millisecond)) if err != nil { t.Fatal(err) } defer db.Close() for i := 0; i < 5; i++ { + // Into the next block, a few log intervals after it begins. + now := time.Now() + time.Sleep(now.Truncate(d).Add(d + 5*time.Millisecond).Sub(now)) for k := 0; k < 10; k++ { if _, err := db.Put(uint64(i*10+k), []byte("value")); err != nil { t.Fatal(err) @@ -129,8 +136,9 @@ func TestCompactLeavesMostlyLiveBlocks(t *testing.T) { if err := db.Delete(uint64(i * 10)); err != nil { t.Fatal(err) } - time.Sleep(2 * d) } + // The last round's block in the past too. + time.Sleep(2 * d) if moved, err := db.Compact(); err != nil || moved != 0 { t.Fatalf("Compact moved %d, %v; want none", moved, err) }