package eventbus import ( "context" "encoding/json" "errors" "os" "path/filepath" "strings" "sync" "testing" "time" ) func TestDurableOutboxWALConcurrentGroupCommitRecoversEveryRecord(t *testing.T) { dir := t.TempDir() wal, err := newDurableOutboxWAL(durableOutboxWALConfig{ Directory: dir, SyncWrites: true, CommitBatch: 64, CommitInterval: 5 * time.Millisecond, }) if err != nil { t.Fatalf("new WAL: %v", err) } const total = 256 refs := make(chan outboxWALRecordRef, total) errs := make(chan error, total) var workers sync.WaitGroup for i := 0; i < total; i++ { workers.Add(1) go func(sequence int) { defer workers.Done() record := walTestRecord(sequence) stored, appendErr := wal.Append(context.Background(), record) if appendErr != nil { errs <- appendErr return } refs <- stored.Ref }(i) } workers.Wait() close(errs) for appendErr := range errs { t.Fatalf("concurrent append: %v", appendErr) } close(refs) for ref := range refs { wal.Release(ref) } if got, _ := wal.Stats(); got != total { t.Fatalf("backlog before restart = %d, want %d", got, total) } if err := wal.Close(); err != nil { t.Fatalf("close first WAL: %v", err) } recovered, err := newDurableOutboxWAL(durableOutboxWALConfig{Directory: dir, SyncWrites: true}) if err != nil { t.Fatalf("reopen WAL: %v", err) } records, err := recovered.ClaimPending(total + 1) if err != nil { t.Fatalf("claim recovered records: %v", err) } if len(records) != total { t.Fatalf("recovered records = %d, want %d", len(records), total) } seen := make(map[string]struct{}, total) for _, record := range records { seen[record.Record.Envelope.EventID] = struct{}{} if err := recovered.Ack(record.Ref); err != nil { t.Fatalf("ack recovered record: %v", err) } } if len(seen) != total { t.Fatalf("unique recovered event ids = %d, want %d", len(seen), total) } if got, _ := recovered.Stats(); got != 0 { t.Fatalf("backlog after ack = %d, want 0", got) } if err := recovered.Close(); err != nil { t.Fatalf("close recovered WAL: %v", err) } } func TestDurableOutboxWALTruncatesIncompleteTrailingFrame(t *testing.T) { dir := t.TempDir() first := encodedWALTestFrame(t, 1) second := encodedWALTestFrame(t, 2) path := filepath.Join(dir, outboxWALFileName(1)) payload := append(append([]byte{}, first...), second[:len(second)/2]...) if err := os.WriteFile(path, payload, 0o640); err != nil { t.Fatalf("write incomplete WAL: %v", err) } wal, err := newDurableOutboxWAL(durableOutboxWALConfig{Directory: dir, SyncWrites: true}) if err != nil { t.Fatalf("recover incomplete WAL: %v", err) } info, err := os.Stat(path) if err != nil { t.Fatalf("stat recovered segment: %v", err) } if got, want := info.Size(), int64(len(first)); got != want { t.Fatalf("truncated size = %d, want %d", got, want) } if got, _ := wal.Stats(); got != 1 { t.Fatalf("recovered backlog = %d, want 1", got) } if err := wal.Close(); err != nil { t.Fatalf("close WAL: %v", err) } } func TestDurableOutboxWALRejectsChecksumCorruption(t *testing.T) { dir := t.TempDir() frame := encodedWALTestFrame(t, 1) frame[len(frame)-1] ^= 0xff path := filepath.Join(dir, outboxWALFileName(1)) if err := os.WriteFile(path, frame, 0o640); err != nil { t.Fatalf("write corrupt WAL: %v", err) } _, err := newDurableOutboxWAL(durableOutboxWALConfig{Directory: dir}) if err == nil || !strings.Contains(err.Error(), "checksum mismatch") { t.Fatalf("corrupt WAL error = %v", err) } } func TestDurableOutboxWALAckDeletesClosedSegment(t *testing.T) { dir := t.TempDir() frameSize := int64(len(encodedWALTestFrame(t, 1))) wal, err := newDurableOutboxWAL(durableOutboxWALConfig{ Directory: dir, SegmentBytes: frameSize, CommitBatch: 1, }) if err != nil { t.Fatalf("new WAL: %v", err) } first, err := wal.Append(context.Background(), walTestRecord(1)) if err != nil { t.Fatalf("append first: %v", err) } firstPath := filepath.Join(dir, outboxWALFileName(first.Ref.segmentID)) second, err := wal.Append(context.Background(), walTestRecord(2)) if err != nil { t.Fatalf("append second: %v", err) } if first.Ref.segmentID == second.Ref.segmentID { t.Fatal("second append should rotate to a new segment") } if err := wal.Ack(first.Ref); err != nil { t.Fatalf("ack first: %v", err) } if _, err := os.Stat(firstPath); !errors.Is(err, os.ErrNotExist) { t.Fatalf("closed acknowledged segment still exists: %v", err) } wal.Release(second.Ref) if err := wal.Close(); err != nil { t.Fatalf("close WAL: %v", err) } } func TestDurableOutboxWALClaimIsBoundedAndNotDuplicated(t *testing.T) { wal, err := newDurableOutboxWAL(durableOutboxWALConfig{Directory: t.TempDir(), CommitBatch: 1}) if err != nil { t.Fatalf("new WAL: %v", err) } for i := 0; i < 10; i++ { stored, err := wal.Append(context.Background(), walTestRecord(i)) if err != nil { t.Fatalf("append %d: %v", i, err) } wal.Release(stored.Ref) } first, err := wal.ClaimPending(3) if err != nil || len(first) != 3 { t.Fatalf("first claim = %d, error = %v", len(first), err) } second, err := wal.ClaimPending(3) if err != nil || len(second) != 3 { t.Fatalf("second claim = %d, error = %v", len(second), err) } claimed := map[outboxWALRecordRef]struct{}{} for _, record := range append(first, second...) { if _, duplicate := claimed[record.Ref]; duplicate { t.Fatalf("record claimed twice: %#v", record.Ref) } claimed[record.Ref] = struct{}{} wal.Release(record.Ref) } if err := wal.Close(); err != nil { t.Fatalf("close WAL: %v", err) } } func TestDurableOutboxWALRejectsAppendAfterClose(t *testing.T) { wal, err := newDurableOutboxWAL(durableOutboxWALConfig{Directory: t.TempDir()}) if err != nil { t.Fatalf("new WAL: %v", err) } if err := wal.Close(); err != nil { t.Fatalf("close WAL: %v", err) } _, err = wal.Append(context.Background(), walTestRecord(1)) if !errors.Is(err, ErrDurableOutboxWALClosed) { t.Fatalf("append after close error = %v", err) } } func walTestRecord(sequence int) durableRecord { env := normalizeDurableEnvelope(durableTestEnvelope()) env.EventID = "wal-event-" + time.Unix(0, int64(sequence)+1).UTC().Format("150405.000000000") env.Sequence = uint16(sequence) return durableRecord{Kind: "raw", Envelope: env} } func encodedWALTestFrame(t *testing.T, sequence int) []byte { t.Helper() payload, err := json.Marshal(walTestRecord(sequence)) if err != nil { t.Fatalf("marshal WAL test record: %v", err) } return encodeOutboxWALFrame(payload) }