From 9f992acf73e00b853a9ff9d3d582a465df5ce143 Mon Sep 17 00:00:00 2001 From: thesyncim Date: Thu, 27 Aug 2026 12:57:08 +0100 Subject: [PATCH] perf(state): keep per-message metadata reads bounded --- internal/state/status.go | 33 +++++++- internal/state/status_test.go | 144 ++++++++++++++++++++++++++++++++++ 2 files changed, 173 insertions(+), 4 deletions(-) create mode 100644 internal/state/status_test.go diff --git a/internal/state/status.go b/internal/state/status.go index 130eb12..c1e3750 100644 --- a/internal/state/status.go +++ b/internal/state/status.go @@ -103,9 +103,34 @@ func (s *Store) Snapshot(ctx context.Context) (Status, error) { // Migration returns current singleton metadata. func (s *Store) Migration(ctx context.Context) (Migration, error) { - status, err := s.Snapshot(ctx) - if err != nil { - return Migration{}, err + s.mu.Lock() + defer s.mu.Unlock() + if s.closed { + return Migration{}, ErrClosed + } + + // CDC reads the cutover boundary for every message; do not scan progress here. + var migration Migration + var createdAt, updatedAt int64 + if err := s.db.QueryRowContext( + ctx, ` + SELECT source_fingerprint, filter_fingerprint, slot_name, snapshot_name, + consistent_point, phase, end_position, created_at, updated_at + FROM migration WHERE id=1`, + ).Scan( + &migration.SourceFingerprint, + &migration.FilterFingerprint, + &migration.SlotName, + &migration.SnapshotName, + &migration.ConsistentPoint, + &migration.Phase, + &migration.EndPosition, + &createdAt, + &updatedAt, + ); err != nil { + return Migration{}, fmt.Errorf("read migration status: %w", err) } - return status.Migration, nil + migration.CreatedAt = fromUnixNano(createdAt) + migration.UpdatedAt = fromUnixNano(updatedAt) + return migration, nil } diff --git a/internal/state/status_test.go b/internal/state/status_test.go new file mode 100644 index 0000000..e046b15 --- /dev/null +++ b/internal/state/status_test.go @@ -0,0 +1,144 @@ +package state + +import ( + "context" + "database/sql" + "errors" + "testing" +) + +func TestMigrationDoesNotReadProgress(t *testing.T) { + t.Parallel() + ctx := t.Context() + store := openTestStore(t, t.TempDir()) + want, err := store.Migration(ctx) + if err != nil { + t.Fatal(err) + } + // Metadata reads must not depend on the dashboard's progress rows. + if _, err := store.db.ExecContext(ctx, "DELETE FROM apply_progress"); err != nil { + t.Fatal(err) + } + got, err := store.Migration(ctx) + if err != nil || got != want { + t.Fatalf("Migration() = %+v, %v; want %+v", got, err, want) + } + if _, err := store.Snapshot(ctx); !errors.Is(err, sql.ErrNoRows) { + t.Fatalf("Snapshot() error = %v; want missing progress row", err) + } +} + +func TestMigrationSeesControlBoundaryAndSurvivesReopen(t *testing.T) { + t.Parallel() + ctx := t.Context() + dir := t.TempDir() + store := openTestStore(t, dir) + if err := store.SetSnapshot(ctx, "slot", "snapshot", "1/A"); err != nil { + t.Fatal(err) + } + control, err := OpenControl(ctx, dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { control.Close() }) + reader, err := OpenReadOnly(ctx, dir) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { reader.Close() }) + for _, s := range []*Store{store, reader} { + before, err := s.Migration(ctx) + if err != nil || before.EndPosition != "" { + t.Fatalf("initial boundary = %+v, %v", before, err) + } + } + if err := control.SetEndPosition(ctx, "2/F"); err != nil { + t.Fatal(err) + } + status, err := control.Snapshot(ctx) + if err != nil || status.Migration.EndPosition != "2/F" { + t.Fatalf("control snapshot = %+v, %v", status, err) + } + for _, s := range []*Store{store, reader} { + got, err := s.Migration(ctx) + if err != nil || got != status.Migration { + t.Fatalf("Migration() = %+v, %v; want %+v", got, err, status.Migration) + } + } + if err := reader.Close(); err != nil { + t.Fatal(err) + } + if err := control.Close(); err != nil { + t.Fatal(err) + } + if err := store.Close(); err != nil { + t.Fatal(err) + } + got, err := openTestStore(t, dir).Migration(ctx) + if err != nil || got != status.Migration { + t.Fatalf("reopened Migration() = %+v, %v; want %+v", got, err, status.Migration) + } +} + +func TestMigrationReadErrors(t *testing.T) { + t.Parallel() + ctx := t.Context() + store := openTestStore(t, t.TempDir()) + canceled, cancel := context.WithCancel(ctx) + cancel() + if _, err := store.Migration(canceled); !errors.Is(err, context.Canceled) { + t.Fatalf("canceled Migration() error = %v", err) + } + if _, err := store.db.ExecContext(ctx, "DELETE FROM migration"); err != nil { + t.Fatal(err) + } + if _, err := store.Migration(ctx); !errors.Is(err, sql.ErrNoRows) { + t.Fatalf("missing Migration() error = %v", err) + } + if err := store.Close(); err != nil { + t.Fatal(err) + } + if _, err := store.Migration(ctx); !errors.Is(err, ErrClosed) { + t.Fatalf("closed Migration() error = %v", err) + } +} + +func BenchmarkMigrationMetadata(b *testing.B) { + ctx := b.Context() + store, err := Open(ctx, b.TempDir(), testFingerprints) + if err != nil { + b.Fatal(err) + } + b.Cleanup(func() { store.Close() }) + // Synthetic inventory with the same object counts as the c3 migration. + for _, statement := range []string{ + `WITH RECURSIVE n(x) AS (VALUES(1) UNION ALL SELECT x+1 FROM n WHERE x<113) + INSERT INTO tables(oid,schema_name,table_name,completed) SELECT x,'public','table_'||x,1 FROM n`, + `WITH RECURSIVE n(x) AS (VALUES(1) UNION ALL SELECT x+1 FROM n WHERE x<276) + INSERT INTO parts(table_oid,part_id,completed) SELECT 1,x,1 FROM n`, + `WITH RECURSIVE n(x) AS (VALUES(1) UNION ALL SELECT x+1 FROM n WHERE x<509) + INSERT INTO indexes(oid,table_oid,name,completed) SELECT x,1,'index_'||x,1 FROM n`, + `WITH RECURSIVE n(x) AS (VALUES(1) UNION ALL SELECT x+1 FROM n WHERE x<12) + INSERT INTO constraints(oid,table_oid,name,completed) SELECT x,1,'constraint_'||x,1 FROM n`, + } { + if _, err := store.db.ExecContext(ctx, statement); err != nil { + b.Fatal(err) + } + } + b.Run("Migration", func(b *testing.B) { + b.ReportAllocs() + for b.Loop() { + if _, err := store.Migration(ctx); err != nil { + b.Fatal(err) + } + } + }) + b.Run("Snapshot", func(b *testing.B) { + b.ReportAllocs() + for b.Loop() { + if _, err := store.Snapshot(ctx); err != nil { + b.Fatal(err) + } + } + }) +}