package wal import ( "errors" "fmt" "os" "path/filepath" "strings" "testing" "github.com/dailz/go-kv/config" ) func testWalConfig() *config.WalConfig { cfg := config.Defaults() return &cfg } func TestSegmentWriterCreation(t *testing.T) { dir := t.TempDir() cfg := testWalConfig() sw, err := NewSegmentWriter(dir, 1, 100, cfg) if err != nil { t.Fatalf("NewSegmentWriter: %v", err) } expectedPath := filepath.Join(dir, "segment-1.wal") if sw.SegmentPath() != expectedPath { t.Errorf("SegmentPath = %q, want %q", sw.SegmentPath(), expectedPath) } if sw.SegmentID() != 1 { t.Errorf("SegmentID = %d, want 1", sw.SegmentID()) } if sw.CurrentOffset() != WalFileHeaderSize { t.Errorf("CurrentOffset = %d, want %d", sw.CurrentOffset(), WalFileHeaderSize) } // Verify the file exists and has correct header. data, err := os.ReadFile(expectedPath) if err != nil { t.Fatalf("ReadFile: %v", err) } if len(data) != WalFileHeaderSize { t.Errorf("file size = %d, want %d (header only)", len(data), WalFileHeaderSize) } hdr, err := DecodeWalHeader(data) if err != nil { t.Fatalf("DecodeWalHeader: %v", err) } if hdr.SegmentID != 1 { t.Errorf("header SegmentID = %d, want 1", hdr.SegmentID) } if hdr.StartSequence != 100 { t.Errorf("header StartSequence = %d, want 100", hdr.StartSequence) } if hdr.BlockSize != cfg.BlockSize { t.Errorf("header BlockSize = %d, want %d", hdr.BlockSize, cfg.BlockSize) } // No .tmp file should remain. tmpPath := filepath.Join(dir, "segment-1.wal.tmp") if _, err := os.Stat(tmpPath); !os.IsNotExist(err) { t.Errorf("temp file %q should not exist", tmpPath) } if err := sw.Close(); err != nil { t.Fatalf("Close: %v", err) } } func TestSegmentWriterAppendBatch(t *testing.T) { dir := t.TempDir() cfg := testWalConfig() sw, err := NewSegmentWriter(dir, 42, 0, cfg) if err != nil { t.Fatalf("NewSegmentWriter: %v", err) } t.Cleanup(func() { sw.Close() }) // Encode a small batch. entries := []*WalEntry{ {OpType: OpPut, ValueKind: VKInline, Key: []byte("key1"), Value: []byte("val1")}, {OpType: OpPut, ValueKind: VKInline, Key: []byte("key2"), Value: []byte("val2")}, } encoded, err := EncodeWalBatch(0, entries) if err != nil { t.Fatalf("EncodeWalBatch: %v", err) } if err := sw.AppendBatch(encoded); err != nil { t.Fatalf("AppendBatch: %v", err) } if err := sw.Close(); err != nil { t.Fatalf("Close: %v", err) } // Read back the file and verify records. data, err := os.ReadFile(sw.SegmentPath()) if err != nil { t.Fatalf("ReadFile: %v", err) } // Skip file header. body := data[WalFileHeaderSize:] // Use FragmentCollector to reassemble. fc := NewFragmentCollector() offset := 0 for offset < len(body) { // Check for trailing zeros (block padding). if body[offset] == 0 { break } rec, consumed, err := DecodePhysicalRecord(body[offset:]) if err != nil { t.Fatalf("DecodePhysicalRecord at offset %d: %v", offset, err) } if err := fc.Append(rec.Type, rec.Payload); err != nil { t.Fatalf("FragmentCollector.Append: %v", err) } offset += consumed } if !fc.IsComplete() { t.Fatal("fragment collector should be complete") } decoded, err := DecodeWalBatch(fc.BatchData()) if err != nil { t.Fatalf("DecodeWalBatch: %v", err) } if decoded.EntryCount != 2 { t.Errorf("EntryCount = %d, want 2", decoded.EntryCount) } if decoded.BaseSequence != 0 { t.Errorf("BaseSequence = %d, want 0", decoded.BaseSequence) } } func TestSegmentWriterMultipleBatches(t *testing.T) { dir := t.TempDir() cfg := testWalConfig() sw, err := NewSegmentWriter(dir, 1, 0, cfg) if err != nil { t.Fatalf("NewSegmentWriter: %v", err) } for i := 0; i < 5; i++ { entries := []*WalEntry{ { OpType: OpPut, ValueKind: VKInline, Key: []byte("key"), Value: []byte("val"), }, } encoded, err := EncodeWalBatch(uint64(i), entries) if err != nil { t.Fatalf("EncodeWalBatch %d: %v", i, err) } if err := sw.AppendBatch(encoded); err != nil { t.Fatalf("AppendBatch %d: %v", i, err) } } if err := sw.Close(); err != nil { t.Fatalf("Close: %v", err) } data, err := os.ReadFile(sw.SegmentPath()) if err != nil { t.Fatalf("ReadFile: %v", err) } body := data[WalFileHeaderSize:] fc := NewFragmentCollector() batchCount := 0 offset := 0 for offset < len(body) { if body[offset] == 0 { break } rec, consumed, err := DecodePhysicalRecord(body[offset:]) if err != nil { t.Fatalf("DecodePhysicalRecord at offset %d: %v", offset, err) } offset += consumed if err := fc.Append(rec.Type, rec.Payload); err != nil { t.Fatalf("FragmentCollector.Append: %v", err) } if fc.IsComplete() { batchCount++ decoded, err := DecodeWalBatch(fc.BatchData()) if err != nil { t.Fatalf("DecodeWalBatch %d: %v", batchCount, err) } if decoded.BaseSequence != uint64(batchCount-1) { t.Errorf("batch %d BaseSequence = %d, want %d", batchCount, decoded.BaseSequence, batchCount-1) } fc.Reset() } } if batchCount != 5 { t.Errorf("decoded %d batches, want 5", batchCount) } } func TestSegmentWriterRemainingPayload(t *testing.T) { dir := t.TempDir() cfg := testWalConfig() sw, err := NewSegmentWriter(dir, 1, 0, cfg) if err != nil { t.Fatalf("NewSegmentWriter: %v", err) } t.Cleanup(func() { sw.Close() }) initialRemaining := sw.RemainingPayload() entries := []*WalEntry{ {OpType: OpPut, ValueKind: VKInline, Key: []byte("k"), Value: []byte("v")}, } encoded, err := EncodeWalBatch(0, entries) if err != nil { t.Fatalf("EncodeWalBatch: %v", err) } if err := sw.AppendBatch(encoded); err != nil { t.Fatalf("AppendBatch: %v", err) } if sw.RemainingPayload() >= initialRemaining { t.Errorf("RemainingPayload should decrease after write, got %d >= %d", sw.RemainingPayload(), initialRemaining) } } func TestSegmentWriterLargeBatch(t *testing.T) { dir := t.TempDir() cfg := testWalConfig() sw, err := NewSegmentWriter(dir, 1, 0, cfg) if err != nil { t.Fatalf("NewSegmentWriter: %v", err) } // Create a batch larger than one block payload (~32KB - 7 bytes). largeValue := make([]byte, MaxWalInlineValueBytes) for i := range largeValue { largeValue[i] = byte(i % 256) } var entries []*WalEntry for i := 0; i < 9; i++ { entries = append(entries, &WalEntry{ OpType: OpPut, ValueKind: VKInline, Key: []byte(fmt.Sprintf("large-key-%d", i)), Value: largeValue, }) } encoded, err := EncodeWalBatch(0, entries) if err != nil { t.Fatalf("EncodeWalBatch: %v", err) } if err := sw.AppendBatch(encoded); err != nil { t.Fatalf("AppendBatch: %v", err) } if err := sw.Close(); err != nil { t.Fatalf("Close: %v", err) } // Read back and verify. data, err := os.ReadFile(sw.SegmentPath()) if err != nil { t.Fatalf("ReadFile: %v", err) } body := data[WalFileHeaderSize:] fc := NewFragmentCollector() offset := 0 fragCount := 0 for offset < len(body) { if body[offset] == 0 { break } rec, consumed, err := DecodePhysicalRecord(body[offset:]) if err != nil { t.Fatalf("DecodePhysicalRecord at offset %d: %v", offset, err) } offset += consumed fragCount++ if err := fc.Append(rec.Type, rec.Payload); err != nil { t.Fatalf("FragmentCollector.Append type=%d: %v", rec.Type, err) } } if !fc.IsComplete() { t.Fatal("fragment collector should be complete after large batch") } if fragCount < 2 { t.Errorf("expected multiple fragments for large batch, got %d", fragCount) } decoded, err := DecodeWalBatch(fc.BatchData()) if err != nil { t.Fatalf("DecodeWalBatch: %v", err) } if decoded.EntryCount != uint32(len(entries)) { t.Errorf("EntryCount = %d, want %d", decoded.EntryCount, len(entries)) } } func TestSegmentWriterSync(t *testing.T) { dir := t.TempDir() cfg := testWalConfig() sw, err := NewSegmentWriter(dir, 1, 0, cfg) if err != nil { t.Fatalf("NewSegmentWriter: %v", err) } if err := sw.Sync(); err != nil { t.Fatalf("Sync: %v", err) } if err := sw.Close(); err != nil { t.Fatalf("Close: %v", err) } } // Regression guard for C6: NewSegmentWriter must return an error when // directory fsync fails, instead of silently succeeding. Per design ยง3.2 // line 272, segment must NOT become active if durable-ready fails. // // dirFsyncFn is a package-level var; tests that override it must not use // t.Parallel(). All wal tests run serially within the package. func TestNewSegmentWriterDirFsyncFailure(t *testing.T) { dir := t.TempDir() cfg := config.Defaults() orig := dirFsyncFn dirFsyncFn = func(string) error { return errors.New("simulated dir fsync failure") } t.Cleanup(func() { dirFsyncFn = orig }) sw, err := NewSegmentWriter(dir, 0, 0, &cfg) if err == nil { if sw != nil { sw.Close() } t.Fatal("NewSegmentWriter: expected error on dir fsync failure, got nil") } if !strings.Contains(err.Error(), "fsync directory") { t.Errorf("error should mention 'fsync directory', got: %v", err) } entries, err := os.ReadDir(dir) if err != nil { t.Fatalf("ReadDir: %v", err) } for _, e := range entries { name := e.Name() if strings.Contains(name, "segment-0") { t.Errorf("segment file should be cleaned up, found: %s", name) } } } // Regression guard for C6: ensure normal path still works after the fix. func TestNewSegmentWriterNormalPathStillWorks(t *testing.T) { dir := t.TempDir() cfg := config.Defaults() sw, err := NewSegmentWriter(dir, 0, 0, &cfg) if err != nil { t.Fatalf("NewSegmentWriter normal path: %v", err) } defer sw.Close() if _, err := os.Stat(sw.SegmentPath()); err != nil { t.Errorf("segment file should exist: %v", err) } } // Regression guard for C6: after a failed NewSegmentWriter due to dir // fsync, retrying with fsync restored must succeed and not leak state. func TestNewSegmentWriterRetryAfterDirFsyncFailure(t *testing.T) { dir := t.TempDir() cfg := config.Defaults() orig := dirFsyncFn dirFsyncFn = func(string) error { return errors.New("simulated") } _, err := NewSegmentWriter(dir, 0, 0, &cfg) if err == nil { t.Fatal("expected first NewSegmentWriter to fail") } dirFsyncFn = orig sw, err := NewSegmentWriter(dir, 0, 0, &cfg) if err != nil { t.Fatalf("retry NewSegmentWriter: %v", err) } defer sw.Close() entries, err := os.ReadDir(dir) if err != nil { t.Fatalf("ReadDir: %v", err) } tmpCount := 0 walCount := 0 for _, e := range entries { if strings.HasSuffix(e.Name(), ".tmp") { tmpCount++ } if strings.HasSuffix(e.Name(), ".wal") { walCount++ } } if tmpCount != 0 { t.Errorf("leftover .tmp files: %d", tmpCount) } if walCount != 1 { t.Errorf("expected exactly 1 .wal file (from successful retry), got %d", walCount) } }