Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 29 additions & 4 deletions internal/state/status.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
144 changes: 144 additions & 0 deletions internal/state/status_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
})
}