From e23c8f9da76e7758f98912b7e247f5194a39645a Mon Sep 17 00:00:00 2001 From: Tommaso Barbugli Date: Fri, 14 Aug 2026 15:42:22 +0200 Subject: [PATCH 1/2] cdc: bound catchup segment work Reuse recovered segment bounds so catchup and pruning do not repeatedly scan large CDC backlogs, while preserving crash recovery and conservative deletion invariants. Co-authored-by: Cursor --- internal/app/app.go | 26 ++- internal/cdc/applier.go | 22 +- internal/cdc/cdc_integration_test.go | 172 +++++++++++++++ internal/cdc/lifecycle_test.go | 309 +++++++++++++++++++++++++++ internal/cdc/maintenance.go | 209 +++++++++++++++++- internal/cdc/reader.go | 4 + internal/cdc/segment.go | 168 ++++++++++++++- internal/cdc/segment_test.go | 141 ++++++++++++ 8 files changed, 1018 insertions(+), 33 deletions(-) diff --git a/internal/app/app.go b/internal/app/app.go index aab359e..6260929 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -605,7 +605,9 @@ func (a App) Run(ctx context.Context, cfg config.Config) (runErr error) { if err := pauseForCrashTest(groupCtx, state.PhaseCatchup); err != nil { return err } - return runApplierToFollow(groupCtx, cfg, store, holder.Snapshot, durable, state.PhaseCatchup) + return runApplierToFollow( + groupCtx, cfg, store, holder.Snapshot, durable, writer.SegmentCatalog(), state.PhaseCatchup, + ) }) if cfg.Metrics != "" { group.Go(func() error { return serveMetrics(groupCtx, cfg.Metrics, store) }) @@ -708,7 +710,9 @@ func (a App) resumePostCopy(ctx context.Context, cfg config.Config, store *state return err } } - return runApplierToFollow(groupCtx, cfg, store, snapshot, durable, phase) + return runApplierToFollow( + groupCtx, cfg, store, snapshot, durable, writer.SegmentCatalog(), phase, + ) }) group.Go(func() error { ticker := time.NewTicker(200 * time.Millisecond) @@ -1329,6 +1333,7 @@ func runApplierToFollow( store *state.Store, snapshot setup.Snapshot, durable *cdc.DurableWatermark, + segments *cdc.SegmentCatalog, phase state.Phase, ) error { migration, err := store.Migration(ctx) @@ -1338,6 +1343,7 @@ func runApplierToFollow( pruner, err := cdc.NewSegmentPruner(cdc.SegmentPrunerConfig{ Directory: filepath.Join(cfg.Dir, "cdc"), Interval: cfg.SegmentPruneInterval, + Catalog: segments, }) if err != nil { return err @@ -1365,13 +1371,13 @@ func runApplierToFollow( return err } if phase == state.PhaseFollow || phase == state.PhaseDrained || phase == state.PhaseCutover { - return runApplierContinuous(ctx, applier, store) + return runApplierWithPruner(ctx, applier, pruner, store) } boundary := durable.Load() applyCtx, cancel := context.WithCancel(ctx) defer cancel() result := make(chan error, 1) - go func() { result <- runApplierContinuous(applyCtx, applier, store) }() + go func() { result <- runApplierWithPruner(applyCtx, applier, pruner, store) }() if err := awaitCatchup(ctx, boundary, applier.WaitUntil, result); err != nil { return err } @@ -1390,6 +1396,18 @@ func runApplierToFollow( } } +func runApplierWithPruner( + ctx context.Context, + applier *cdc.Applier, + pruner *cdc.SegmentPruner, + store *state.Store, +) error { + group, groupCtx := errgroup.WithContext(ctx) + group.Go(func() error { return pruner.Run(groupCtx) }) + group.Go(func() error { return runApplierContinuous(groupCtx, applier, store) }) + return group.Wait() +} + // awaitCatchup waits for target apply progress to reach boundary while watching // the applier that produces it. The wait polls a row on the target that only the // applier advances, so an applier exit has to end the wait as well: reading the diff --git a/internal/cdc/applier.go b/internal/cdc/applier.go index aba1c9b..56a6113 100644 --- a/internal/cdc/applier.go +++ b/internal/cdc/applier.go @@ -110,14 +110,16 @@ func (a *Applier) Run(ctx context.Context) error { // WaitUntil blocks until the authoritative target progress reaches boundary. // Only Applier.Run advances that progress, so the caller must supervise Run // concurrently and stop waiting if it exits; otherwise this polls forever. +// +// boundary must already be a transaction EndLSN. Catchup passes the durable +// watermark, which the persister published from a commit. A manual cutover LSN +// is normalized before it is stored. This wait must not scan the segment +// directory: NormalizeEndPosition reads every staged transaction, and after a +// long copy that is the whole backlog. func (a *Applier) WaitUntil(ctx context.Context, boundary LSN) error { if boundary == 0 { return nil } - // Resolve the boundary at most once. It cannot move afterwards, and - // resolving it per poll re-decoded the whole staged stream each time. - effectiveBoundary := boundary - resolved := false for { conn, err := postgres.Connect(ctx, a.config.ConnString) if err != nil { @@ -128,17 +130,7 @@ func (a *Applier) WaitUntil(ctx context.Context, boundary LSN) error { if readErr != nil { return fmt.Errorf("cdc: read catch-up progress: %w", readErr) } - if !resolved { - if durable := a.config.Durable.Load(); durable >= boundary { - resolution, err := NormalizeEndPosition(a.config.Directory, boundary, durable) - if err != nil { - return err - } - effectiveBoundary = resolution.Boundary - resolved = true - } - } - if LSN(progress) >= effectiveBoundary { + if LSN(progress) >= boundary { return nil } timer := time.NewTimer(a.config.PollInterval) diff --git a/internal/cdc/cdc_integration_test.go b/internal/cdc/cdc_integration_test.go index d090b20..e614cf6 100644 --- a/internal/cdc/cdc_integration_test.go +++ b/internal/cdc/cdc_integration_test.go @@ -7,6 +7,7 @@ import ( "errors" "fmt" "os" + "path/filepath" "slices" "strings" "sync" @@ -224,6 +225,7 @@ func TestPG17LiveWALStageApplyCrashRetry(t *testing.T) { pruner, err := NewSegmentPruner(SegmentPrunerConfig{ Directory: directory, Interval: time.Nanosecond, + Catalog: writer.SegmentCatalog(), }) if err != nil { t.Fatal(err) @@ -460,6 +462,176 @@ func (c *collectingSampler) all() []string { return slices.Compact(seen) } +// TestPG17WaitUntilDoesNotScanStagedSegments is the catchup wait: the boundary +// is already the durable EndLSN. Scanning the segment directory to "normalize" +// it would decode the whole backlog after a long copy, which is how a shard +// sat on the first file at the memory ceiling with apply idle. +func TestPG17WaitUntilDoesNotScanStagedSegments(t *testing.T) { + target := pgtest.Start(t, 17) + ctx := context.Background() + conn := target.Connect(t) + if err := postgres.EnsureProgressTable(ctx, conn); err != nil { + t.Fatal(err) + } + if err := postgres.UpdateProgress(ctx, conn, "catchup", 0x100); err != nil { + t.Fatal(err) + } + watermark := new(DurableWatermark) + watermark.Publish(0x100) + applier, err := NewApplier(ApplierConfig{ + ConnString: target.URI, + Directory: filepath.Join(t.TempDir(), "empty-cdc"), + StreamID: "catchup", + Durable: watermark, + PollInterval: 5 * time.Millisecond, + }) + if err != nil { + t.Fatal(err) + } + waitCtx, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + if err := applier.WaitUntil(waitCtx, 0x100); err != nil { + t.Fatal(err) + } +} + +func TestPG17ApplierStartsBeforeReadingUnappliedSuffix(t *testing.T) { + target := pgtest.Start(t, 17) + ctx := context.Background() + targetSQL := target.Connect(t) + if _, err := targetSQL.Exec(ctx, ` + CREATE TABLE public.catchup_probe ( + id bigint PRIMARY KEY, + value text NOT NULL + )`); err != nil { + t.Fatal(err) + } + + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + var lastEnd LSN + for i := 1; i <= 256; i++ { + value := fmt.Sprint(i) + row := Tuple{ + {Kind: DatumText, Data: []byte(value)}, + {Kind: DatumText, Data: []byte("value-" + value)}, + } + transaction := Transaction{ + CommitLSN: LSN(i * 0x10), + EndLSN: LSN(i*0x10 + 1), + CommitTime: time.Unix(int64(i), 0).UTC(), + Relations: []Relation{{ + OID: 4242, + Namespace: "public", + Name: "catchup_probe", + ReplicaIdentity: 'd', + Columns: []Column{ + {Name: "id", Type: 20, Flags: 1}, + {Name: "value", Type: 25}, + }, + }}, + Changes: []Change{{ + RelationOID: 4242, + Kind: ChangeInsert, + New: &row, + }}, + } + lastEnd = transaction.EndLSN + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + ranges := writer.SegmentCatalog().snapshot() + last := ranges[len(ranges)-1] + file, err := os.OpenFile(last.Path, os.O_WRONLY, 0) + if err != nil { + t.Fatal(err) + } + if _, err := file.WriteAt([]byte{0xff}, frameHeaderSize); err != nil { + _ = file.Close() + t.Fatal(err) + } + if err := file.Close(); err != nil { + t.Fatal(err) + } + + const streamID = "catchup-starts-before-suffix" + const generation = "generation-1" + if err := EnsureStreamProgressIdentity(ctx, targetSQL, StreamIdentityConfig{ + StreamID: streamID, Generation: generation, FreshSetup: true, + }); err != nil { + t.Fatal(err) + } + if err := postgres.UpdateProgress(ctx, targetSQL, streamID, 0); err != nil { + t.Fatal(err) + } + if err := EnsureStreamProgressIdentity(ctx, targetSQL, StreamIdentityConfig{ + StreamID: streamID, Generation: generation, FreshSetup: true, + }); err != nil { + t.Fatal(err) + } + watermark := new(DurableWatermark) + watermark.Publish(lastEnd) + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Nanosecond, + Catalog: writer.SegmentCatalog(), + }) + if err != nil { + t.Fatal(err) + } + applier, err := NewApplier(ApplierConfig{ + ConnString: target.URI, + Directory: directory, + StreamID: streamID, + StreamGeneration: generation, + TargetHasCopiedData: true, + Durable: watermark, + PollInterval: time.Millisecond, + AfterProgress: pruner.OnProgress, + ReaderSpillDirectory: filepath.Join(t.TempDir(), "reader-spill"), + }) + if err != nil { + t.Fatal(err) + } + applyCtx, stop := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- applier.Run(applyCtx) }() + deadline := time.Now().Add(30 * time.Second) + var progress pglogrepl.LSN + for time.Now().Before(deadline) { + progress, _, err = postgres.ReadProgress(ctx, targetSQL, streamID) + if err != nil { + stop() + t.Fatal(err) + } + if progress > 0 { + break + } + select { + case err := <-done: + stop() + t.Fatalf("applier reached corrupt suffix before first progress: %v", err) + default: + } + time.Sleep(time.Millisecond) + } + if progress == 0 { + stop() + t.Fatal("target progress did not move before the unapplied suffix") + } + stop() + if err := <-done; err != nil && !errors.Is(err, context.Canceled) { + t.Fatal(err) + } +} + func waitFor(t testing.TB, timeout time.Duration, condition func() bool) { t.Helper() deadline := time.Now().Add(timeout) diff --git a/internal/cdc/lifecycle_test.go b/internal/cdc/lifecycle_test.go index 65796ad..453946c 100644 --- a/internal/cdc/lifecycle_test.go +++ b/internal/cdc/lifecycle_test.go @@ -2,9 +2,12 @@ package cdc import ( "context" + "encoding/binary" "errors" + "hash/crc32" "os" "path/filepath" + "strings" "testing" "time" ) @@ -66,6 +69,312 @@ func TestSegmentPrunerBoundsSustainedAppliedSegments(t *testing.T) { } } +func TestCatalogPrunerDeletesInBatchesAndKeepsSafety(t *testing.T) { + t.Parallel() + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + var lastEnd LSN + for i := 1; i <= 5; i++ { + transaction := testTransaction(LSN(i * 0x10)) + lastEnd = transaction.EndLSN + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Hour, + Catalog: writer.SegmentCatalog(), + MaxRemoveSegments: 2, + }) + if err != nil { + t.Fatal(err) + } + if err := pruner.OnProgress(context.Background(), lastEnd+1); err != nil { + t.Fatal(err) + } + if got := len(writer.SegmentCatalog().snapshot()); got != 3 { + t.Fatalf("catalog after first batch=%d, want 3", got) + } + if err := pruner.OnProgress(context.Background(), lastEnd+1); err != nil { + t.Fatal(err) + } + ranges := writer.SegmentCatalog().snapshot() + if len(ranges) != 1 || ranges[0].LastEnd != lastEnd { + t.Fatalf("catalog after second batch=%#v, want last safety segment", ranges) + } + segments, err := listSegments(directory) + if err != nil { + t.Fatal(err) + } + if len(segments) != 1 || segments[0].start != 0x50 { + t.Fatalf("remaining segments=%#v, want safety segment 0x50", segments) + } + if _, err := Recover(directory); err != nil { + t.Fatalf("recover after catalog prune: %v", err) + } +} + +func TestSegmentPrunerRunDrainsBatchesWithoutNewProgress(t *testing.T) { + t.Parallel() + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + var lastEnd LSN + for i := 1; i <= 7; i++ { + transaction := testTransaction(LSN(i * 0x10)) + lastEnd = transaction.EndLSN + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Hour, + Catalog: writer.SegmentCatalog(), + MaxRemoveSegments: 2, + }) + if err != nil { + t.Fatal(err) + } + if err := pruner.OnProgress(context.Background(), lastEnd+1); err != nil { + t.Fatal(err) + } + if got := len(writer.SegmentCatalog().snapshot()); got != 5 { + t.Fatalf("catalog after callback batch=%d, want 5", got) + } + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- pruner.Run(ctx) }() + deadline := time.Now().Add(2 * time.Second) + for len(writer.SegmentCatalog().snapshot()) != 1 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if got := len(writer.SegmentCatalog().snapshot()); got != 1 { + cancel() + t.Fatalf("catalog after background batches=%d, want 1", got) + } + cancel() + if err := <-done; !errors.Is(err, context.Canceled) { + t.Fatalf("pruner Run error=%v, want context cancellation", err) + } + segments, err := listSegments(directory) + if err != nil { + t.Fatal(err) + } + if len(segments) != 1 || segments[0].start != 0x70 { + t.Fatalf("remaining segments=%#v, want final safety segment", segments) + } +} + +func TestCatalogPrunerRejectsSegmentChangedAfterValidation(t *testing.T) { + t.Parallel() + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + for i := 1; i <= 3; i++ { + transaction := testTransaction(LSN(i * 0x10)) + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + ranges := writer.SegmentCatalog().snapshot() + file, err := os.OpenFile(ranges[0].Path, os.O_WRONLY|os.O_APPEND, 0) + if err != nil { + t.Fatal(err) + } + if _, err := file.Write([]byte{0}); err != nil { + _ = file.Close() + t.Fatal(err) + } + if err := file.Close(); err != nil { + t.Fatal(err) + } + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Nanosecond, + Catalog: writer.SegmentCatalog(), + }) + if err != nil { + t.Fatal(err) + } + err = pruner.OnProgress(context.Background(), 0x40) + if err == nil || !strings.Contains(err.Error(), "changed after validation") { + t.Fatalf("changed segment prune error=%v", err) + } + if got := len(writer.SegmentCatalog().snapshot()); got != 3 { + t.Fatalf("catalog after rejected prune=%d, want 3", got) + } +} + +func TestCatalogPrunerRefusesAnUnexplainedMissingSegment(t *testing.T) { + t.Parallel() + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + for i := 1; i <= 3; i++ { + transaction := testTransaction(LSN(i * 0x10)) + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + ranges := writer.SegmentCatalog().snapshot() + if err := os.Remove(ranges[0].Path); err != nil { + t.Fatal(err) + } + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Nanosecond, + Catalog: writer.SegmentCatalog(), + }) + if err != nil { + t.Fatal(err) + } + err = pruner.OnProgress(context.Background(), 0x40) + if err == nil || !strings.Contains(err.Error(), "missing from disk") { + t.Fatalf("missing segment prune error=%v", err) + } + if got := len(writer.SegmentCatalog().snapshot()); got != 3 { + t.Fatalf("catalog after missing segment=%d, want 3", got) + } + if _, err := os.Stat(ranges[1].Path); err != nil { + t.Fatalf("pruner deleted another segment after missing file: %v", err) + } +} + +func TestCatalogPrunerRevalidatesAChangedFinalizedSegment(t *testing.T) { + t.Parallel() + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + for i := 1; i <= 3; i++ { + transaction := testTransaction(LSN(i * 0x10)) + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + first := writer.SegmentCatalog().snapshot()[0] + extra := testTransaction(0x15) + extra.EndLSN = 0x16 + payload, err := MarshalTransaction(&extra) + if err != nil { + t.Fatal(err) + } + var header [frameHeaderSize]byte + binary.LittleEndian.PutUint32(header[:4], uint32(len(payload))) + binary.LittleEndian.PutUint32(header[4:], crc32.Checksum(payload, castagnoliTable)) + file, err := os.OpenFile(first.Path, os.O_WRONLY|os.O_APPEND, 0) + if err != nil { + t.Fatal(err) + } + if err := writeFull(file, header[:]); err != nil { + _ = file.Close() + t.Fatal(err) + } + if err := writeFull(file, payload); err != nil { + _ = file.Close() + t.Fatal(err) + } + if err := file.Close(); err != nil { + t.Fatal(err) + } + + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Hour, + Catalog: writer.SegmentCatalog(), + }) + if err != nil { + t.Fatal(err) + } + if err := pruner.OnProgress(context.Background(), 0x40); err != nil { + t.Fatal(err) + } + ranges := writer.SegmentCatalog().snapshot() + if len(ranges) != 3 || ranges[0].LastCommit != 0x15 || ranges[0].LastEnd != 0x16 { + t.Fatalf("refreshed catalog=%#v", ranges) + } + if err := pruner.OnProgress(context.Background(), 0x40); err != nil { + t.Fatal(err) + } + ranges = writer.SegmentCatalog().snapshot() + if len(ranges) != 1 || ranges[0].StartCommit != 0x30 { + t.Fatalf("catalog after refreshed prune=%#v", ranges) + } + if _, err := Recover(directory); err != nil { + t.Fatalf("recover after refreshed prune: %v", err) + } +} + +func TestCatalogPruneAllowsOpenReaderToReachSafetySegment(t *testing.T) { + t.Parallel() + directory := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: directory, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + for i := 1; i <= 4; i++ { + transaction := testTransaction(LSN(i * 0x10)) + if err := writer.Append(&transaction); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + reader, err := NewReader(directory, 0x10, writer.DurableEndLSN()) + if err != nil { + t.Fatal(err) + } + defer reader.Close() + transaction, err := reader.Next() + if err != nil || transaction.CommitLSN != 0x20 { + t.Fatalf("first reader transaction=%x err=%v, want 20", transaction.CommitLSN, err) + } + pruner, err := NewSegmentPruner(SegmentPrunerConfig{ + Directory: directory, + Interval: time.Nanosecond, + Catalog: writer.SegmentCatalog(), + }) + if err != nil { + t.Fatal(err) + } + if err := pruner.OnProgress(context.Background(), 0x40); err != nil { + t.Fatal(err) + } + transaction, err = reader.Next() + if err != nil || transaction.CommitLSN != 0x30 { + t.Fatalf("reader after unlink=%x err=%v, want safety transaction 30", transaction.CommitLSN, err) + } +} + func TestSegmentPrunerErrorsAreActionable(t *testing.T) { t.Parallel() directory := t.TempDir() diff --git a/internal/cdc/maintenance.go b/internal/cdc/maintenance.go index c10cd52..6d0eb67 100644 --- a/internal/cdc/maintenance.go +++ b/internal/cdc/maintenance.go @@ -4,6 +4,8 @@ import ( "context" "errors" "fmt" + "os" + "path/filepath" "sync" "time" ) @@ -13,8 +15,10 @@ import ( type ProgressCallback func(context.Context, LSN) error type SegmentPrunerConfig struct { - Directory string - Interval time.Duration + Directory string + Interval time.Duration + Catalog *SegmentCatalog + MaxRemoveSegments int } // SegmentPruner provides an applier progress callback that periodically prunes @@ -22,8 +26,12 @@ type SegmentPrunerConfig struct { type SegmentPruner struct { directory string interval time.Duration + catalog *SegmentCatalog + maxRemove int mu sync.Mutex next time.Time + latest LSN + wake chan struct{} } func NewSegmentPruner(config SegmentPrunerConfig) (*SegmentPruner, error) { @@ -36,20 +44,207 @@ func NewSegmentPruner(config SegmentPrunerConfig) (*SegmentPruner, error) { if config.Interval == 0 { config.Interval = time.Minute } - return &SegmentPruner{directory: config.Directory, interval: config.Interval}, nil + if config.MaxRemoveSegments < 0 { + return nil, errors.New("cdc: segment prune removal batch must not be negative") + } + if config.MaxRemoveSegments == 0 { + config.MaxRemoveSegments = 128 + } + if config.Catalog != nil && filepath.Clean(config.Catalog.directory) != filepath.Clean(config.Directory) { + return nil, errors.New("cdc: segment pruner catalog directory does not match") + } + return &SegmentPruner{ + directory: config.Directory, + interval: config.Interval, + catalog: config.Catalog, + maxRemove: config.MaxRemoveSegments, + wake: make(chan struct{}, 1), + }, nil } // OnProgress is suitable for ApplierConfig.AfterProgress. func (p *SegmentPruner) OnProgress(_ context.Context, applied LSN) error { p.mu.Lock() defer p.mu.Unlock() + if applied > p.latest { + p.latest = applied + } now := time.Now() if !p.next.IsZero() && now.Before(p.next) { return nil } - if _, err := Prune(p.directory, applied); err != nil { - return fmt.Errorf("cdc: prune applied segments through %x: %w", applied, err) + more, err := p.pruneLocked(p.latest, now) + if more { + p.wakeLocked() + } + return err +} + +// Run continues bounded pruning batches independently of apply progress. It +// must be supervised alongside the applier so maintenance failures stop the +// migration instead of becoming detached background errors. +func (p *SegmentPruner) Run(ctx context.Context) error { + ticker := time.NewTicker(p.interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + case <-p.wake: + } + for { + p.mu.Lock() + more, err := p.pruneLocked(p.latest, time.Now()) + p.mu.Unlock() + if err != nil { + return err + } + if !more { + break + } + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + } + } +} + +func (p *SegmentPruner) pruneLocked(applied LSN, now time.Time) (bool, error) { + var ( + more bool + err error + ) + if p.catalog != nil { + _, more, err = p.catalog.prune(applied, p.maxRemove) + } else { + _, err = Prune(p.directory, applied) + } + if err != nil { + return false, fmt.Errorf("cdc: prune applied segments through %x: %w", applied, err) + } + if more { + p.next = time.Time{} + } else { + p.next = now.Add(p.interval) + } + return more, nil +} + +func (p *SegmentPruner) wakeLocked() { + select { + case p.wake <- struct{}{}: + default: + } +} + +// prune removes a bounded oldest prefix of applied finalized segments using +// bounds already validated by recovery or writer rotation. The newest eligible +// segment is retained as the same safety segment kept by scan-based Prune. +func (c *SegmentCatalog) prune(applied LSN, maxRemove int) ([]string, bool, error) { + finalized := c.snapshot() + entries, err := os.ReadDir(c.directory) + if err != nil { + return nil, false, fmt.Errorf("read segment directory before catalog prune: %w", err) + } + present := make(map[string]struct{}, len(entries)) + for _, entry := range entries { + if entry.IsDir() { + continue + } + if _, _, ok := parseSegmentName(entry.Name()); ok { + present[filepath.Join(c.directory, entry.Name())] = struct{}{} + } + } + for _, segment := range finalized { + if _, ok := present[segment.Path]; !ok { + // A writer only adds files and this pruner is serialized, so a + // missing cataloged file is not a benign rotation race. Do not + // delete anything else from an unexplained directory state. + return nil, false, fmt.Errorf( + "cataloged segment %s is missing from disk", filepath.Base(segment.Path), + ) + } + } + eligible := 0 + for eligible < len(finalized) && finalized[eligible].LastEnd < applied { + eligible++ + } + if eligible <= 1 { + return nil, false, nil + } + candidates := finalized[:eligible-1] + more := len(candidates) > maxRemove + if more { + candidates = candidates[:maxRemove] + } + + for i, segment := range candidates { + info, err := os.Stat(segment.Path) + if err != nil { + return nil, more, fmt.Errorf("stat cataloged segment %s: %w", filepath.Base(segment.Path), err) + } + if !info.Mode().IsRegular() || info.Size() != segment.ValidatedSize { + refreshed, err := revalidateCatalogedSegment(finalized, i) + if err != nil { + return nil, more, fmt.Errorf( + "cataloged segment %s changed after validation: %w", + filepath.Base(segment.Path), err, + ) + } + if err := c.replaceFinalized(refreshed); err != nil { + return nil, more, err + } + // Recompute the eligible prefix from the refreshed catalog before + // deleting anything. This call remains conservative and bounded to + // one fallback segment scan. + return nil, true, nil + } + } + + removed := make([]string, 0, len(candidates)) + var removeErr error + for _, segment := range candidates { + if err := os.Remove(segment.Path); err != nil { + removeErr = fmt.Errorf("remove segment %s: %w", filepath.Base(segment.Path), err) + break + } + removed = append(removed, segment.Path) + } + if len(removed) == 0 { + return nil, more, removeErr + } + directoryErr := syncDirectory(c.directory) + if directoryErr == nil { + c.removeFinalized(removed) + } + if removeErr != nil || directoryErr != nil { + return removed, more, errors.Join(removeErr, directoryErr) + } + return removed, more, nil +} + +func revalidateCatalogedSegment(finalized []SegmentRange, index int) (SegmentRange, error) { + if index < 0 || index >= len(finalized) { + return SegmentRange{}, errors.New("cdc: catalog segment index is out of range") + } + var previousCommit, previousEnd LSN + if index > 0 { + previousCommit = finalized[index-1].LastCommit + previousEnd = finalized[index-1].LastEnd + } + segment := finalized[index] + scan, err := scanSegment(segment.Path, false, previousCommit, previousEnd, nil) + if err != nil { + return SegmentRange{}, err } - p.next = now.Add(p.interval) - return nil + return SegmentRange{ + Path: segment.Path, + StartCommit: segment.StartCommit, + LastCommit: scan.lastCommitLSN, + LastEnd: scan.lastEndLSN, + ValidatedSize: scan.size, + }, nil } diff --git a/internal/cdc/reader.go b/internal/cdc/reader.go index 2b758a8..79af8ea 100644 --- a/internal/cdc/reader.go +++ b/internal/cdc/reader.go @@ -397,7 +397,11 @@ func Prune(directory string, appliedLSN LSN) ([]string, error) { previousEnd = scan.lastEndLSN if scan.lastEndLSN < appliedLSN { candidates = append(candidates, segment.path) + continue } + // EndLSNs are globally monotonic, so every later segment is also + // ineligible. This cold-path fallback must not read an unapplied suffix. + break } if len(candidates) <= 1 { return nil, nil diff --git a/internal/cdc/segment.go b/internal/cdc/segment.go index 662fda5..9a7a298 100644 --- a/internal/cdc/segment.go +++ b/internal/cdc/segment.go @@ -35,6 +35,7 @@ type Writer struct { dir string rotationBytes int64 + catalog *SegmentCatalog file *os.File partialPath string size int64 @@ -57,6 +58,116 @@ type Recovery struct { TruncatedBytes int64 } +// SegmentRange is the immutable LSN range of one finalized segment. It is +// populated only after the segment has been validated by recovery or durably +// finalized by Writer. +type SegmentRange struct { + Path string + StartCommit LSN + LastCommit LSN + LastEnd LSN + ValidatedSize int64 +} + +// SegmentCatalog caches validated finalized-segment bounds for pruning. Disk is +// authoritative: OpenWriter rebuilds the catalog through recovery after every +// restart, and the mutable partial tail is never included. +type SegmentCatalog struct { + mu sync.RWMutex + directory string + finalized []SegmentRange +} + +func newSegmentCatalog(directory string, finalized []SegmentRange) *SegmentCatalog { + return &SegmentCatalog{ + directory: directory, + finalized: append([]SegmentRange(nil), finalized...), + } +} + +func (c *SegmentCatalog) snapshot() []SegmentRange { + if c == nil { + return nil + } + c.mu.RLock() + defer c.mu.RUnlock() + return append([]SegmentRange(nil), c.finalized...) +} + +func (c *SegmentCatalog) addFinalized(segment SegmentRange) error { + if c == nil { + return errors.New("cdc: finalized segment catalog is missing") + } + c.mu.Lock() + defer c.mu.Unlock() + if len(c.finalized) != 0 { + previous := c.finalized[len(c.finalized)-1] + if segment.StartCommit <= previous.StartCommit || + segment.LastCommit <= previous.LastCommit || + segment.LastEnd <= previous.LastEnd { + return fmt.Errorf( + "cdc: finalized segment %s range %x/%x does not follow %s range %x/%x", + filepath.Base(segment.Path), segment.LastCommit, segment.LastEnd, + filepath.Base(previous.Path), previous.LastCommit, previous.LastEnd, + ) + } + } + c.finalized = append(c.finalized, segment) + return nil +} + +func (c *SegmentCatalog) removeFinalized(paths []string) { + if c == nil || len(paths) == 0 { + return + } + removed := make(map[string]struct{}, len(paths)) + for _, path := range paths { + removed[path] = struct{}{} + } + c.mu.Lock() + defer c.mu.Unlock() + kept := c.finalized[:0] + for _, segment := range c.finalized { + if _, ok := removed[segment.Path]; !ok { + kept = append(kept, segment) + } + } + c.finalized = kept +} + +func (c *SegmentCatalog) replaceFinalized(replacement SegmentRange) error { + if c == nil { + return errors.New("cdc: finalized segment catalog is missing") + } + c.mu.Lock() + defer c.mu.Unlock() + for i := range c.finalized { + if c.finalized[i].Path != replacement.Path { + continue + } + if replacement.StartCommit != c.finalized[i].StartCommit { + return fmt.Errorf("cdc: cataloged segment %s changed start LSN", filepath.Base(replacement.Path)) + } + if i > 0 { + previous := c.finalized[i-1] + if replacement.LastCommit <= previous.LastCommit || replacement.LastEnd <= previous.LastEnd { + return fmt.Errorf("cdc: refreshed segment %s no longer follows %s", + filepath.Base(replacement.Path), filepath.Base(previous.Path)) + } + } + if i+1 < len(c.finalized) { + next := c.finalized[i+1] + if replacement.LastCommit >= next.LastCommit || replacement.LastEnd >= next.LastEnd { + return fmt.Errorf("cdc: refreshed segment %s no longer precedes %s", + filepath.Base(replacement.Path), filepath.Base(next.Path)) + } + } + c.finalized[i] = replacement + return nil + } + return fmt.Errorf("cdc: refreshed segment %s is absent from catalog", filepath.Base(replacement.Path)) +} + // OpenWriter recovers the segment directory and opens its partial tail. func OpenWriter(config WriterConfig) (*Writer, Recovery, error) { if config.Directory == "" { @@ -77,7 +188,7 @@ func OpenWriter(config WriterConfig) (*Writer, Recovery, error) { if err := mkdirAllDurable(config.Directory, 0o750); err != nil { return nil, Recovery{}, err } - recovery, err := Recover(config.Directory) + recovery, finalized, err := recoverDirectory(config.Directory) if err != nil { return nil, Recovery{}, err } @@ -85,6 +196,7 @@ func OpenWriter(config WriterConfig) (*Writer, Recovery, error) { w := &Writer{ dir: config.Directory, rotationBytes: config.RotationBytes, + catalog: newSegmentCatalog(config.Directory, finalized), lastCommitLSN: recovery.LastCommitLSN, lastEndLSN: recovery.DurableLSN, pendingEndLSN: recovery.DurableLSN, @@ -112,6 +224,15 @@ func OpenWriter(config WriterConfig) (*Writer, Recovery, error) { return w, recovery, nil } +// SegmentCatalog returns the writer's validated finalized-segment catalog. +// The catalog remains current as this writer rotates new segments. +func (w *Writer) SegmentCatalog() *SegmentCatalog { + if w == nil { + return nil + } + return w.catalog +} + // Append writes one complete transaction frame. It does not make the // transaction durable until Sync is called, unless it triggers rotation. func (w *Writer) Append(tx *Transaction) error { @@ -313,6 +434,12 @@ func (w *Writer) rotateLocked() error { if w.file == nil { return nil } + finalized := SegmentRange{ + StartCommit: segmentStart(w.partialPath), + LastCommit: w.lastCommitLSN, + LastEnd: w.lastEndLSN, + ValidatedSize: w.size, + } if err := w.syncLocked(); err != nil { return err } @@ -337,9 +464,21 @@ func (w *Writer) rotateLocked() error { if err := w.directorySync(w.dir); err != nil { return err } + finalized.Path = finalPath + if err := w.catalog.addFinalized(finalized); err != nil { + return err + } return nil } +func segmentStart(path string) LSN { + start, _, ok := parseSegmentName(filepath.Base(path)) + if !ok { + return 0 + } + return start +} + // Finalize fsyncs and renames the current tail to .seg. func (w *Writer) Finalize() error { return w.Rotate() @@ -427,31 +566,46 @@ func mkdirAllDurable(directory string, mode os.FileMode) error { // Recover verifies finalized segments and truncates a torn or corrupt partial // tail at its first invalid frame. func Recover(directory string) (Recovery, error) { + recovery, _, err := recoverDirectory(directory) + return recovery, err +} + +func recoverDirectory(directory string) (Recovery, []SegmentRange, error) { segments, err := listSegments(directory) if err != nil { if errors.Is(err, os.ErrNotExist) { - return Recovery{}, nil + return Recovery{}, nil, nil } - return Recovery{}, err + return Recovery{}, nil, err } var result Recovery + finalized := make([]SegmentRange, 0, len(segments)) var previousCommit LSN var previousEnd LSN partialSeen := false for _, segment := range segments { if segment.partial { if partialSeen { - return Recovery{}, errors.New("cdc: multiple partial segments") + return Recovery{}, nil, errors.New("cdc: multiple partial segments") } partialSeen = true result.PartialPath = segment.path } else if partialSeen { - return Recovery{}, errors.New("cdc: finalized segment follows partial tail") + return Recovery{}, nil, errors.New("cdc: finalized segment follows partial tail") } scan, scanErr := scanSegment(segment.path, segment.partial, previousCommit, previousEnd, nil) if scanErr != nil { - return Recovery{}, scanErr + return Recovery{}, nil, scanErr + } + if !segment.partial { + finalized = append(finalized, SegmentRange{ + Path: segment.path, + StartCommit: segment.start, + LastCommit: scan.lastCommitLSN, + LastEnd: scan.lastEndLSN, + ValidatedSize: scan.size, + }) } previousCommit = scan.lastCommitLSN previousEnd = scan.lastEndLSN @@ -459,7 +613,7 @@ func Recover(directory string) (Recovery, error) { } result.LastCommitLSN = previousCommit result.DurableLSN = previousEnd - return result, nil + return result, finalized, nil } type segmentFile struct { diff --git a/internal/cdc/segment_test.go b/internal/cdc/segment_test.go index 3e8815c..07c5bbf 100644 --- a/internal/cdc/segment_test.go +++ b/internal/cdc/segment_test.go @@ -466,6 +466,106 @@ func TestRotationAndReaderAfterLSN(t *testing.T) { } } +func TestSegmentCatalogTracksRotationAndRecovery(t *testing.T) { + t.Parallel() + dir := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: dir, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + for _, pair := range []struct { + commit LSN + end LSN + }{{0x10, 0x18}, {0x20, 0x2f}, {0x30, 0x45}} { + tx := testTransaction(pair.commit) + tx.EndLSN = pair.end + if err := writer.Append(&tx); err != nil { + t.Fatal(err) + } + } + ranges := writer.SegmentCatalog().snapshot() + if len(ranges) != 3 { + t.Fatalf("catalog ranges=%d, want 3", len(ranges)) + } + for i, want := range []struct { + start LSN + end LSN + }{{0x10, 0x18}, {0x20, 0x2f}, {0x30, 0x45}} { + if ranges[i].StartCommit != want.start || ranges[i].LastEnd != want.end { + t.Fatalf("catalog range %d=%x/%x, want %x/%x", + i, ranges[i].StartCommit, ranges[i].LastEnd, want.start, want.end) + } + info, err := os.Stat(ranges[i].Path) + if err != nil { + t.Fatal(err) + } + if info.Size() != ranges[i].ValidatedSize { + t.Fatalf("catalog size %d=%d, want %d", i, ranges[i].ValidatedSize, info.Size()) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + + reopened, recovery, err := OpenWriter(WriterConfig{Directory: dir, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + defer reopened.Close() + if recovery.LastCommitLSN != 0x30 || recovery.DurableLSN != 0x45 { + t.Fatalf("recovery LSNs=%x/%x, want 30/45", recovery.LastCommitLSN, recovery.DurableLSN) + } + recovered := reopened.SegmentCatalog().snapshot() + if !reflect.DeepEqual(recovered, ranges) { + t.Fatalf("recovered catalog=%#v, want %#v", recovered, ranges) + } +} + +func TestSegmentCatalogFinalizesRecoveredPartialRange(t *testing.T) { + t.Parallel() + dir := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{ + Directory: dir, RotationBytes: int64(^uint64(0) >> 1), + }) + if err != nil { + t.Fatal(err) + } + first := testTransaction(0x10) + first.EndLSN = 0x18 + if err := writer.Append(&first); err != nil { + t.Fatal(err) + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + + reopened, recovery, err := OpenWriter(WriterConfig{ + Directory: dir, RotationBytes: int64(^uint64(0) >> 1), + }) + if err != nil { + t.Fatal(err) + } + if recovery.PartialPath == "" || len(reopened.SegmentCatalog().snapshot()) != 0 { + t.Fatalf("recovered partial=%q catalog=%#v", recovery.PartialPath, reopened.SegmentCatalog().snapshot()) + } + second := testTransaction(0x20) + second.EndLSN = 0x2f + if err := reopened.Append(&second); err != nil { + t.Fatal(err) + } + if err := reopened.Finalize(); err != nil { + t.Fatal(err) + } + ranges := reopened.SegmentCatalog().snapshot() + if len(ranges) != 1 || ranges[0].StartCommit != 0x10 || + ranges[0].LastCommit != 0x20 || ranges[0].LastEnd != 0x2f { + t.Fatalf("finalized recovered partial catalog=%#v", ranges) + } + if err := reopened.Close(); err != nil { + t.Fatal(err) + } +} + func TestRecoveryTruncatesTornTailAndContinues(t *testing.T) { t.Parallel() dir := t.TempDir() @@ -696,6 +796,47 @@ func TestPruneKeepsOneSafetySegment(t *testing.T) { } } +func TestPruneStopsBeforeUnappliedCorruptSuffix(t *testing.T) { + t.Parallel() + dir := t.TempDir() + writer, _, err := OpenWriter(WriterConfig{Directory: dir, RotationBytes: 1}) + if err != nil { + t.Fatal(err) + } + for _, lsn := range []LSN{0x10, 0x20, 0x30} { + tx := testTransaction(lsn) + if err := writer.Append(&tx); err != nil { + t.Fatal(err) + } + } + if err := writer.Close(); err != nil { + t.Fatal(err) + } + segments, err := listSegments(dir) + if err != nil { + t.Fatal(err) + } + file, err := os.OpenFile(segments[1].path, os.O_WRONLY, 0) + if err != nil { + t.Fatal(err) + } + if _, err := file.WriteAt([]byte{0xff}, frameHeaderSize); err != nil { + _ = file.Close() + t.Fatal(err) + } + if err := file.Close(); err != nil { + t.Fatal(err) + } + + removed, err := Prune(dir, 0x10) + if err != nil { + t.Fatalf("prune read unapplied corrupt suffix: %v", err) + } + if len(removed) != 0 { + t.Fatalf("pruned %d unapplied segments", len(removed)) + } +} + func TestFinalizeRenamesPartialSegment(t *testing.T) { t.Parallel() dir := t.TempDir() From 31dbd339a708472c79de97172f4aadabee6f0309 Mon Sep 17 00:00:00 2001 From: Tommaso Barbugli Date: Fri, 14 Aug 2026 15:49:32 +0200 Subject: [PATCH 2/2] build: use patched Go 1.25.13 toolchain Update the pinned standard library to resolve the vulnerabilities reported by CI. Co-authored-by: Cursor --- go.mod | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/go.mod b/go.mod index 0d5dafa..ac6e4e4 100644 --- a/go.mod +++ b/go.mod @@ -2,7 +2,7 @@ module github.com/GetStream/pgmigrate go 1.25.0 -toolchain go1.25.12 +toolchain go1.25.13 require ( github.com/jackc/pglogrepl v0.0.0-20260401131349-e37c41485510