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..ed3804f --- /dev/null +++ b/memdb/compact_test.go @@ -0,0 +1,244 @@ +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. 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 = 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) + } + } + // One of ten deleted: the block stays mostly live. + if err := db.Delete(uint64(i * 10)); err != nil { + t.Fatal(err) + } + } + // 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) + } +} + +// 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)) +}