diff --git a/.omo/plans/fix-c2-c3-wal-recovery.md b/.omo/plans/fix-c2-c3-wal-recovery.md new file mode 100644 index 0000000..5cf48b4 --- /dev/null +++ b/.omo/plans/fix-c2-c3-wal-recovery.md @@ -0,0 +1,641 @@ +# C→A 修复方案:C2+C3 数据安全 Patch + +## TL;DR + +> **目标**:消除 Phase 1 WAL recovery 的两个数据丢失路径(C2 + C3),按 Oracle 修订意见先更新审核报告,再统一修复。 +> +> **交付**: +> - 更新 `docs/audit-3.2.md`,加入 Oracle 的 4 条修订意见 +> - 删除 `wal/recover.go` 中 2 处 `manifest.Save` 调用 + 1 处 CURRENT fallback +> - 重写 1 个测试 + 新增 3 个测试 +> - 单次 commit 提交 +> +> **预估工时**:1.5-3 小时 +> **风险**:低(纯删除 + 测试调整,无新逻辑) + +--- + +## Context + +### 为什么 C2 和 C3 必须一起修 + +Oracle 验证(`bg_ef425776`)发现一个我原审核漏掉的关键事实: + +- Phase 1 没有 flush,**MANIFEST 在正常运行中始终是 0** +- `resolveRecoverySegmentID` 看到 MANIFEST=0 就 fallback 到 CURRENT +- `segment_manager.go:47,89` 每次 create/rotate segment 都顺手写 CURRENT +- 所以 Phase 1 默认状态下,首次 Recover 就会从 CURRENT 指向的 active segment 开始,跳过更早的 segment + +→ 单独修 C3(删 `manifest.Save`)没用,C2 路径照样丢数据。 + +### Oracle 的其他修订 + +1. **C4** 多一个失败模式:`Recover` 总是 truncate `segments[last]`,但 corruption 可能出现在非尾段,导致截断错误的 segment +2. **M1**(findValidOffset)应升级为 High:只校验物理记录不跟踪 batch 边界,对 `First+Middle*` 没 `Last` 的情况返回错误的截断点 +3. **测试方向**:连续两次 Recover 的"幂等"测试不能一刀切,要分干净 WAL 和尾部损坏两种场景 + +--- + +## 执行计划 + +### Phase C:更新审核报告(15 分钟) + +**目标**:把 Oracle 的修订固化到 `docs/audit-3.2.md`,避免下次重新审计时遗漏。 + +#### C.1 修订 C2 条目 + +在 C2 的"后果"段后追加: + +> **Oracle 修订(bg_ef425776)**:C2 在 Phase 1 比 C3 更严重。Phase 1 没有 flush,MANIFEST 一直保持 0;`segment_manager.go:47,89` 每次创建/轮转 segment 都会写 CURRENT 指向 active segment。**默认状态下首次 Recover 就会触发 fallback**,不需要"CURRENT 缺失/落后"这种特殊场景。修 C3 之前必须先修 C2,否则数据丢失窗口依然存在。 + +#### C.2 修订 C3 条目 + +在 C3 的"修复"段后追加: + +> **Oracle 修订(bg_ef425776)**:单修 C3 不足以解决 Phase 1 数据丢失,必须和 C2 一起修。测试方向需要区分干净 WAL 和尾部损坏 WAL 两种幂等性。 + +#### C.3 修订 C4 条目 + +在 C4 的"后果"段后追加: + +> **Oracle 修订(bg_ef425776)**:C4 还有一个更严重的失败模式。`wal/recover.go:61` 总是对 `segments[len(segments)-1]` 调用 `truncateSegment`,但 `RecoverFromSegments` 的 `TailCorruptionError` 可能来自非尾段。结果:真正损坏的 segment 不动,最后一段的有效数据被错误截掉。 + +#### C.4 升级 M1 → H8 + +把 M1(`findValidOffset` 与原始解析不一致)从 Medium 升到 High,改编号 H8,加入具体证据: + +> **Oracle 修订(bg_ef425776)**:升级为 High。`findValidOffset` 只用 `DecodePhysicalRecord` 校验物理记录,**不跟踪完整 batch 边界**。对尾部 `First + Middle*` 没 `Last` 的情况,`recovery.go:133` 在最后一个完整 batch 结尾报尾部损坏,而 `findValidOffset` 可能返回不完整 fragment 之后的 EOF,导致截断点错位、残留半截 fragment、下次启动反复 repair。 + +#### C.5 修订优先级表 + +把"必须立刻修的"表格更新为: + +| 编号 | 一句话 | 风险等级 | Phase 1 是否默认触发 | +|------|--------|----------|---------------------| +| C1 | CRC 多项式错 | 格式不符 | N/A | +| **C2** | Recovery CURRENT 兜底 | **恢复起点错误** | **是(默认触发)** | +| **C3** | Recovery 写 MANIFEST | 数据丢失 | 是(修了 C2 才能完全止血) | +| C4 | 非尾段 fragment 当尾段损坏 | 截断错误 segment | 偶发 | +| C5 | 截断无 fsync | DB 状态不一致 | 偶发 | +| C6 | dir fsync 静默 | 已确认写入消失 | 偶发 | +| C7 | Put/Close 竞态 | 偶发 panic | 偶发 | + +并更新底部优先级: + +> 修订后优先级:**C2 + C3(一起)→ C6 → C5 + H8(一起,同一文件)→ C4 → C1 → C7** + +#### C.6 验证 + +```bash +# 报告应该读起来前后一致,无矛盾 +grep -c "Oracle 修订" docs/audit-3.2.md # 期望 4 +``` + +--- + +### Phase A:修复 C2+C3(45-90 分钟) + +**目标**:消除 Phase 1 默认状态下的两条数据丢失路径。 + +#### A.1 修改 `wal/recover.go` + +**变更 1:删除两处 `manifest.Save` 调用** + +文件:`wal/recover.go` + +删除尾部损坏路径(约 line 76-79): +```go +// 删除: +// Update MANIFEST with new recovery state. +if saveErr := manifest.Save(dir, result.NextSegmentID); saveErr != nil { + return nil, fmt.Errorf("wal: recover: save manifest after truncation: %w", saveErr) +} +``` + +删除成功路径(约 line 99-102): +```go +// 删除: +// Update MANIFEST with new recovery state. +if saveErr := manifest.Save(dir, result.NextSegmentID); saveErr != nil { + return nil, fmt.Errorf("wal: recover: save manifest: %w", saveErr) +} +``` + +**变更 2:删除 CURRENT fallback** + +文件:`wal/recover.go`,函数 `resolveRecoverySegmentID` + +修改前: +```go +func resolveRecoverySegmentID(dir string) (uint64, error) { + mf, err := manifest.Load(dir) + if err != nil { + return 0, fmt.Errorf("load manifest: %w", err) + } + if mf.RecoverySegmentID > 0 { + return mf.RecoverySegmentID, nil + } + + // MANIFEST had 0 (fresh DB or not yet written). Try CURRENT. + if segID, ok := manifest.ReadCurrent(dir); ok { + return segID, nil + } + + return 0, nil +} +``` + +修改后: +```go +// resolveRecoverySegmentID returns the recovery start segment ID from MANIFEST. +// MANIFEST is the only authoritative source of recovery start per design §3.2 +// line 600-06. CURRENT is a write-side hint and must NOT be used here. +func resolveRecoverySegmentID(dir string) (uint64, error) { + mf, err := manifest.Load(dir) + if err != nil { + return 0, fmt.Errorf("load manifest: %w", err) + } + return mf.RecoverySegmentID, nil +} +``` + +**变更 3:更新 import 注释** + +把 `recover.go:27` 注释从 `// Step 1: Determine recovery segment ID from MANIFEST or CURRENT.` 改成 `// Step 1: Determine recovery segment ID from MANIFEST.` + +**变更 4:检查 manifest import 是否还需要** + +`wal/recover.go` 还使用 `manifest.Load`,所以 import 保留。 + +#### A.2 检查 manifest.WriteCurrent 是否还有调用方 + +```bash +grep -rn "manifest.WriteCurrent\|manifest.ReadCurrent" --include="*.go" +``` + +预期: +- `wal/segment_manager.go:48,90` 仍然调用 `WriteCurrent`(合法,CURRENT 作为写入侧 hint) +- `manifest/current.go` 自身定义 +- `wal/recover.go` 之前调用 `ReadCurrent` 的地方已删除 + +→ CURRENT 文件保留,只是 recovery 不再读它。符合设计 §3.2 line 276。 + +#### A.3 修改 `wal/recover_test.go` + +> **Oracle 修订(bg_2e86d33b)**: +> - BLOCKING:原计划的可选 e2e `TestOpenTwiceKeepsData`(两次 open)不足以验证 C3。C3 数据丢失发生在**第二次 restart**,需要**三次 open**才能抓到。已升级为必选。 +> - NICE-TO-HAVE:注释中"NewSegmentWriter writes CURRENT"是错的,只有 SegmentManager 写。修正注释。 +> - NICE-TO-HAVE:MANIFEST 已存在的 subcase,防止"覆盖已有 MANIFEST"的回归。 +> - 注意:测试中避免强制多 segment 轮转,因为 `segment_manager.go:64-66` 有无关 bug(C8,见审核报告),多 segment recovery 会失败。 + +**变更 1:重写 `TestRecoverUpdatesManifest` → `TestRecoverDoesNotUpdateManifest`** + +改名 `TestRecoverDoesNotUpdateManifest`,断言相反: + +```go +func TestRecoverDoesNotUpdateManifest(t *testing.T) { + dir := t.TempDir() + + // Write test data. NewSegmentWriter creates the .wal file via the + // durable-ready protocol but does NOT write CURRENT (only SegmentManager + // does) and does NOT write MANIFEST. + writeTestSegment(t, dir, 0, 50, [][]*WalEntry{ + {makePutEntry("a", "b")}, + {makePutEntry("c", "d")}, + }) + + // Snapshot MANIFEST state before recovery. For a fresh DB, MANIFEST does + // not exist on disk; manifest.Load returns a zero-value Manifest. + beforeMF, err := manifest.Load(dir) + if err != nil { + t.Fatalf("manifest.Load before recover: %v", err) + } + beforeExists := fileExists(t, filepath.Join(dir, "MANIFEST")) + + replayer := &mockReplayer{} + result, err := Recover(dir, replayer) + if err != nil { + t.Fatalf("Recover: %v", err) + } + + // MANIFEST must be unchanged. + afterMF, err := manifest.Load(dir) + if err != nil { + t.Fatalf("manifest.Load after recover: %v", err) + } + if afterMF.RecoverySegmentID != beforeMF.RecoverySegmentID { + t.Errorf("MANIFEST RecoverySegmentID changed: %d → %d", + beforeMF.RecoverySegmentID, afterMF.RecoverySegmentID) + } + afterExists := fileExists(t, filepath.Join(dir, "MANIFEST")) + if beforeExists != afterExists { + t.Errorf("MANIFEST file existence changed: before=%v after=%v", + beforeExists, afterExists) + } + + // RecoveryResult.NextSegmentID is in-memory only — not persisted. + _ = result // NextSegmentID is allowed to differ from MANIFEST.RecoverySegmentID. +} + +// Subtest: MANIFEST already exists (e.g. from a prior flush in future phases). +// Recovery must not overwrite or delete it. Catches accidental writes that +// fresh-DB subtest above cannot detect (since fresh DB has no MANIFEST). +func TestRecoverPreservesExistingManifest(t *testing.T) { + dir := t.TempDir() + + writeTestSegment(t, dir, 0, 50, [][]*WalEntry{ + {makePutEntry("a", "b")}, + }) + + // Simulate a prior checkpoint having advanced MANIFEST to segment 0. + // (Even though Phase 1 has no flush, future phases will. This test + // guards the invariant going forward.) + if err := manifest.Save(dir, 0); err != nil { + t.Fatalf("manifest.Save setup: %v", err) + } + beforeBytes, err := os.ReadFile(filepath.Join(dir, "MANIFEST")) + if err != nil { + t.Fatalf("ReadFile MANIFEST: %v", err) + } + + replayer := &mockReplayer{} + if _, err := Recover(dir, replayer); err != nil { + t.Fatalf("Recover: %v", err) + } + + afterBytes, err := os.ReadFile(filepath.Join(dir, "MANIFEST")) + if err != nil { + t.Fatalf("ReadFile MANIFEST after recover: %v", err) + } + if !bytes.Equal(beforeBytes, afterBytes) { + t.Errorf("MANIFEST bytes changed:\n before=%q\n after=%q", + string(beforeBytes), string(afterBytes)) + } +} + +func fileExists(t *testing.T, path string) bool { + t.Helper() + _, err := os.Stat(path) + if err == nil { + return true + } + if os.IsNotExist(err) { + return false + } + t.Fatalf("stat %s: %v", path, err) + return false +} +``` + +**变更 2:新增 `TestRecoverIdempotentClean`** + +```go +// TestRecoverIdempotentClean verifies that recovering a clean WAL twice +// produces identical results. This is the "no MANIFEST writes" guarantee +// from design §3.2 line 280. +func TestRecoverIdempotentClean(t *testing.T) { + dir := t.TempDir() + + writeTestSegment(t, dir, 0, 0, [][]*WalEntry{ + {makePutEntry("k1", "v1")}, + {makePutEntry("k2", "v2")}, + }) + + replayer1 := &mockReplayer{} + result1, err := Recover(dir, replayer1) + if err != nil { + t.Fatalf("Recover (1st): %v", err) + } + + replayer2 := &mockReplayer{} + result2, err := Recover(dir, replayer2) + if err != nil { + t.Fatalf("Recover (2nd): %v", err) + } + + if result1.NextSequence != result2.NextSequence { + t.Errorf("NextSequence differs: %d vs %d", result1.NextSequence, result2.NextSequence) + } + if result1.NextSegmentID != result2.NextSegmentID { + t.Errorf("NextSegmentID differs: %d vs %d", result1.NextSegmentID, result2.NextSegmentID) + } + if result1.Truncated || result2.Truncated { + t.Errorf("Truncated should be false for clean WAL: r1=%v r2=%v", + result1.Truncated, result2.Truncated) + } + if !reflect.DeepEqual(replayer1.puts, replayer2.puts) { + t.Errorf("replayed puts differ:\n r1=%#v\n r2=%#v", replayer1.puts, replayer2.puts) + } +} +``` + +**变更 3:新增 `TestRecoverIdempotentAfterTruncation`** + +```go +// TestRecoverIdempotentAfterTruncation verifies that after the first recovery +// truncates a corrupted tail, the second recovery sees a clean WAL with the +// same NextSequence. Per Oracle: the first call reports Truncated=true, the +// second reports Truncated=false but identical replay state. +func TestRecoverIdempotentAfterTruncation(t *testing.T) { + dir := t.TempDir() + + filePath := writeTestSegment(t, dir, 0, 100, [][]*WalEntry{ + {makePutEntry("good1", "before-corruption")}, + {makePutEntry("good2", "also-before")}, + }) + appendFileBytes(t, filePath, []byte{0xDE, 0xAD, 0xBE, 0xEF}) + + // First recovery: detects tail corruption, truncates, returns Truncated=true. + replayer1 := &mockReplayer{} + result1, err := Recover(dir, replayer1) + if err != nil { + t.Fatalf("Recover (1st): %v", err) + } + if !result1.Truncated { + t.Fatal("1st Recover: Truncated = false, want true") + } + + // Second recovery: file has been truncated, no corruption remains. + replayer2 := &mockReplayer{} + result2, err := Recover(dir, replayer2) + if err != nil { + t.Fatalf("Recover (2nd): %v", err) + } + if result2.Truncated { + t.Error("2nd Recover: Truncated = true, want false (tail already repaired)") + } + if result1.NextSequence != result2.NextSequence { + t.Errorf("NextSequence differs: %d vs %d", result1.NextSequence, result2.NextSequence) + } + if !reflect.DeepEqual(replayer1.puts, replayer2.puts) { + t.Errorf("replayed puts differ:\n r1=%#v\n r2=%#v", replayer1.puts, replayer2.puts) + } +} +``` + +**变更 4:新增 `TestRecoverIgnoresCurrentFallback`**(C2 回归测试) + +> **Momus 修订(bg_3dfb54a9)**:`writeTestSegment` 走 `NewSegmentWriter`,**不写 CURRENT**(只有 `SegmentManager` 写)。测试必须**显式**调用 `manifest.WriteCurrent` 模拟 Phase 1 默认状态。 + +```go +// TestRecoverIgnoresCurrentFallback verifies that CURRENT is not used as a +// recovery start fallback when MANIFEST.RecoverySegmentID == 0. Per design +// §3.2 line 604-06, CURRENT is only a write-side hint and must not affect +// recovery start. Without this guarantee, Phase 1 default state (MANIFEST=0, +// CURRENT pointing to active segment) causes recovery to skip older segments. +// +// Note: writeTestSegment uses NewSegmentWriter directly, which does NOT write +// CURRENT (only SegmentManager does). We must write CURRENT explicitly to +// simulate the Phase 1 default state where SegmentManager has been rotating +// segments. +func TestRecoverIgnoresCurrentFallback(t *testing.T) { + dir := t.TempDir() + + // Write 3 segments with valid batches. + writeTestSegment(t, dir, 0, 0, [][]*WalEntry{ + {makePutEntry("seg0-k1", "v1")}, + }) + writeTestSegment(t, dir, 1, 1, [][]*WalEntry{ + {makePutEntry("seg1-k1", "v1")}, + }) + writeTestSegment(t, dir, 2, 2, [][]*WalEntry{ + {makePutEntry("seg2-k1", "v1")}, + }) + + // Simulate Phase 1 default: SegmentManager has been rotating, so CURRENT + // exists and points to the last segment (2). MANIFEST still does not exist + // because no flush has happened yet. + if err := manifest.WriteCurrent(dir, 2); err != nil { + t.Fatalf("WriteCurrent: %v", err) + } + + // Sanity: CURRENT exists and points to segment 2. + currentSegID, ok := manifest.ReadCurrent(dir) + if !ok || currentSegID != 2 { + t.Fatalf("CURRENT setup wrong: segID=%d ok=%v", currentSegID, ok) + } + // Sanity: MANIFEST.RecoverySegmentID == 0 (fresh DB). + mf, err := manifest.Load(dir) + if err != nil { + t.Fatalf("Load: %v", err) + } + if mf.RecoverySegmentID != 0 { + t.Fatalf("MANIFEST.RecoverySegmentID = %d, want 0", mf.RecoverySegmentID) + } + + // Recover must start from segment 0, not 2. + replayer := &mockReplayer{} + result, err := Recover(dir, replayer) + if err != nil { + t.Fatalf("Recover: %v", err) + } + + // All 3 entries must be replayed. + wantPuts := []replayPut{ + {key: "seg0-k1", value: "v1", seq: 0}, + {key: "seg1-k1", value: "v1", seq: 1}, + {key: "seg2-k1", value: "v1", seq: 2}, + } + if !reflect.DeepEqual(replayer.puts, wantPuts) { + t.Errorf("puts = %#v, want %#v", replayer.puts, wantPuts) + } + if result.NextSequence != 3 { + t.Errorf("NextSequence = %d, want 3", result.NextSequence) + } + if result.NextSegmentID != 3 { + t.Errorf("NextSegmentID = %d, want 3", result.NextSegmentID) + } +} +``` + +#### A.4 验证步骤 + +按顺序执行: + +```bash +# 1. 编译通过 +go build ./... + +# 2. wal 包测试全绿 +go test ./wal/... -count=1 -v + +# 3. 全仓测试不回归 +go test ./... -count=1 + +# 4. vet +go vet ./... + +# 5. lint(如果环境装了 golangci-lint;没装就跳过,不阻塞) +if command -v golangci-lint >/dev/null 2>&1; then + golangci-lint run ./wal/... ./manifest/... +else + echo "golangci-lint not installed, skipping" +fi + +# 6. 重点跑这次新增的 4+1 个测试(Oracle 建议) +go test ./wal -run 'TestRecover(DoesNotUpdateManifest|PreservesExistingManifest|IdempotentClean|IdempotentAfterTruncation|IgnoresCurrentFallback)$' -count=1 -v +go test . -run 'TestOpenThreeTimesKeepsData$' -count=1 -v + +# 7. race 检测(Oracle 建议,wal 包有 goroutine) +go test -race ./wal/... -count=1 +``` + +#### A.5 必选 e2e:`TestOpenThreeTimesKeepsData`(Oracle 升级) + +> **Oracle 修订(bg_2e86d33b)BLOCKING**:原计划 `TestOpenTwiceKeepsData` 两次 open 不足以验证 C3。C3 的数据丢失发生在**第二次 restart**(第一次 recovery 写了 MANIFEST,第二次 recovery 才会跳过 segment)。必须用**三次 open**。 + +放在 `db_test.go`(package `go_kv`): + +```go +// TestOpenThreeTimesKeepsData is the end-to-end regression for C2+C3. +// +// C3's data-loss bug manifests on the SECOND restart after writes: +// - Open 1: write data, close. (Phase 1: no MANIFEST write yet from flush.) +// - Open 2: recovery (buggy code) writes MANIFEST=NextSegmentID, then +// opens new writer. Data still visible because the in-memory memtable +// was rebuilt from WAL. +// - Open 3: recovery reads advanced MANIFEST, skips old segments, +// data NOT replayed → data permanently invisible. +// +// Two opens cannot catch this; three opens can. +// +// Single-segment only: do NOT force rotation, because segment_manager +// has an unrelated C8 bug (passes byte offset as startSequence) that +// breaks multi-segment recovery. That bug is tracked separately. +func TestOpenThreeTimesKeepsData(t *testing.T) { + dir := t.TempDir() + + cfg := config.Defaults() + // Default MaxSegmentSize=64MB is plenty for these writes; no rotation. + + // Open 1: write keys. + db1, err := Open(dir, &cfg) + if err != nil { + t.Fatalf("Open 1: %v", err) + } + keys := []string{"k1", "k2", "k3", "k4", "k5"} + for _, k := range keys { + if err := db1.Put([]byte(k), []byte("v-"+k)); err != nil { + t.Fatalf("Put %s: %v", k, err) + } + } + if err := db1.Close(); err != nil { + t.Fatalf("Close 1: %v", err) + } + + // Open 2: verify, then close. (Buggy code: this is where MANIFEST + // gets advanced. Fixed code: MANIFEST stays unchanged.) + db2, err := Open(dir, &cfg) + if err != nil { + t.Fatalf("Open 2: %v", err) + } + for _, k := range keys { + r := db2.Get([]byte(k)) + if !r.Found { + t.Errorf("Open 2: key %s not found", k) + } + } + if err := db2.Close(); err != nil { + t.Fatalf("Close 2: %v", err) + } + + // Open 3: verify again. This is where C3's data loss would manifest + // on buggy code (MANIFEST was advanced in Open 2, recovery now skips + // the original segments). + db3, err := Open(dir, &cfg) + if err != nil { + t.Fatalf("Open 3: %v", err) + } + defer db3.Close() + for _, k := range keys { + r := db3.Get([]byte(k)) + if !r.Found { + t.Errorf("Open 3: key %s not found (C3 regression: data lost)", k) + } + } +} +``` + +#### A.6 Commit + +单次 commit,message: + +``` +fix(wal): C2+C3 stop recovery from advancing MANIFEST or using CURRENT + +Phase 1 default state had two data-loss paths in WAL recovery: + +1. (C2) resolveRecoverySegmentID fell back to CURRENT when MANIFEST=0. + Since segment_manager writes CURRENT on every segment create/rotate, + the first recovery in Phase 1 (MANIFEST always 0 without flush) would + start from the active segment, skipping earlier unflushed segments. + +2. (C3) Recover called manifest.Save after every recovery, advancing + recoverySegmentID past segments that were still the only durable copy + of their data (no SSTable flush yet). Next restart would filter those + segments out and permanently lose the data. + +Per design §3.2 line 280, recovery must not update MANIFEST; per line 604-06, +CURRENT must not be used as recovery start. Both fixes are required together +— fixing C3 alone leaves C2's data-loss window open. + +Changes: +- wal/recover.go: remove manifest.Save calls on both success and tail-repair + paths; remove CURRENT fallback in resolveRecoverySegmentID. + RecoveryResult.NextSegmentID is now in-memory only (consumed by DB.Open + to seed the new writer, but never persisted to MANIFEST). +- wal/recover_test.go: rewrite TestRecoverUpdatesManifest as + TestRecoverDoesNotUpdateManifest; add TestRecoverPreservesExistingManifest, + TestRecoverIdempotentClean, TestRecoverIdempotentAfterTruncation, + TestRecoverIgnoresCurrentFallback. +- db_test.go: add TestOpenThreeTimesKeepsData (three opens to catch C3's + second-restart data loss). + +Refs docs/audit-3.2.md C2+C3 (with Oracle revisions from bg_ef425776 and +bg_2e86d33b). +``` + +--- + +## 验收清单 + +- [ ] Phase C:`docs/audit-3.2.md` 包含 4 处"Oracle 修订"标注 +- [ ] Phase A.1:`wal/recover.go` 中 `manifest.Save` 出现次数 = 0 +- [ ] Phase A.1:`wal/recover.go` 中 `manifest.ReadCurrent` 出现次数 = 0 +- [ ] Phase A.3:5 个测试名(1 改 + 4 新)全部存在 +- [ ] Phase A.5:`TestOpenThreeTimesKeepsData`(必选)通过 +- [ ] `go test ./wal/... -count=1` 全绿 +- [ ] `go test . -count=1` 全绿(含 db_test.go) +- [ ] `go test ./... -count=1` 全绿 +- [ ] `go test -race ./wal/... -count=1` 全绿 +- [ ] `go vet ./...` 无新增警告 +- [ ] 单次 commit,message 引用 audit C2+C3 + +--- + +## 不在本次范围内(后续 issue) + +| 编号 | 为什么不放进来 | +|------|---------------| +| **C8(新)** | `segment_manager.go:64-66` 把 `active.CurrentOffset()`(字节偏移)当 `newStartSequence` 传给 rotate。多 segment recovery 会因 startSequence 不匹配而失败。和 C2+C3 完全独立,但相关测试必须避免强制轮转 | +| C4 | 涉及 `RecoverFromSegments` 接口变更(isLast 参数),改动面更大,独立做 | +| C5 + H8 | 都在 truncation 路径,应该一起做(fsync + 删空 segment + dir fsync + findValidOffset batch 边界),但和 C2+C3 不耦合 | +| C6 | 涉及 segment_writer + segment_manager,独立做 | +| C1 | CRC 改动跨多个文件,需要测试向量,独立做 | +| C7 | 并发竞态,独立做 | + +--- + +## 修订记录 + +- **v1(原始)**:C→A 方案初稿,送 Momus 审 +- **v1.1(Momus 修订 bg_3dfb54a9)**: + - Blocking:`TestRecoverIgnoresCurrentFallback` setup 假设 `writeTestSegment` 会写 CURRENT,但实际只有 `SegmentManager` 写。修订:测试中显式调用 `manifest.WriteCurrent(dir, 2)` 模拟 Phase 1 默认状态 + - Minor:`golangci-lint run ... 2>/dev/null || true` 太宽容,改为 `command -v` 探测,没装就跳过 +- **v1.2(Oracle 修订 bg_2e86d33b)**: + - Blocking:`TestOpenTwiceKeepsData`(两次 open)不足以验证 C3,C3 数据丢失发生在第二次 restart。升级为必选 `TestOpenThreeTimesKeepsData`(三次 open)。测试中避免强制多 segment 轮转,因为 C8 会让多 segment recovery 失败 + - 新增 `TestRecoverPreservesExistingManifest` 子测试,覆盖"覆盖已有 MANIFEST"的回归 + - 修正 `TestRecoverDoesNotUpdateManifest` 的注释(NewSegmentWriter 不写 CURRENT) + - 加重点测试命令 + `-race` + - commit message 补一句 `RecoveryResult.NextSegmentID` 现在是 in-memory + - 发现无关 bug C8(segment_manager 轮转 startSequence 错),加入"不在本次范围"表 diff --git a/db_test.go b/db_test.go index d2800f8..d3ca253 100644 --- a/db_test.go +++ b/db_test.go @@ -133,6 +133,58 @@ func TestDBMultiplePuts(t *testing.T) { } } +// Regression guard for C2+C3: data must survive three Open/Close cycles. +// C3's data-loss bug manifests on the SECOND restart after writes — two +// opens cannot catch it. +// +// Single-segment only: do NOT force rotation, because segment_manager has +// an unrelated C8 bug (passes byte offset as startSequence) that breaks +// multi-segment recovery. That bug is tracked separately. +func TestOpenThreeTimesKeepsData(t *testing.T) { + dir := t.TempDir() + + keys := []string{"k1", "k2", "k3", "k4", "k5"} + + db1, err := Open(dir, nil) + if err != nil { + t.Fatalf("Open 1: %v", err) + } + for _, k := range keys { + if err := db1.Put([]byte(k), []byte("v-"+k)); err != nil { + t.Fatalf("Put %s: %v", k, err) + } + } + if err := db1.Close(); err != nil { + t.Fatalf("Close 1: %v", err) + } + + verify := func(label string, db *DB) { + t.Helper() + for _, k := range keys { + r := db.Get([]byte(k)) + if !r.Found { + t.Errorf("%s: key %s not found", label, k) + } + } + } + + db2, err := Open(dir, nil) + if err != nil { + t.Fatalf("Open 2: %v", err) + } + verify("Open 2", db2) + if err := db2.Close(); err != nil { + t.Fatalf("Close 2: %v", err) + } + + db3, err := Open(dir, nil) + if err != nil { + t.Fatalf("Open 3: %v", err) + } + defer db3.Close() + verify("Open 3", db3) +} + func openTestDB(t *testing.T, dir string) *DB { t.Helper() diff --git a/docs/audit-3.2.md b/docs/audit-3.2.md new file mode 100644 index 0000000..71e9ca4 --- /dev/null +++ b/docs/audit-3.2.md @@ -0,0 +1,298 @@ +# WAL 代码审核报告(对照 `docs/design.md` § 3.2) + +审核范围:`/wal/...`、`/memtable/...`、`/manifest/...`、`/config/...`、`db.go`。 +审核时间:基于 commit `e34de4a`(fix: implement group commit collection window)。 + +下面按严重程度排列,每条给出 **设计要求 → 代码现状 → 后果 → 修复建议**。 + +--- + +## 🔴 Critical(破坏持久化 / 恢复正确性,必须修) + +### C1. CRC 多项式错误:用 IEEE 而非 crc32c + +- **设计要求**(§3.2 Physical Record、WAL File Header):crc32c(Castagnoli,0x82F63B78),用于识别 torn write / partial write。 +- **代码现状**: + - `wal/record.go:30,55` → `crc32.ChecksumIEEE` + - `wal/header.go:44,83` → `crc32.ChecksumIEEE` +- **后果**:CRC 校验值与设计文档不一致;如果未来要做跨实现兼容(其他客户端 / 工具按 crc32c 校验),所有 segment 都会被判损坏。当前自洽但与规范脱钩。 +- **修复**:改用 `hash/crc32.Castagnoli`(即 `crc32.MakeTable(crc32.Castagnoli)`),所有 `ChecksumIEEE` 替换为 `Checksum(data, castagnoliTable)`。常数需要加测试固定。 + +--- + +### C2. Recovery 用 CURRENT 作为兜底,违反"MANIFEST 是唯一权威" + +- **设计要求**(§3.2 CURRENT / MANIFEST 权威性,line 600-66): + > MANIFEST 是 recovery 起点和 checkpoint 状态的权威源……CURRENT 只表示写入侧上次尝试记录的 active WAL segment hint……**不能作为 recovery 起点、终点或排除 segment 的依据**。 +- **代码现状**:`wal/recover.go:109-124` + ```go + if mf.RecoverySegmentID > 0 { return mf.RecoverySegmentID, nil } + if segID, ok := manifest.ReadCurrent(dir); ok { return segID, nil } // ← 违规 + return 0, nil + ``` +- **后果**:CURRENT 是 best-effort 写入、可能落后 / 指向已被截断或未 durable-ready 的 segment。用它做 recovery 起点,要么漏恢复(如果它指向比 MANIFEST 更新的 segment,而那个 segment 实际上没 durable-ready),要么恢复出空集(如果它指向已被 MANIFEST 覆盖删除的旧 segment)。这条路径在 MANIFEST.RecoverySegmentID==0 时(首次创建 DB 或刚 flush 后尚未写过 MANIFEST)会被触发。 +- **修复**:删掉 CURRENT 兜底,MANIFEST.RecoverySegmentID==0 时直接返回 0;CURRENT 仅在写入侧作为 hint 顺手写。 + +> **Oracle 修订(bg_ef425776)**:C2 在 Phase 1 比 C3 更严重。Phase 1 没有 flush,MANIFEST 一直保持 0;`segment_manager.go:47,90` 每次创建/轮转 segment 都会写 CURRENT 指向 active segment。**默认状态下首次 Recover 就会触发 fallback**,不需要"CURRENT 缺失/落后"这种特殊场景。修 C3 之前必须先修 C2,否则数据丢失窗口依然存在。 + +--- + +### C3. Recovery 主动写 MANIFEST,违反"recovery repair 不更新 MANIFEST" + +- **设计要求**(§3.2 line 280): + > MANIFEST 只作为 checkpoint / recovery 起点元数据;**recovery repair 不更新 MANIFEST**,也不通过 MANIFEST 记录恢复终点。 +- **代码现状**:`wal/recover.go:77,100` + ```go + if saveErr := manifest.Save(dir, result.NextSegmentID); saveErr != nil { ... } + ``` + 无论恢复成功还是尾部截断,都会把 MANIFEST.RecoverySegmentID 改写为"最后一个 segment + 1"。 +- **后果**:**直接数据丢失**。设想 segment-5 是 active,写了若干 batch 后崩溃,恢复时只重放了部分 batch 并把尾部截断。代码把 MANIFEST 推进到 segment-6。**下次启动时 recovery 直接从 segment-6 开始,跳过了 segment-5 中已经 durable 的 batch**。设计明确要求 MANIFEST 只能由 checkpoint(MemTable flush 完成 + SSTable 元数据落盘)推进。 +- **修复**:删掉 recovery 中的 `manifest.Save`。MANIFEST 推进只能发生在 flush 完成后(未来 §3.3 实现)。 + +> **Oracle 修订(bg_ef425776)**:单修 C3 不足以解决 Phase 1 数据丢失,必须和 C2 一起修。测试方向需要区分干净 WAL 和尾部损坏 WAL 两种幂等性。 + +--- + +### C4. 非尾段 fragment 状态被当成尾段损坏 + +- **设计要求**(§3.2 Fragment 状态机 line 731): + > Segment 边界不是合法的 fragment 边界……如果扫描到 WAL 尾部时仍处于 CollectingFragments 状态……丢弃该 incomplete batch。设计同时明确:**CollectingFragments + 不是最后恢复 segment → 视为 WAL 中间损坏,报错**。 +- **代码现状**:`wal/recovery.go:133-138` 在 `ReplaySegmentFile` 末尾,只要处于 `FragmentCollecting` 就返回 `TailCorruptionError`,**不区分当前 segment 是否是最后一个**;`RecoverFromSegments` 也直接透传。 +- **后果**:如果 segment-5 中间出现 `First + Middle*` 但没有 `Last`(中间损坏),但 segment-6 还存在,代码会把 segment-5 的中间损坏当尾部截断处理,截掉 segment-6 的有效数据。设计要求此时硬报错。 +- **修复**:`ReplaySegmentFile` 需要接收 `isLast bool` 参数,或者在 `RecoverFromSegments` 的 segment 循环里检查:非最后段返回 `CollectingFragments` 必须返回普通 corruption 错误,不是 `TailCorruptionError`。 + +> **Oracle 修订(bg_ef425776)**:C4 还有一个更严重的失败模式。`wal/recover.go:61` 总是对 `segments[len(segments)-1]` 调用 `truncateSegment`,但 `RecoverFromSegments` 的 `TailCorruptionError` 可能来自非尾段。结果:真正损坏的 segment 不动,最后一段的有效数据被错误截掉。 + +--- + +### C5. 尾部截断不 fsync、不删空 segment、不 fsync 目录,且失败被吞 + +- **设计要求**(§3.2 尾部截断持久化 line 786-800): + > 1. ftruncate 当前 active segment 到 lastCompleteBatchEnd + > 2. **fsync 被截断的 segment** + > 3. 删除 startSequence == expectedSequence 且不含任何 complete batch 的后续空 segment + > 4. **fsync WAL directory** + > 全部成功后 recovery 才能进入恢复完成状态。**若任一步失败,recovery 必须报错,DB 不得进入可写状态**。 +- **代码现状**:`wal/recover.go:60-69` + ```go + validOffset, truncErr := findValidOffset(lastSeg.FilePath) + if truncErr != nil { + result.TruncateError = fmt.Errorf("%w (find valid offset: %v)", err, truncErr) // 吞掉 + } else if truncErr := truncateSegment(lastSeg.FilePath, validOffset); truncErr != nil { + result.TruncateError = fmt.Errorf("%w (truncate: %v)", err, truncErr) // 吞掉 + } + ``` + `truncateSegment` 只是 `os.Truncate`,没有 fsync 文件、没有删后续空 segment、没有 fsync 目录,失败也只是写进 `result.TruncateError` 然后**正常返回成功**。 +- **后果**:崩溃恢复后 WAL 尾部可能再次暴露已被"截断"的脏字节,下次启动会重复 repair 或 repair 出不同边界;DB 已经接受新写入,破坏了"恢复后的 WAL 状态即权威"这一不变量。 +- **修复**:把 truncation 拆成独立函数,按设计 4 步串行执行,任何一步失败 `return err`,DB.Open 必须失败。 + +--- + +### C6. SegmentWriter 创建时目录 fsync 失败被静默忽略 + +- **设计要求**(§3.2 WAL 元数据持久化协议 line 248-273):durable-ready 协议第 5 步 `fsync WAL directory` 是 segment 进入 durable-ready 的硬条件。目录 fsync 失败时 segment **不得**承载可确认写入;若已有 Batch 依赖该 segment,进入 write-stopped。 +- **代码现状**:`wal/segment_writer.go:86-89` + ```go + if dirFD, derr := os.Open(dir); derr == nil { + dirFD.Sync() // 错误被忽略 + dirFD.Close() + } + ``` +- **后果**:rename 已发生但目录元数据未落盘。掉电后恢复可能看不到这个 segment 文件,但 WAL writer 已经向调用方确认了该 Batch 成功("不丢已确认写入"被破坏)。这是 §3.2 Always 策略下最严重的 durability 漏洞之一。 +- **修复**:目录 fsync 错误必须返回,触发 write-stopped;同时显式管理 durable-ready 状态机(新字段 `durableReady bool`),未就绪时禁止 AppendBatch。 + +--- + +### C7. Put/Delete vs Close 存在 send-on-closed-channel 竞态 + +- **代码现状**:`wal/writer.go:314-326`(Put)先检查 `writeStopped`,再 `queue.Submit`;而 `Close` 先 `writeStopped.Store(true)` 再 `queue.Close()`。两者之间没有同步。 +- **后果**:并发调用 Put 和 Close 时,Put 通过 writeStopped 检查后、Submit 之前,Close 把 channel 关掉,Put 的 `cq.ch <- req` 触发 panic。这不是 §3.2 设计直接约束,但破坏写入路径稳定性。 +- **修复**:Close 用 RWMutex 保护,Submit 用 RLock 检查 closed 标志;或者把 close 时机延后到所有 in-flight submit 完成。 + +--- + +### C8. Segment 轮转时把字节偏移当 sequence number 传给新 segment(Oracle bg_2e86d33b 发现) + +- **设计要求**(§3.2 WAL File Header、Segment 连续性校验 line 639-663):每个 segment 的 `startSequence` 必须严格衔接前一 segment 最后一个 batch 后的 sequence;recovery 时 `segment.StartSequence != expectedSequence` 必须报错。 +- **代码现状**:`wal/segment_manager.go:64-66` + ```go + if sm.active.RemainingPayload() < worstCaseSize { + if err := sm.rotate(sm.active.CurrentOffset()); err != nil { ... } + } + ``` + `rotate(newStartSequence uint64)` 的参数名是 `newStartSequence`,但传入的 `sm.active.CurrentOffset()` 返回的是**字节偏移**(从 `WalFileHeaderSize=32` 累加),不是 sequence number。 +- **后果**:发生 segment 轮转后,新 segment 的 header 里 `startSequence = 上一 segment 的字节偏移`。Recovery 时 `wal/recovery.go:155-157`: + ```go + if segment.StartSequence != nextSequence { error } + ``` + 立刻失败。**Phase 1 多 segment 场景的 recovery 实际上是坏的**。 +- **修复**:`SegmentManager` 需要跟踪当前 `nextSequence`(或从 WalWriter 拿),rotate 时把它传下去,而不是 `CurrentOffset()`。 +- **对 C2+C3 patch 的影响**:相关 e2e 测试(`TestOpenThreeTimesKeepsData`)必须**单 segment**,不能强制轮转,否则会撞上 C8 而不是 C2+C3。 + +--- + +## 🟠 High(破坏格式约束或重要不变量) + +### H1. Physical Record 不校验 `length > 0` + +- **设计要求**(§3.2 Block 边界处理、Physical Record 解析规则 line 389, 681):`length` 必须 `> 0`。 +- **代码现状**:`wal/record.go:38-67` `DecodePhysicalRecord` 只校验 `length <= len(data)-headerSize`,没校验下界。`wal/record_parser.go:69` 也没校验。 +- **后果**:攻击者 / 损坏数据可注入 `length=0` 的 record,绕过 CRC(payload 为空时 CRC 只覆盖 length+type),恢复出空 batch 进而触发 `entry count is zero` 之类的错误,被误判为"尾部损坏可截断",实际上属于中间损坏。 +- **修复**:`DecodePhysicalRecord` 加 `if length == 0 { return ErrZeroLength }`,并在 parser 把它当中间损坏(非尾部不截断)。 + +--- + +### H2. Physical Record 不校验 `type != Invalid (0)` + +- **设计要求**(§3.2 fragment 类型表 line 366-372):type=0 是 invalid,用于损坏检测。 +- **代码现状**:`wal/record.go` 解析时不校验 type;`FragmentCollector.Append` 的 default 分支会拒绝,但 `ParseBlock` 收集 records 时不会。 +- **后果**:损坏的 type=0 record 会被加入 records 列表,传给 fragment collector 后才报错;这把"物理层损坏"推迟到"batch 层错误",错误分类可能错(按设计应该硬错,但实际可能被当 CRC 失败归为尾部损坏)。 +- **修复**:`DecodePhysicalRecord` 校验 `recType ∈ {1,2,3,4}`,否则返回错误。 + +--- + +### H3. `worstCaseSize` 估算不足,segment 可能写超 `MaxSegmentSize` + +- **设计要求**(§3.2 Segment Rotation 约束、config 不变量 line 392-451):判断是否需要 rotate 时必须基于"真实最坏情况"开销:`numRecords * prHeaderSize + worstCaseBlockPadding`,对默认配置最大 batch 是 903 + 7 = 910 bytes 开销。 +- **代码现状**:`wal/segment_manager.go:62` + ```go + worstCaseSize := uint64(len(encodedBatch)) + uint64(PhysicalRecordHeaderSize) + uint64(PhysicalRecordHeaderSize) + ``` + 只加 14 bytes(2 × 7)。对接近 4MB 的大 batch,少估了约 896 bytes。 +- **后果**:当 active segment 剩余 payload 在 `[encodedBatchSize + 14, encodedBatchSize + 910]` 区间时,代码认为放得下,实际写出后超过 `MaxSegmentSize`。ValidateBatchLimits 只校验"空 segment 能放下",没校验"当前剩余能放下"。下一批次才会触发 rotate,期间 segment 文件实际大小超出配置上限。 +- **修复**:把 `ValidateBatchLimits` 中已经算过的 `numRecords`/`overhead`/`padding` 计算抽成公共函数,segment_manager 用同样的公式判断 remaining。或者更稳:rotate 阈值改为 `if active.RemainingPayload() < minFreshSegmentPayload { rotate }`,即只要剩余不足以容纳最大合法 batch 就 rotate。 + +--- + +### H4. `MaxImmutableCount` 配置存在但从不强制 + +- **设计要求**(§3.2 写入流程 step ⑤ line 119): + > 若 Immutable MemTable 队列已达上限,该 Batch 必须在 WAL write 之前等待后台 flush 释放容量。 +- **代码现状**:`wal/writer.go:244-265` `reserveMemTable` 只在 active 满时 rotate,**完全不检查 immutable 队列长度**。`MemTableList.rotateActive` 无脑 append。 +- **后果**:写入速度快于 flush 时,immutable 队列无限增长,OOM;同时也违反了"WAL write 前等待 flush"的设计约束。 +- **修复**:`reserveMemTable` 在 rotate 前检查 `len(immutable) >= MaxImmutableCount`,是则阻塞等待条件变量(需要 flush goroutine 来唤醒)。Phase 1 没有 flush,至少应该在超限时返回 `ErrWriteStopped` 或阻塞,而不是无脑 append。 + +--- + +### H5. Per-batch `(segmentID, endOffset, endSequence)` 跟踪缺失,durableSequence 推进模型错误 + +- **设计要求**(§3.2 line 205-245):每个完整 Batch 必须记录 `(segmentID, endOffset, endSequence)`;后台 fsync worker snapshot `(segmentID, endOffset, endSequence)`;fsync 成功后只能推进到满足"segment 已 durable-ready + offset 覆盖 + endSequence 连续"的最大 batch。跨 segment 推进还需要目标 segment 的 header + rename + 目录都已 fsync。 +- **代码现状**:`wal/sequence.go:68-78` `MarkDurable` 只是 CAS 把 durableSequence 推到 seq;`wal/writer.go:228` 每次 `processBatch` fsync 后立刻调用 `MarkDurable(lastSequence)`。 +- **后果**:当前 Phase 1 只有 Always 模式 + 单 writer 串行 fsync,**功能上恰好正确**(每次 fsync 覆盖且只覆盖当前 batch,durableSequence == publishedSequence)。但一旦未来引入 Periodic / 异步 fsync / 多 batch 合并 fsync,模型立刻崩溃。设计明确要求"durableSequence 不能只根据'最近一次 fsync 成功'模糊推进"。 +- **修复**:哪怕 Phase 1,也要按设计记录每 batch 的 `(segmentID, endOffset, endSequence)`,并实现 fsyncSnapshot 比较逻辑。否则是对未来扩展的债务性违约。 + +--- + +### H6. MemTable publish 机制不符合 release/acquire 模型 + +- **设计要求**(§3.2 step ⑧-⑩ line 121): + > MemTable skiplist / arena 节点必须先通过**原子发布机制**写入读路径可见结构(例如 `atomic.Pointer` store-release,或等价的 release publish),并且 Batch 内所有 entry 的节点都完成发布后,才能用 `atomic.Uint64.Store` 推进 `publishedSequence`。普通读者必须先 `Load` 当前 `publishedSequence`,再遍历 MemTable;读到 entry 后仍以 `entry.sequence <= loadedPublishedSequence` 判断可见性。 +- **代码现状**:`memtable/memtable.go:87-127` `Publish`: + 1. 在 `skiplist.mu` 下遍历收集 pending 节点; + 2. **释放锁后**,对每个 entry 调用 `skiplist.Put(..., pending=false)` 创建**新节点**替换旧节点; + 3. `mt.published.Store(upToSequence)`。 + + 读路径 `skiplist.Get` 用 `node.pending` 判断可见性,**完全不读 `publishedSequence`**。 +- **后果**: + 1. 读路径看到的可见性边界是"pending 标志位",不是 `publishedSequence`。Phase 1 单 key autocommit 下语义等价,但设计要求的是后者。 + 2. Publish 创建新节点替换旧节点,旧 pending 节点成为 GC 垃圾;多个 batch 同时 publish(理论上)会竞争同一 key 的替换路径。 + 3. 弱内存序架构上,新节点写入(通过 `atomic.Pointer.Store`)发生在 `published.Store` 之前,**顺序对的**;但设计要求的 publish-then-load-sequence 模式没有被代码使用,未来引入 MVCC 的 `commitSequence` / `visibleCommitSequence` 时会失配。 + 4. MemTable.published 字段写完后**没人读**,是死代码。 +- **修复**:要么按设计实现:skiplist 节点带 sequence,Put 时直接 published(不 pending),读路径 Load `publishedSequence` 后过滤;要么保留 pending 模型但在设计文档里明确改写,并删除/利用 MemTable.published。 + +--- + +### H7. `arena.Reserve` 只是检查,不是真正预留 + +- **设计要求**(§3.2 step ⑤ line 119):按 Batch 内所有 entry 的最大内存占用(key、value 或 ValueLogPointer、**skiplist 节点**、arena 对齐与**层高开销**)计算需要预留的 Arena 字节数。 +- **代码现状**:`memtable/arena.go:75-84` `Reserve` 只读 offset 不修改;`memtable/memtable.go:55-70` 每条 entry 估算 `metadataOverhead = 32`,没考虑 `maxLevel * sizeof(pointer) = 20 * 8 = 160` bytes 的层高开销。skiplist 节点实际存在堆上(`newSkipNode` 不用 arena)。 +- **后果**:Reserve 是个粗略的 budget gate,实际 heap 用量可能超过 MemTableSize;多线程并发 reserve 还可能 double-count(两个 batch 都看到 remaining 够,都写入,实际超限)。Phase 1 不用 arena 存节点,所以不会触发 ErrArenaFull,但容量预算失效。 +- **修复**:要么把 Reserve 改成真 atomic advance offset(commit 预留),要么显式承认 Phase 1 是软限制并在设计里标注;估算公式按最坏层高(maxLevel=20)算。 + +--- + +## 🟡 Medium(设计偏离但不立即致错) + +### H8. `findValidOffset` 不跟踪 batch 边界,截断点错位(Oracle bg_ef425776 升级 M1) + +- **设计要求**(§3.2 尾部截断持久化 line 786):截断目标必须是"最后一个完整 WAL Batch 的结束位置 `lastCompleteBatchEnd`"。 +- **代码现状**:`wal/recover.go:137-203` 的 `findValidOffset` 只用 `DecodePhysicalRecord` 校验物理记录,**不跟踪完整 batch 边界 / fragment 状态机**。对尾部 `First + Middle*` 没 `Last` 的情况: + - `recovery.go:133` 在最后一个完整 batch 结尾报尾部损坏 + - `findValidOffset` 可能返回不完整 fragment 之后的 EOF,**比真正应截断的位置更靠后** +- **后果**:残留半截 fragment 在文件尾部。下次启动 recovery 还会在同位置触发尾部损坏,**反复 repair,每次结果可能不同**。 +- **修复**:让 `ParseRecordsFromFile` 直接返回 `lastValidOffset`(基于完整 batch 边界),作为 truncation 的权威依据;删掉独立的 `findValidOffset` 重复解析。 + +> 原 M1("两次解析可能不一致")已升级为 H8,因为 Oracle 给出具体证据:不跟踪 batch 边界会让截断点错位,导致反复 repair。 + +### M2. `EncodeWalBatch(0, entries)` 预校验后立刻丢弃,再编一次 + +- `wal/writer.go:192-204` 先用 baseSequence=0 编码一次只为校验,然后再用真实 baseSequence 编码。功能对(baseSequence 不影响 size),但白编一次 4MB buffer。 +- **修复**:抽 `validateEncodedSize(entries)` 函数,只算 size 不分配 buffer;或者直接信任 `ValidateBatchLimits` 已经覆盖的检查(其实已经覆盖了),删掉这次预编码。 + +### M3. `SyncMode` 是字符串 `"always"`,不是设计里的策略枚举 + +- 设计明确三种策略 `Always` / `Periodic` / `Never`,Phase 1 只支持 Always。代码用裸字符串,类型不安全,未来加 Periodic 时容易漏改地方。 +- **修复**:定义 `type SyncMode int` + 常量,配置序列化层做字符串映射。 + +### M4. `CommitQueue` 容量 = `MaxBatchEntries`,可能过大 + +- `wal/writer.go:82` `queueCapacity := max(1, int(cfg.MaxBatchEntries))` 默认 10000,意味着 channel 缓冲 10000 个 `*CommitRequest`。每个 request 至少带 1 个 entry 的 key/value clone。高并发下内存占用被低估。 +- **修复**:独立配置项 `CommitQueueCapacity`,默认更小(例如 1024)。 + +### M5. `SegmentWriter` 没有 `durableReady` 状态字段 + +- 设计要求 segment 显式进入 durable-ready 状态才能承载可确认写入。代码隐式假设"构造函数返回即可写",没有状态机字段。 +- **修复**:加 `durableReady bool` 字段,`AppendBatch` 前断言;构造函数中所有 fsync 步聚通过后才置 true。 + +### M6. `GroupCommitDelay` 计时器语义偏离设计 + +- 设计:"500µs 或 32KB,先到者触发"。代码 `runLoop` 每次从 channel 读到第一个 req 后启 timer,500µs 内继续收集。这相当于"从第一个 req 起等待 500µs",与设计一致。但 timer 在每个 batch 循环重建,如果 channel 持续有 req,低负载时其实仍能 500µs 触发,OK。**不算 bug,但建议注释清楚**。 + +### M7. `MemTable.aborted` 是 `sync.Map` 且从不清理 + +- `memtable/memtable.go:36`。每次 `Abort(seq)` 写入,从不删除。长时间运行 + 频繁 abort 会内存泄漏。 +- **修复**:Publish 时把已 publish 的 sequence 从 aborted 里删掉;或者用 bitmap。 + +### M8. `processRemaining` 在 write-stopped 后仍会处理队列 + +- `wal/writer.go:165-173` 在 close 时 drain 队列。但 `processBatch` 内部会先检查 `writeStopped`,如果已停止会直接 sendError。所以 drain 不会真正写 WAL,OK。**但需要测试覆盖这个路径**。 + +### M9. 错误消息中"segment 损坏"和"尾部损坏"混用,难审计 + +- 多处 `fmt.Errorf` 没用 `%w` 包装 `ErrWALCorrupted`,调用者难以 `errors.Is`。 +- **修复**:定义 `ErrWALCorrupted` / `ErrWALTailCorrupted` 哨兵错误,所有 corruption 路径用 `%w` 包装。 + +--- + +## 🟢 Low / 建议 + +- **L1**:`wal/header.go` 把 `walFileHeaderSize` 等小写常量和 `wal/header.go` 与 `wal/constants.go` 中重复的大写常量合并,避免两套维护。 +- **L2**:`wal/segment_writer.go:42` `os.O_EXCL` 在并发 Open 同一 tmp 名时直接失败,错误信息可加上"可能是上次崩溃残留"。 +- **L3**:`memtable/skiplist.go:209` `rand.Float64()` 非并发安全(Go 1.20+ 默认 auto-seed 后是安全的,但显式用 `math/rand/v2` 更清楚)。 +- **L4**:`config/config.go:124` `SyncMode != "always"` 错误消息建议列出合法值。 +- **L5**:测试覆盖建议补:① crc32c 固定向量;② 非尾段 CollectingFragments 必须 hard error;③ MANIFEST 不被 recovery 修改;④ Put(k, []) 与 Delete(k) 的 Get 结果在恢复后仍能区分;⑤ dir fsync 失败时 segment 不可写。 + +--- + +## 总结(必须立刻修的) + +| 编号 | 一句话 | 风险等级 | +|------|--------|----------| +| **C1** | CRC 多项式错(IEEE → crc32c) | 格式不符规范 | +| **C2** | Recovery 用 CURRENT 兜底 | 恢复起点错误 | +| **C3** | Recovery 写 MANIFEST | **下次启动跳过有效 segment,丢数据** | +| **C4** | 非尾段 CollectingFragments 当尾段损坏 | 中间损坏被静默截断 | +| **C5** | 截断无 fsync / 不删空 segment / 失败被吞 | **重复 repair,DB 状态不一致** | +| **C6** | SegmentWriter dir fsync 失败被吞 | **掉电后已确认写入消失** | +| **C7** | Put vs Close 的 channel 竞态 | 偶发 panic | +| **C8** | Segment 轮转把字节偏移当 startSequence | **多 segment recovery 直接失败** | + +C3 和 C6 是最严重的:C3 直接导致数据丢失,C6 直接破坏 Always 模式的"不丢已确认写入"承诺。C8 是隐藏炸弹:单 segment 时一切正常,第一次轮转后就坏。建议优先修复 C1-C8 这 8 条,再处理 H1-H7。 + +--- + +## 修复优先级建议 + +按数据丢失风险排:**C3 → C6 → C5 → C2 → C4 → C1 → C7**,再处理 H1-H7。 + +> **Oracle 修订(bg_2e86d33b)**:实际上 C2 和 C3 必须**捆绑**修复(C2 的触发条件在 Phase 1 默认状态下就成立,比 C3 更宽松),单独修 C3 不止血。修订后优先级:**C2 + C3(一起)→ C6 → C5 + H8(一起,同一文件)→ C4 → C8 → C1 → C7**。 diff --git a/wal/recover.go b/wal/recover.go index ec05bae..907b452 100644 --- a/wal/recover.go +++ b/wal/recover.go @@ -24,7 +24,7 @@ func Recover(dir string, replayer BatchReplayer) (*RecoveryResult, error) { return nil, fmt.Errorf("wal: recover: replayer is nil") } - // Step 1: Determine recovery segment ID from MANIFEST or CURRENT. + // Step 1: Determine recovery segment ID from MANIFEST. recoverySegmentID, err := resolveRecoverySegmentID(dir) if err != nil { return nil, fmt.Errorf("wal: recover: resolve segment id: %w", err) @@ -73,10 +73,10 @@ func Recover(dir string, replayer BatchReplayer) (*RecoveryResult, error) { // the replayer interface doesn't expose a count. result.ReplayedEntries = 0 // caller can inspect replayer directly - // Update MANIFEST with new recovery state. - if saveErr := manifest.Save(dir, result.NextSegmentID); saveErr != nil { - return nil, fmt.Errorf("wal: recover: save manifest after truncation: %w", saveErr) - } + // Per design §3.2 line 280, recovery repair must NOT update MANIFEST. + // The truncated WAL state is persisted via ftruncate (see C5 for the + // remaining fsync gaps). MANIFEST can only advance via checkpoint + // (MemTable flush) in future phases. return result, nil } @@ -96,31 +96,22 @@ func Recover(dir string, replayer BatchReplayer) (*RecoveryResult, error) { result.NextSegmentID = segments[len(segments)-1].SegmentID + 1 } - // Update MANIFEST with new recovery state. - if saveErr := manifest.Save(dir, result.NextSegmentID); saveErr != nil { - return nil, fmt.Errorf("wal: recover: save manifest: %w", saveErr) - } + // Per design §3.2 line 280, recovery must NOT update MANIFEST. + // RecoveryResult.NextSegmentID is in-memory only, consumed by DB.Open to + // seed the new WalWriter. MANIFEST stays at its pre-recovery value. return result, nil } -// resolveRecoverySegmentID determines the starting segment ID for recovery. -// It tries MANIFEST first, then falls back to CURRENT, then defaults to 0. +// resolveRecoverySegmentID returns the recovery start segment ID from MANIFEST. +// MANIFEST is the only authoritative source of recovery start per design §3.2 +// line 600-06. CURRENT is a write-side hint and must NOT be used here. func resolveRecoverySegmentID(dir string) (uint64, error) { mf, err := manifest.Load(dir) if err != nil { return 0, fmt.Errorf("load manifest: %w", err) } - if mf.RecoverySegmentID > 0 { - return mf.RecoverySegmentID, nil - } - - // MANIFEST had 0 (fresh DB or not yet written). Try CURRENT. - if segID, ok := manifest.ReadCurrent(dir); ok { - return segID, nil - } - - return 0, nil + return mf.RecoverySegmentID, nil } // truncateSegment truncates the file at filePath to validOffset bytes, diff --git a/wal/recover_test.go b/wal/recover_test.go index fcf64a1..de026e6 100644 --- a/wal/recover_test.go +++ b/wal/recover_test.go @@ -1,6 +1,7 @@ package wal import ( + "bytes" "os" "path/filepath" "reflect" @@ -124,34 +125,211 @@ func TestRecoverWithTailCorruption(t *testing.T) { } } -func TestRecoverUpdatesManifest(t *testing.T) { +func TestRecoverDoesNotUpdateManifest(t *testing.T) { dir := t.TempDir() - // Write test data. writeTestSegment(t, dir, 0, 50, [][]*WalEntry{ {makePutEntry("a", "b")}, {makePutEntry("c", "d")}, }) + beforeMF, err := manifest.Load(dir) + if err != nil { + t.Fatalf("manifest.Load before recover: %v", err) + } + beforeExists := fileExists(t, filepath.Join(dir, "MANIFEST")) + + replayer := &mockReplayer{} + if _, err := Recover(dir, replayer); err != nil { + t.Fatalf("Recover: %v", err) + } + + afterMF, err := manifest.Load(dir) + if err != nil { + t.Fatalf("manifest.Load after recover: %v", err) + } + if afterMF.RecoverySegmentID != beforeMF.RecoverySegmentID { + t.Errorf("MANIFEST RecoverySegmentID changed: %d -> %d", + beforeMF.RecoverySegmentID, afterMF.RecoverySegmentID) + } + afterExists := fileExists(t, filepath.Join(dir, "MANIFEST")) + if beforeExists != afterExists { + t.Errorf("MANIFEST file existence changed: before=%v after=%v", + beforeExists, afterExists) + } +} + +// Regression guard for C3: covers the case where MANIFEST already exists. +// The fresh-DB test above cannot catch accidental overwrites of an existing +// MANIFEST. +func TestRecoverPreservesExistingManifest(t *testing.T) { + dir := t.TempDir() + + writeTestSegment(t, dir, 0, 50, [][]*WalEntry{ + {makePutEntry("a", "b")}, + }) + + if err := manifest.Save(dir, 0); err != nil { + t.Fatalf("manifest.Save setup: %v", err) + } + beforeBytes, err := os.ReadFile(filepath.Join(dir, "MANIFEST")) + if err != nil { + t.Fatalf("ReadFile MANIFEST: %v", err) + } + + replayer := &mockReplayer{} + if _, err := Recover(dir, replayer); err != nil { + t.Fatalf("Recover: %v", err) + } + + afterBytes, err := os.ReadFile(filepath.Join(dir, "MANIFEST")) + if err != nil { + t.Fatalf("ReadFile MANIFEST after recover: %v", err) + } + if !bytes.Equal(beforeBytes, afterBytes) { + t.Errorf("MANIFEST bytes changed:\n before=%q\n after=%q", + string(beforeBytes), string(afterBytes)) + } +} + +// Regression guard for C3: design §3.2 line 280 requires recovery to be +// idempotent on a clean WAL (no MANIFEST side effects). +func TestRecoverIdempotentClean(t *testing.T) { + dir := t.TempDir() + + writeTestSegment(t, dir, 0, 0, [][]*WalEntry{ + {makePutEntry("k1", "v1")}, + {makePutEntry("k2", "v2")}, + }) + + replayer1 := &mockReplayer{} + result1, err := Recover(dir, replayer1) + if err != nil { + t.Fatalf("Recover (1st): %v", err) + } + + replayer2 := &mockReplayer{} + result2, err := Recover(dir, replayer2) + if err != nil { + t.Fatalf("Recover (2nd): %v", err) + } + + if result1.NextSequence != result2.NextSequence { + t.Errorf("NextSequence differs: %d vs %d", result1.NextSequence, result2.NextSequence) + } + if result1.NextSegmentID != result2.NextSegmentID { + t.Errorf("NextSegmentID differs: %d vs %d", result1.NextSegmentID, result2.NextSegmentID) + } + if result1.Truncated || result2.Truncated { + t.Errorf("Truncated should be false for clean WAL: r1=%v r2=%v", + result1.Truncated, result2.Truncated) + } + if !reflect.DeepEqual(replayer1.puts, replayer2.puts) { + t.Errorf("replayed puts differ:\n r1=%#v\n r2=%#v", replayer1.puts, replayer2.puts) + } +} + +// Regression guard for C3: after the first recovery truncates a corrupted +// tail, the second recovery must observe stable state (Truncated=false, same +// NextSequence, same replayed entries). +func TestRecoverIdempotentAfterTruncation(t *testing.T) { + dir := t.TempDir() + + filePath := writeTestSegment(t, dir, 0, 100, [][]*WalEntry{ + {makePutEntry("good1", "before-corruption")}, + {makePutEntry("good2", "also-before")}, + }) + appendFileBytes(t, filePath, []byte{0xDE, 0xAD, 0xBE, 0xEF}) + + replayer1 := &mockReplayer{} + result1, err := Recover(dir, replayer1) + if err != nil { + t.Fatalf("Recover (1st): %v", err) + } + if !result1.Truncated { + t.Fatal("1st Recover: Truncated = false, want true") + } + + replayer2 := &mockReplayer{} + result2, err := Recover(dir, replayer2) + if err != nil { + t.Fatalf("Recover (2nd): %v", err) + } + if result2.Truncated { + t.Error("2nd Recover: Truncated = true, want false (tail already repaired)") + } + if result1.NextSequence != result2.NextSequence { + t.Errorf("NextSequence differs: %d vs %d", result1.NextSequence, result2.NextSequence) + } + if !reflect.DeepEqual(replayer1.puts, replayer2.puts) { + t.Errorf("replayed puts differ:\n r1=%#v\n r2=%#v", replayer1.puts, replayer2.puts) + } +} + +// Regression guard for C2: design §3.2 line 604-06 forbids using CURRENT as +// recovery start. writeTestSegment uses NewSegmentWriter directly, which does +// NOT write CURRENT (only SegmentManager does), so we write CURRENT explicitly +// to simulate the Phase 1 default state. +func TestRecoverIgnoresCurrentFallback(t *testing.T) { + dir := t.TempDir() + + writeTestSegment(t, dir, 0, 0, [][]*WalEntry{ + {makePutEntry("seg0-k1", "v1")}, + }) + writeTestSegment(t, dir, 1, 1, [][]*WalEntry{ + {makePutEntry("seg1-k1", "v1")}, + }) + writeTestSegment(t, dir, 2, 2, [][]*WalEntry{ + {makePutEntry("seg2-k1", "v1")}, + }) + + if err := manifest.WriteCurrent(dir, 2); err != nil { + t.Fatalf("WriteCurrent: %v", err) + } + + currentSegID, ok := manifest.ReadCurrent(dir) + if !ok || currentSegID != 2 { + t.Fatalf("CURRENT setup wrong: segID=%d ok=%v", currentSegID, ok) + } + mf, err := manifest.Load(dir) + if err != nil { + t.Fatalf("Load: %v", err) + } + if mf.RecoverySegmentID != 0 { + t.Fatalf("MANIFEST.RecoverySegmentID = %d, want 0", mf.RecoverySegmentID) + } + replayer := &mockReplayer{} result, err := Recover(dir, replayer) if err != nil { t.Fatalf("Recover: %v", err) } - // Verify MANIFEST file exists and has correct content. - mf, err := manifest.Load(dir) - if err != nil { - t.Fatalf("manifest.Load after recover: %v", err) + wantPuts := []replayPut{ + {key: "seg0-k1", value: "v1", seq: 0}, + {key: "seg1-k1", value: "v1", seq: 1}, + {key: "seg2-k1", value: "v1", seq: 2}, } - if mf.RecoverySegmentID != result.NextSegmentID { - t.Errorf("MANIFEST RecoverySegmentID = %d, want %d", - mf.RecoverySegmentID, result.NextSegmentID) + if !reflect.DeepEqual(replayer.puts, wantPuts) { + t.Errorf("puts = %#v, want %#v", replayer.puts, wantPuts) } - - // Verify MANIFEST file path. - manifestPath := filepath.Join(dir, "MANIFEST") - if _, err := os.Stat(manifestPath); os.IsNotExist(err) { - t.Error("MANIFEST file should exist after Recover") + if result.NextSequence != 3 { + t.Errorf("NextSequence = %d, want 3", result.NextSequence) + } + if result.NextSegmentID != 3 { + t.Errorf("NextSegmentID = %d, want 3", result.NextSegmentID) } } + +func fileExists(t *testing.T, path string) bool { + t.Helper() + _, err := os.Stat(path) + if err == nil { + return true + } + if os.IsNotExist(err) { + return false + } + t.Fatalf("stat %s: %v", path, err) + return false +}