From d7f211ace9b39483ed8ba859fee247110c316d2c Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Sat, 26 Sep 2026 15:17:27 +1000 Subject: [PATCH 1/3] feat(copier): keyset chunker over the proven primary key Chunks are sized in rows and cut from the live table; the first is open below and the last open above so every key belongs to exactly one chunk (CO-4 coverage) and a watermark at the largest int64 means copy complete. Boundary parameters are declared bigint so bounds outside a smallint or integer key's range can be sent while the key index still serves the query. --- SAFETY.md | 2 +- docs/copy-and-swap-design.md | 7 +- docs/invariants.md | 6 +- pkg/copier/chunker.go | 241 ++++++++++++++++++++ pkg/copier/chunker_integration_test.go | 299 +++++++++++++++++++++++++ pkg/copier/chunker_test.go | 105 +++++++++ 6 files changed, 656 insertions(+), 4 deletions(-) create mode 100644 pkg/copier/chunker.go create mode 100644 pkg/copier/chunker_integration_test.go create mode 100644 pkg/copier/chunker_test.go diff --git a/SAFETY.md b/SAFETY.md index ab01189..0909dbb 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -21,7 +21,7 @@ The invariant registry (invariant IDs referenced below) lives in | `pkg/preflight` — precondition verifier, refusals | ✅ core | exists; copy-and-swap target proof declaration exists | ST-6, RF-1..RF-5 | | `pkg/executor` — bounded optimistic attempt; native concurrent index build with invalid-index recovery; native sequence executor for the safer idioms | ✅ core | exists (Phase 1: attempt-under-budget; Phase 3.1: concurrent index build; Phase 3.2: sequence executor) | LK-2 (attempt bound + the CONCURRENTLY wait-policy exception), CO-9 (qualified proof reads), ST-9 (create owner verified, never repaired) | | `pkg/checksum` — chunk verifier, continuous checker, repair | ✅ core | types and proof-type declarations exist; verifier planned | CO-1, CO-2, CO-3 | -| `pkg/copier` — shadow-table chunked copy | ✅ core | contract types exist; copier planned | CO-4, LK-3 | +| `pkg/copier` — shadow-table chunked copy | ✅ core | contract types and the keyset `Chunker` exist (row-count chunks over the proven key, first chunk open below and last open above, time-targeted sizing); copy step planned | CO-4 (chunk coverage), LK-3 | | `pkg/applier` — change apply, buffer, flush scheduling | ✅ core | package contract exists; applier planned | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/decode` — logical decoding, LSN/position accounting, per-column presence | ✅ core | contract types exist; decoder planned | ST-4, CO-4, CO-8 | | `pkg/checkpoint` — durable resume state | ✅ core | checkpoint contract exists; persistence planned | ST-1, ST-2 | diff --git a/docs/copy-and-swap-design.md b/docs/copy-and-swap-design.md index 07b2f40..9789abb 100644 --- a/docs/copy-and-swap-design.md +++ b/docs/copy-and-swap-design.md @@ -274,7 +274,10 @@ checkpoint intervals bounded. **Alternative considered → deferred.** Replica-lag-based throttling needs topology-specific observation and policy. -**Where enforced.** `pkg/copier` and `pkg/decode`; LK-3, ST-3. +**Where enforced.** `pkg/copier` (`Chunker`: chunks are sized in rows and cut by keyset from the +live table, so sparse and dense key spaces yield equal work per chunk; each timing feedback moves +the next chunk's row count toward the target by at most a factor of two, within a configured floor +and ceiling) and `pkg/decode`; LK-3, ST-3. ### D13 — Recover unique-secondary-key moves batch-wide @@ -363,7 +366,7 @@ decoding but adds write-path availability and amplification costs. | --- | --- | --- | | `pkg/dbconn` | Produces `TableLock`. | LK-1 | | `pkg/preflight` | Produces `CopySwapTarget`, the copy-and-swap route's proof (the table facts `PreflightedTable` carries plus the v1 shape, replica identity, dependent-object, decoding, and headroom checks above); owns Tier-3 refusals. | ST-6, RF-1..RF-3 | -| `pkg/copier` | Produces `Chunk` and `Watermark`. | CO-4, LK-3 | +| `pkg/copier` | Produces `Chunk` and `Watermark`; `Chunker` (built only from a `CopySwapTarget`) cuts consecutive chunks that tile the whole int64 key space — first open below, last open above — so every key a row can carry belongs to exactly one chunk and a watermark at the largest value means the copy is complete. | CO-4, LK-3 | | `pkg/checksum` | Produces `VerifiedShadow` and `CleanWatermark`; their constructors are private to this package. | CO-1, CO-2, CO-3 | | `pkg/decode` | Produces `ChangeEvent`, including per-column presence and `OldKey` for an UPDATE that moved the primary key. | ST-3, ST-4, CO-4, CO-8 | | `pkg/applier` | Applies presence-aware events from the per-key buffer. | CO-4, CO-5, CO-6, CO-8, LK-3 | diff --git a/docs/invariants.md b/docs/invariants.md index 733f59b..f392f31 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -90,7 +90,11 @@ complete unchanged-TOAST markers; the tombstone form is not available to the v1 ([copy-and-swap D13](copy-and-swap-design.md#d13--recover-unique-secondary-key-moves-batch-wide)). Full statement and the races these resolve: [low-level-design § copy and apply ordering](low-level-design.md#copy-and-apply-ordering-the-core-correctness-subtlety). -*Enforced:* copier/applier SQL shapes + flush scheduling that defers any flush overlapping an +*Enforced today:* `pkg/copier` `Chunker` — chunks are consecutive closed ranges that tile the +whole int64 key space (first open below, last open above), so every key a row can carry belongs +to exactly one chunk and "above the watermark" always names a chunk the copier will still read +(coverage, resume-from-watermark, empty-table, and cross-type key tests). *Planned enforcement:* +copier/applier SQL shapes + flush scheduling that defers any flush overlapping an in-flight chunk's key range (mutual exclusion, not tombstone retention). *Test obligation:* a marker-bearing UPDATE for a key inside an in-flight chunk asserts the flush waits for the chunk and the row is then completed from the copied shadow row, never an absent-row abort; a diff --git a/pkg/copier/chunker.go b/pkg/copier/chunker.go new file mode 100644 index 0000000..fe15e57 --- /dev/null +++ b/pkg/copier/chunker.go @@ -0,0 +1,241 @@ +package copier + +import ( + "context" + "errors" + "fmt" + "math" + "sync" + "time" + + "github.com/jackc/pgx/v5" + + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" +) + +// Chunk sizing defaults. A chunk is sized in rows, not key width, so a +// sparse key space still produces chunks of predictable work. +const ( + // DefaultTargetChunkTime is the copy duration each chunk is sized toward. + DefaultTargetChunkTime = 500 * time.Millisecond + // DefaultInitialChunkRows is the first chunk's row count, before any + // timing feedback has arrived. + DefaultInitialChunkRows int64 = 1000 + // DefaultMinChunkRows is the floor timing feedback can shrink a chunk to. + DefaultMinChunkRows int64 = 100 + // DefaultMaxChunkRows is the ceiling timing feedback can grow a chunk to. + DefaultMaxChunkRows int64 = 100_000 +) + +// A single feedback step moves the chunk size by at most this factor in +// either direction, so one anomalous chunk cannot swing the next one wildly. +const ( + maxGrowthPerStep = 2.0 + maxShrinkPerStep = 0.5 +) + +var ( + // ErrInvalidChunkerOptions reports options that cannot describe a + // bounded chunk: a non-positive target time or row count, or a floor + // above the ceiling. + ErrInvalidChunkerOptions = errors.New("invalid chunker options") + // ErrInvariantViolation is the sentinel a forged or empty proof wraps. + ErrInvariantViolation = errors.New("invariant violation") +) + +// ChunkerOptions bounds chunk sizing. Zero values take the defaults above. +type ChunkerOptions struct { + // TargetChunkTime is the copy duration each chunk is sized toward (D12). + TargetChunkTime time.Duration + // InitialRows is the first chunk's row count. + InitialRows int64 + // MinRows is the smallest chunk feedback may produce. + MinRows int64 + // MaxRows is the largest chunk feedback may produce. + MaxRows int64 +} + +func (o ChunkerOptions) withDefaults() ChunkerOptions { + if o.TargetChunkTime == 0 { + o.TargetChunkTime = DefaultTargetChunkTime + } + if o.InitialRows == 0 { + o.InitialRows = DefaultInitialChunkRows + } + if o.MinRows == 0 { + o.MinRows = DefaultMinChunkRows + } + if o.MaxRows == 0 { + o.MaxRows = DefaultMaxChunkRows + } + return o +} + +func (o ChunkerOptions) validate() error { + if o.TargetChunkTime <= 0 { + return fmt.Errorf("%w: target chunk time %s must be positive", ErrInvalidChunkerOptions, o.TargetChunkTime) + } + for _, rows := range []struct { + name string + value int64 + }{{"initial rows", o.InitialRows}, {"minimum rows", o.MinRows}, {"maximum rows", o.MaxRows}} { + if rows.value <= 0 { + return fmt.Errorf("%w: %s %d must be positive", ErrInvalidChunkerOptions, rows.name, rows.value) + } + } + if o.MinRows > o.MaxRows { + return fmt.Errorf("%w: minimum rows %d exceeds maximum rows %d", ErrInvalidChunkerOptions, o.MinRows, o.MaxRows) + } + if o.InitialRows < o.MinRows || o.InitialRows > o.MaxRows { + return fmt.Errorf("%w: initial rows %d is outside [%d, %d]", ErrInvalidChunkerOptions, o.InitialRows, o.MinRows, o.MaxRows) + } + return nil +} + +// Chunker cuts the proven table's primary-key space into consecutive closed +// ranges. Every key a row can carry belongs to exactly one chunk: the first +// chunk is open below (its lower bound is the smallest int64) and the last +// is open above (its upper bound is the largest), so a row that arrives +// under a key outside the table's bounds at the time chunking began is still +// covered by some chunk. That coverage is what lets the applier discard a +// captured change whose key lies above the copier's watermark: the copier +// will read the live row when it reaches that key (CO-4). +// +// Boundaries are cut by keyset: the upper bound of the next chunk is the +// key that makes the chunk hold exactly the current row count, read from +// the live table when Next is called. Sparse or dense key spaces therefore +// yield chunks of the same row count rather than the same key width. +// +// A Chunker is safe for concurrent use: Next serializes callers so chunks +// stay consecutive, and Feedback may arrive from any worker. +type Chunker struct { + target preflight.CopySwapTarget + opts ChunkerOptions + + mu sync.Mutex + next int64 // lower bound of the next chunk; meaningful only while !done + rows int64 // current chunk size in rows + done bool +} + +// NewChunker returns a chunker that resumes after from, or covers the whole +// key space when from is the zero watermark. It rejects a zero proof. +func NewChunker(target preflight.CopySwapTarget, from Watermark, opts ChunkerOptions) (*Chunker, error) { + // INV: ST-6 + if target.Table() == "" { + return nil, fmt.Errorf("%w: copy-and-swap target proof is empty", ErrInvariantViolation) + } + opts = opts.withDefaults() + if err := opts.validate(); err != nil { + return nil, err + } + c := &Chunker{target: target, opts: opts, rows: opts.InitialRows} + c.next, c.done = startAfter(from) + return c, nil +} + +// startAfter returns the lower bound that follows the watermark. Nothing +// copied means the key space is covered from its smallest value; a +// watermark at the largest value means nothing is left to cover. +func startAfter(from Watermark) (lower int64, done bool) { + if !from.Valid() { + return math.MinInt64, false + } + if from.Value() == math.MaxInt64 { + return 0, true + } + return from.Value() + 1, false +} + +// Rows returns the row count the next chunk will be cut to. +func (c *Chunker) Rows() int64 { + c.mu.Lock() + defer c.mu.Unlock() + return c.rows +} + +// Next cuts the next chunk from the live table with one bounded query on +// db. ok is false once the key space is covered; the last chunk returned +// before that has the largest int64 as its upper bound. +func (c *Chunker) Next(ctx context.Context, db dbconn.RowQuerier) (chunk Chunk, ok bool, err error) { + c.mu.Lock() + defer c.mu.Unlock() + if c.done { + return Chunk{}, false, nil + } + upper, err := c.boundary(ctx, db, c.next, c.rows) + if err != nil { + return Chunk{}, false, err + } + chunk, err = NewChunk(c.next, upper) + if err != nil { + return Chunk{}, false, fmt.Errorf("%w: chunk boundary: %w", ErrInvariantViolation, err) + } + if upper == math.MaxInt64 { + c.done = true + } else { + c.next = upper + 1 + } + return chunk, true, nil +} + +// boundary returns the key that closes a chunk of rows starting at lower: +// the rows-th key at or after lower. When fewer keys remain the chunk is +// the final one and is open above. +func (c *Chunker) boundary(ctx context.Context, db dbconn.RowQuerier, lower, rows int64) (int64, error) { + var upper *int64 + err := db.QueryRow(ctx, boundarySQL(c.target), lower, rows-1).Scan(&upper) + if errors.Is(err, pgx.ErrNoRows) { + return math.MaxInt64, nil + } + if err != nil { + return 0, fmt.Errorf("cut chunk boundary on %s.%s from %d: %w", c.target.Schema(), c.target.Table(), lower, err) + } + if upper == nil { + // The key column is a primary key and can never be NULL; a NULL here is + // a table that is not the one the proof described. + return 0, fmt.Errorf("%w: chunk boundary on %s.%s from %d is NULL", ErrInvariantViolation, c.target.Schema(), c.target.Table(), lower) + } + return *upper, nil +} + +// boundarySQL is the keyset boundary query: with $1 the chunk's lower bound +// and $2 one less than the chunk's row count, it returns the key that makes +// the chunk hold exactly that many rows, or no row when fewer remain. The +// parameters are declared bigint whatever the key's integer type: without +// the cast PostgreSQL would infer $1 as the key's type and a bound outside a +// smallint or integer key's range could not be sent at all. The integer +// operator family compares bigint against every integer key type, so the +// primary-key index still serves the query. +func boundarySQL(target preflight.CopySwapTarget) string { + key := pgx.Identifier{target.PKColumn()}.Sanitize() + return "SELECT " + key + + " FROM " + pgx.Identifier{target.Schema(), target.Table()}.Sanitize() + + " WHERE " + key + " >= $1::bigint" + + " ORDER BY " + key + + " OFFSET $2::bigint LIMIT 1" +} + +// Feedback reports how long the last chunk took to copy so the next one is +// sized toward the target time (D12). One step changes the size by at most +// a factor of two in either direction, within the configured floor and +// ceiling. +func (c *Chunker) Feedback(elapsed time.Duration) { + c.mu.Lock() + defer c.mu.Unlock() + c.rows = nextRows(c.rows, elapsed, c.opts) +} + +// nextRows scales rows by target/elapsed, clamped per step and to the +// configured bounds. A non-positive elapsed is treated as instantaneous, +// which grows the chunk by the maximum step. +func nextRows(rows int64, elapsed time.Duration, opts ChunkerOptions) int64 { + ratio := maxGrowthPerStep + if elapsed > 0 { + ratio = float64(opts.TargetChunkTime) / float64(elapsed) + } + ratio = math.Min(maxGrowthPerStep, math.Max(maxShrinkPerStep, ratio)) + scaled := int64(math.Round(float64(rows) * ratio)) + return min(opts.MaxRows, max(opts.MinRows, scaled)) +} diff --git a/pkg/copier/chunker_integration_test.go b/pkg/copier/chunker_integration_test.go new file mode 100644 index 0000000..4f20aec --- /dev/null +++ b/pkg/copier/chunker_integration_test.go @@ -0,0 +1,299 @@ +package copier + +import ( + "context" + "fmt" + "math" + "strings" + "testing" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/internal/testutil" + "github.com/block/pg-sprite/pkg/dbconn" + "github.com/block/pg-sprite/pkg/preflight" +) + +// chunkerFixture is a throwaway schema on a superuser pool. The superuser +// is a SET-usable member of every role, so the copy-and-swap proof the +// chunker demands is minted without provisioning. +type chunkerFixture struct { + pool *pgxpool.Pool + schema string +} + +func newChunkerFixture(t *testing.T) chunkerFixture { + t.Helper() + pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: testutil.StartPostgres(t)}) + require.NoError(t, err) + t.Cleanup(pool.Close) + return chunkerFixture{pool: pool, schema: testutil.NewSchema(t, pool)} +} + +// exec runs SQL with %s standing for the fixture schema. +func (f chunkerFixture) exec(t *testing.T, sql string) { + t.Helper() + _, err := f.pool.Exec(t.Context(), strings.ReplaceAll(sql, "%s", f.schema)) + require.NoError(t, err) +} + +// prove mints the copy-and-swap proof for table. +func (f chunkerFixture) prove(t *testing.T, table string) preflight.CopySwapTarget { + t.Helper() + role, err := preflight.CheckPrivileges(t.Context(), f.pool, f.schema, table, preflight.Requirement{Tier: preflight.TierCopyAndSwap}) + require.NoError(t, err) + target, err := preflight.CheckCopySwapShape(t.Context(), f.pool, f.schema, table, role) + require.NoError(t, err) + return target +} + +// sparseKeys is the key set every chunking test below cuts. Its gaps make +// key-width and row-count chunking give different answers, and the +// negative key proves the first chunk really is open below. +const sparseKeys = "(-5), (1), (2), (3), (10), (11), (12), (13), (20), (100), (101)" + +func (f chunkerFixture) createSparse(t *testing.T, keyType string) preflight.CopySwapTarget { + t.Helper() + f.exec(t, fmt.Sprintf(` + CREATE TABLE %%s.sparse ( + id %s PRIMARY KEY, + note text + )`, keyType)) + f.exec(t, "INSERT INTO %s.sparse (id) VALUES "+sparseKeys) + return f.prove(t, "sparse") +} + +// drain cuts chunks until the chunker reports the key space covered. +func drain(t *testing.T, c *Chunker, db dbconn.RowQuerier) []Chunk { + t.Helper() + var chunks []Chunk + for { + chunk, ok, err := c.Next(t.Context(), db) + require.NoError(t, err) + if !ok { + return chunks + } + chunks = append(chunks, chunk) + require.Less(t, len(chunks), 100, "chunking must terminate") + } +} + +// countIn is the behavioral oracle for a chunk: the rows the copy step will +// read for it, using the same bigint-parameter discipline as the chunker. +func (f chunkerFixture) countIn(t *testing.T, target preflight.CopySwapTarget, chunk Chunk) int64 { + t.Helper() + key := pgx.Identifier{target.PKColumn()}.Sanitize() + var n int64 + err := f.pool.QueryRow(t.Context(), + "SELECT count(*) FROM "+pgx.Identifier{target.Schema(), target.Table()}.Sanitize()+ + " WHERE "+key+" BETWEEN $1::bigint AND $2::bigint", + chunk.Lower(), chunk.Upper()).Scan(&n) + require.NoError(t, err) + return n +} + +func assertChunk(t *testing.T, chunk Chunk, lower, upper int64) { + t.Helper() + assert.Equal(t, lower, chunk.Lower(), "chunk lower bound") + assert.Equal(t, upper, chunk.Upper(), "chunk upper bound") +} + +// assertCoverage checks that consecutive chunks tile the whole int64 key +// space with no gap and no overlap, and that together they hold every row. +func (f chunkerFixture) assertCoverage(t *testing.T, target preflight.CopySwapTarget, chunks []Chunk, totalRows int64) { + t.Helper() + require.NotEmpty(t, chunks) + assert.Equal(t, int64(math.MinInt64), chunks[0].Lower(), "the first chunk is open below") + assert.Equal(t, int64(math.MaxInt64), chunks[len(chunks)-1].Upper(), "the last chunk is open above") + var covered int64 + for i, chunk := range chunks { + if i > 0 { + assert.Equal(t, chunks[i-1].Upper()+1, chunk.Lower(), "chunk %d starts right after chunk %d", i, i-1) + } + covered += f.countIn(t, target, chunk) + } + assert.Equal(t, totalRows, covered, "every row belongs to exactly one chunk") +} + +func TestChunkerCutsByRowCountNotKeyWidth(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, "bigint") + + c, err := NewChunker(target, Watermark{}, ChunkerOptions{InitialRows: 4, MinRows: 4, MaxRows: 4}) + require.NoError(t, err) + + chunks := drain(t, c, f.pool) + require.Len(t, chunks, 3) + // Four keys each: {-5,1,2,3} then {10,11,12,13}; the final three keys + // {20,100,101} are fewer than a chunk, so the last chunk is open above. + assertChunk(t, chunks[0], math.MinInt64, 3) + assertChunk(t, chunks[1], 4, 13) + assertChunk(t, chunks[2], 14, math.MaxInt64) + assert.Equal(t, int64(4), f.countIn(t, target, chunks[0])) + assert.Equal(t, int64(4), f.countIn(t, target, chunks[1])) + assert.Equal(t, int64(3), f.countIn(t, target, chunks[2])) + f.assertCoverage(t, target, chunks, 11) + + // Once covered, the chunker stays exhausted. + _, ok, err := c.Next(t.Context(), f.pool) + require.NoError(t, err) + assert.False(t, ok) +} + +func TestChunkerRowCountDividingTableStillClosesAbove(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, "bigint") + + // Eleven rows cut in elevens: the first chunk ends exactly on the last + // key, and a further, empty chunk is still needed so that keys inserted + // above it while the copy runs belong to some chunk. + c, err := NewChunker(target, Watermark{}, ChunkerOptions{InitialRows: 11, MinRows: 11, MaxRows: 11}) + require.NoError(t, err) + + chunks := drain(t, c, f.pool) + require.Len(t, chunks, 2) + assertChunk(t, chunks[0], math.MinInt64, 101) + assertChunk(t, chunks[1], 102, math.MaxInt64) + assert.Equal(t, int64(0), f.countIn(t, target, chunks[1])) + f.assertCoverage(t, target, chunks, 11) +} + +func TestChunkerEmptyTableIsOneOpenChunk(t *testing.T) { + f := newChunkerFixture(t) + f.exec(t, ` + CREATE TABLE %s.empty ( + id bigint PRIMARY KEY + )`) + target := f.prove(t, "empty") + + c, err := NewChunker(target, Watermark{}, ChunkerOptions{}) + require.NoError(t, err) + + chunks := drain(t, c, f.pool) + require.Len(t, chunks, 1) + assertChunk(t, chunks[0], math.MinInt64, math.MaxInt64) +} + +func TestChunkerResumesAfterWatermark(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, "bigint") + + // A watermark of 3 means the chunk ending at 3 was copied; the resumed + // chunker starts at 4 and never revisits the copied keys. + c, err := NewChunker(target, NewWatermark(3), ChunkerOptions{InitialRows: 4, MinRows: 4, MaxRows: 4}) + require.NoError(t, err) + + chunks := drain(t, c, f.pool) + require.Len(t, chunks, 2) + assertChunk(t, chunks[0], 4, 13) + assertChunk(t, chunks[1], 14, math.MaxInt64) + + // A watermark at the top of the key space means the copy is complete. + finished, err := NewChunker(target, NewWatermark(math.MaxInt64), ChunkerOptions{}) + require.NoError(t, err) + assert.Empty(t, drain(t, finished, f.pool)) +} + +func TestChunkerFeedbackResizesTheNextChunk(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, "bigint") + + c, err := NewChunker(target, Watermark{}, ChunkerOptions{InitialRows: 2, MinRows: 2, MaxRows: 8}) + require.NoError(t, err) + + first, ok, err := c.Next(t.Context(), f.pool) + require.NoError(t, err) + require.True(t, ok) + assertChunk(t, first, math.MinInt64, 1) + + // The first chunk copied in a quarter of the target time: the next one + // doubles (the per-step cap), so it holds four keys {2,3,10,11}. + c.Feedback(DefaultTargetChunkTime / 4) + assert.Equal(t, int64(4), c.Rows()) + second, ok, err := c.Next(t.Context(), f.pool) + require.NoError(t, err) + require.True(t, ok) + assertChunk(t, second, 2, 11) + + // The second chunk took three times the target: the next one shrinks by + // half to two keys {12,13}. + c.Feedback(3 * DefaultTargetChunkTime) + assert.Equal(t, int64(2), c.Rows()) + third, ok, err := c.Next(t.Context(), f.pool) + require.NoError(t, err) + require.True(t, ok) + assertChunk(t, third, 12, 13) +} + +// TestChunkerSmallKeyTypes proves the bigint parameter discipline: bounds +// far outside a smallint or integer key's range are sent without error and +// the primary-key index still serves the boundary query. +func TestChunkerSmallKeyTypes(t *testing.T) { + for _, keyType := range []string{"smallint", "integer"} { + t.Run(keyType, func(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, keyType) + + c, err := NewChunker(target, Watermark{}, ChunkerOptions{InitialRows: 4, MinRows: 4, MaxRows: 4}) + require.NoError(t, err) + chunks := drain(t, c, f.pool) + require.Len(t, chunks, 3) + assertChunk(t, chunks[0], math.MinInt64, 3) + assertChunk(t, chunks[1], 4, 13) + assertChunk(t, chunks[2], 14, math.MaxInt64) + f.assertCoverage(t, target, chunks, 11) + + f.assertBoundaryUsesPrimaryKeyIndex(t, target) + }) + } +} + +// assertBoundaryUsesPrimaryKeyIndex plans the boundary query with sequential +// scans disabled: a plan that still walks the heap would mean the bigint +// comparison cannot use the key's index, and chunking would read the whole +// table per chunk. +func (f chunkerFixture) assertBoundaryUsesPrimaryKeyIndex(t *testing.T, target preflight.CopySwapTarget) { + t.Helper() + tx, err := f.pool.Begin(t.Context()) + require.NoError(t, err) + defer func() { assert.NoError(t, tx.Rollback(context.WithoutCancel(t.Context()))) }() + + _, err = tx.Exec(t.Context(), "SET LOCAL enable_seqscan = off") + require.NoError(t, err) + rows, err := tx.Query(t.Context(), "EXPLAIN (FORMAT TEXT) "+boundarySQL(target), math.MinInt64, 3) + require.NoError(t, err) + lines, err := pgx.CollectRows(rows, pgx.RowTo[string]) + require.NoError(t, err) + plan := strings.Join(lines, "\n") + assert.Contains(t, plan, "Index", "boundary query plan:\n%s", plan) + assert.NotContains(t, plan, "Seq Scan", "boundary query plan:\n%s", plan) +} + +// TestBoundarySQLFrozen pins the boundary query text on a minted proof. The +// behavioral tests above prove what the query does; this one makes any +// change to it a deliberate, reviewed edit. +func TestBoundarySQLFrozen(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, "bigint") + + want := fmt.Sprintf(`SELECT "id" FROM "%s"."sparse" WHERE "id" >= $1::bigint ORDER BY "id" OFFSET $2::bigint LIMIT 1`, f.schema) + assert.Equal(t, want, boundarySQL(target)) +} + +func TestChunkerNextReportsQueryFailure(t *testing.T) { + f := newChunkerFixture(t) + target := f.createSparse(t, "bigint") + c, err := NewChunker(target, Watermark{}, ChunkerOptions{}) + require.NoError(t, err) + + // The proof outlives the table it described: the boundary query fails + // and the failure is returned, not swallowed as an exhausted key space. + f.exec(t, "DROP TABLE %s.sparse") + _, ok, err := c.Next(t.Context(), f.pool) + require.Error(t, err) + assert.False(t, ok) + assert.ErrorContains(t, err, "cut chunk boundary on "+f.schema+".sparse") +} diff --git a/pkg/copier/chunker_test.go b/pkg/copier/chunker_test.go new file mode 100644 index 0000000..7f86fe5 --- /dev/null +++ b/pkg/copier/chunker_test.go @@ -0,0 +1,105 @@ +package copier + +import ( + "math" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/preflight" +) + +func TestNewChunkerRejectsEmptyProof(t *testing.T) { + _, err := NewChunker(preflight.CopySwapTarget{}, Watermark{}, ChunkerOptions{}) + require.ErrorIs(t, err, ErrInvariantViolation) +} + +func TestChunkerOptionsDefaults(t *testing.T) { + opts := ChunkerOptions{}.withDefaults() + assert.Equal(t, DefaultTargetChunkTime, opts.TargetChunkTime) + assert.Equal(t, DefaultInitialChunkRows, opts.InitialRows) + assert.Equal(t, DefaultMinChunkRows, opts.MinRows) + assert.Equal(t, DefaultMaxChunkRows, opts.MaxRows) + require.NoError(t, opts.validate()) + + // An explicit value is kept; only zero takes the default. + partial := ChunkerOptions{MaxRows: 7}.withDefaults() + assert.Equal(t, int64(7), partial.MaxRows) + assert.Equal(t, DefaultMinChunkRows, partial.MinRows) +} + +func TestChunkerOptionsValidate(t *testing.T) { + valid := ChunkerOptions{TargetChunkTime: time.Second, InitialRows: 10, MinRows: 5, MaxRows: 20} + require.NoError(t, valid.validate()) + + cases := map[string]struct { + mutate func(*ChunkerOptions) + detail string + }{ + "negative target time": {func(o *ChunkerOptions) { o.TargetChunkTime = -time.Second }, "target chunk time -1s must be positive"}, + "zero initial rows": {func(o *ChunkerOptions) { o.InitialRows = 0 }, "initial rows 0 must be positive"}, + "negative min rows": {func(o *ChunkerOptions) { o.MinRows = -1 }, "minimum rows -1 must be positive"}, + "zero max rows": {func(o *ChunkerOptions) { o.MaxRows = 0 }, "maximum rows 0 must be positive"}, + "floor above ceiling": {func(o *ChunkerOptions) { o.MinRows = 21 }, "minimum rows 21 exceeds maximum rows 20"}, + "initial below floor": {func(o *ChunkerOptions) { o.InitialRows = 4 }, "initial rows 4 is outside [5, 20]"}, + "initial above ceiling": {func(o *ChunkerOptions) { o.InitialRows = 21 }, "initial rows 21 is outside [5, 20]"}, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + opts := valid + tc.mutate(&opts) + err := opts.validate() + require.ErrorIs(t, err, ErrInvalidChunkerOptions) + assert.EqualError(t, err, "invalid chunker options: "+tc.detail) + }) + } +} + +func TestStartAfter(t *testing.T) { + lower, done := startAfter(Watermark{}) + assert.Equal(t, int64(math.MinInt64), lower, "nothing copied covers the key space from its smallest value") + assert.False(t, done) + + lower, done = startAfter(NewWatermark(41)) + assert.Equal(t, int64(42), lower) + assert.False(t, done) + + // A copied-through watermark of zero is a real position, not the zero + // watermark: the next chunk starts at one. + lower, done = startAfter(NewWatermark(0)) + assert.Equal(t, int64(1), lower) + assert.False(t, done) + + // The last chunk is open above, so a watermark at the largest key means + // the copy is complete; there is no key after it to overflow into. + _, done = startAfter(NewWatermark(math.MaxInt64)) + assert.True(t, done) +} + +func TestNextRows(t *testing.T) { + opts := ChunkerOptions{TargetChunkTime: time.Second, InitialRows: 1000, MinRows: 100, MaxRows: 5000} + + // Half the target time doubles; double the target time halves. + assert.Equal(t, int64(2000), nextRows(1000, 500*time.Millisecond, opts)) + assert.Equal(t, int64(500), nextRows(1000, 2*time.Second, opts)) + // Exactly on target leaves the size alone. + assert.Equal(t, int64(1000), nextRows(1000, time.Second, opts)) + // A fractional ratio scales proportionally and rounds. + assert.Equal(t, int64(1250), nextRows(1000, 800*time.Millisecond, opts)) + + // One step never moves by more than a factor of two, however far off + // the chunk was. + assert.Equal(t, int64(2000), nextRows(1000, time.Millisecond, opts), "a very fast chunk grows by the maximum step only") + assert.Equal(t, int64(500), nextRows(1000, time.Minute, opts), "a very slow chunk shrinks by the maximum step only") + + // The configured bounds win over the step clamp. + assert.Equal(t, int64(5000), nextRows(4000, 500*time.Millisecond, opts)) + assert.Equal(t, int64(100), nextRows(150, 2*time.Second, opts)) + + // A non-positive elapsed is treated as instantaneous, not as a division + // by zero or a shrink. + assert.Equal(t, int64(2000), nextRows(1000, 0, opts)) + assert.Equal(t, int64(2000), nextRows(1000, -time.Second, opts)) +} From 4b5fce532aa7e54f24936a78a5ac673de7f10543 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Sat, 26 Sep 2026 16:33:56 +1000 Subject: [PATCH 2/3] copier: expose the cut frontier and size feedback from the timed chunk The landed watermark lags the chunks already read under concurrent workers, so the applier's CO-4 discard rule needs Chunker.Cut, not the watermark. Feedback scales the measured chunk's own size so concurrent reports agree instead of compounding; INV markers, ST-6 ids, and one-bound defaults fitting address the remaining review findings. --- SAFETY.md | 2 +- docs/copy-and-swap-design.md | 6 +- docs/invariants.md | 16 +++- pkg/copier/chunker.go | 88 +++++++++++++++---- pkg/copier/chunker_integration_test.go | 25 +++++- pkg/copier/chunker_test.go | 113 ++++++++++++++++++++++++- pkg/copier/types.go | 8 ++ pkg/copier/types_test.go | 1 + 8 files changed, 227 insertions(+), 32 deletions(-) diff --git a/SAFETY.md b/SAFETY.md index 0909dbb..01fd67c 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -21,7 +21,7 @@ The invariant registry (invariant IDs referenced below) lives in | `pkg/preflight` — precondition verifier, refusals | ✅ core | exists; copy-and-swap target proof declaration exists | ST-6, RF-1..RF-5 | | `pkg/executor` — bounded optimistic attempt; native concurrent index build with invalid-index recovery; native sequence executor for the safer idioms | ✅ core | exists (Phase 1: attempt-under-budget; Phase 3.1: concurrent index build; Phase 3.2: sequence executor) | LK-2 (attempt bound + the CONCURRENTLY wait-policy exception), CO-9 (qualified proof reads), ST-9 (create owner verified, never repaired) | | `pkg/checksum` — chunk verifier, continuous checker, repair | ✅ core | types and proof-type declarations exist; verifier planned | CO-1, CO-2, CO-3 | -| `pkg/copier` — shadow-table chunked copy | ✅ core | contract types and the keyset `Chunker` exist (row-count chunks over the proven key, first chunk open below and last open above, time-targeted sizing); copy step planned | CO-4 (chunk coverage), LK-3 | +| `pkg/copier` — shadow-table chunked copy | ✅ core | contract types and the keyset `Chunker` exist (row-count chunks over the proven key, first chunk open below and last open above, a cut frontier for the applier's discard rule, time-targeted sizing); copy step planned | CO-4 (chunk coverage), LK-3 | | `pkg/applier` — change apply, buffer, flush scheduling | ✅ core | package contract exists; applier planned | CO-4, CO-5, CO-6, CO-8, LK-3 | | `pkg/decode` — logical decoding, LSN/position accounting, per-column presence | ✅ core | contract types exist; decoder planned | ST-4, CO-4, CO-8 | | `pkg/checkpoint` — durable resume state | ✅ core | checkpoint contract exists; persistence planned | ST-1, ST-2 | diff --git a/docs/copy-and-swap-design.md b/docs/copy-and-swap-design.md index 9789abb..9971363 100644 --- a/docs/copy-and-swap-design.md +++ b/docs/copy-and-swap-design.md @@ -275,9 +275,9 @@ checkpoint intervals bounded. observation and policy. **Where enforced.** `pkg/copier` (`Chunker`: chunks are sized in rows and cut by keyset from the -live table, so sparse and dense key spaces yield equal work per chunk; each timing feedback moves -the next chunk's row count toward the target by at most a factor of two, within a configured floor -and ceiling) and `pkg/decode`; LK-3, ST-3. +live table, so sparse and dense key spaces yield equal work per chunk; each timing feedback scales +the measured chunk's own row count toward the target by at most a factor of two, within a configured +floor and ceiling, so concurrent workers' reports do not compound) and `pkg/decode`; LK-3, ST-3. ### D13 — Recover unique-secondary-key moves batch-wide diff --git a/docs/invariants.md b/docs/invariants.md index f392f31..2a56bab 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -90,12 +90,20 @@ complete unchanged-TOAST markers; the tombstone form is not available to the v1 ([copy-and-swap D13](copy-and-swap-design.md#d13--recover-unique-secondary-key-moves-batch-wide)). Full statement and the races these resolve: [low-level-design § copy and apply ordering](low-level-design.md#copy-and-apply-ordering-the-core-correctness-subtlety). +Under concurrent copy workers two positions matter and only one of them is the discard rule's: +the **cut frontier** (`Chunker.Cut`, the highest key of any chunk the copier has started reading) +and the **landed watermark** (`Watermark`, the contiguous prefix of landed chunks, which is what +is checkpointed and resumed from). Chunks land out of order, so a key can lie above the landed +watermark yet inside a chunk already read; discarding a change for it would lose it. The applier +therefore discards only for keys above the cut frontier, applies for keys in landed chunks, and +defers for keys in in-flight chunks. With one worker the two positions coincide. *Enforced today:* `pkg/copier` `Chunker` — chunks are consecutive closed ranges that tile the whole int64 key space (first open below, last open above), so every key a row can carry belongs -to exactly one chunk and "above the watermark" always names a chunk the copier will still read -(coverage, resume-from-watermark, empty-table, and cross-type key tests). *Planned enforcement:* -copier/applier SQL shapes + flush scheduling that defers any flush overlapping an -in-flight chunk's key range (mutual exclusion, not tombstone retention). *Test obligation:* a +to exactly one chunk, and `Cut` reports the frontier so that "not yet cut" always names a chunk +the copier will still read (coverage, resume-from-watermark, empty-table, frontier-after-each-cut, +and cross-type key tests). *Planned enforcement:* copier/applier SQL shapes, the in-flight chunk +registry, and flush scheduling that defers any flush overlapping an in-flight chunk's key range +(mutual exclusion, not tombstone retention). *Test obligation:* a marker-bearing UPDATE for a key inside an in-flight chunk asserts the flush waits for the chunk and the row is then completed from the copied shadow row, never an absent-row abort; a key-moving UPDATE that straddles the watermark (`UPDATE t SET id = 5000 WHERE id = 5`, watermark diff --git a/pkg/copier/chunker.go b/pkg/copier/chunker.go index fe15e57..5fd4129 100644 --- a/pkg/copier/chunker.go +++ b/pkg/copier/chunker.go @@ -44,7 +44,11 @@ var ( ErrInvariantViolation = errors.New("invariant violation") ) -// ChunkerOptions bounds chunk sizing. Zero values take the defaults above. +// ChunkerOptions bounds chunk sizing. Zero values take the defaults above, +// fitted to whatever bounds the caller did set: a floor or ceiling given on +// its own pulls the other defaults inside it, so setting one bound never +// makes the defaults contradict it. Explicitly set values are validated as +// given. type ChunkerOptions struct { // TargetChunkTime is the copy duration each chunk is sized toward (D12). TargetChunkTime time.Duration @@ -60,14 +64,26 @@ func (o ChunkerOptions) withDefaults() ChunkerOptions { if o.TargetChunkTime == 0 { o.TargetChunkTime = DefaultTargetChunkTime } - if o.InitialRows == 0 { - o.InitialRows = DefaultInitialChunkRows - } + // Bounds first, each fitted inside the other when only one was given, + // then the initial size fitted inside both. Only positive bounds are + // fitted to; a non-positive one is left for validate to refuse as given. if o.MinRows == 0 { o.MinRows = DefaultMinChunkRows + if o.MaxRows > 0 { + o.MinRows = min(o.MinRows, o.MaxRows) + } } if o.MaxRows == 0 { - o.MaxRows = DefaultMaxChunkRows + o.MaxRows = max(DefaultMaxChunkRows, o.MinRows) + } + if o.InitialRows == 0 { + o.InitialRows = DefaultInitialChunkRows + if o.MinRows > 0 { + o.InitialRows = max(o.InitialRows, o.MinRows) + } + if o.MaxRows > 0 { + o.InitialRows = min(o.InitialRows, o.MaxRows) + } } return o } @@ -98,9 +114,16 @@ func (o ChunkerOptions) validate() error { // chunk is open below (its lower bound is the smallest int64) and the last // is open above (its upper bound is the largest), so a row that arrives // under a key outside the table's bounds at the time chunking began is still -// covered by some chunk. That coverage is what lets the applier discard a -// captured change whose key lies above the copier's watermark: the copier -// will read the live row when it reaches that key (CO-4). +// covered by some chunk (CO-4 coverage). +// +// Coverage alone tells the applier which chunk a key belongs to, not whether +// the copier has read that key yet. The keys the copier will still read are +// exactly those in chunks Next has not yet returned; a chunk Next has +// returned may be in flight or landed whatever the watermark says, because +// concurrent workers land chunks out of order and the watermark advances only +// over the contiguous landed prefix. Cut reports the frontier between the two, +// so the applier discards a captured change only for a key beyond it; the +// low watermark alone is not that signal. // // Boundaries are cut by keyset: the upper bound of the next chunk is the // key that makes the chunk hold exactly the current row count, read from @@ -124,7 +147,7 @@ type Chunker struct { func NewChunker(target preflight.CopySwapTarget, from Watermark, opts ChunkerOptions) (*Chunker, error) { // INV: ST-6 if target.Table() == "" { - return nil, fmt.Errorf("%w: copy-and-swap target proof is empty", ErrInvariantViolation) + return nil, fmt.Errorf("%w (ST-6): copy-and-swap target proof is empty", ErrInvariantViolation) } opts = opts.withDefaults() if err := opts.validate(); err != nil { @@ -139,6 +162,7 @@ func NewChunker(target preflight.CopySwapTarget, from Watermark, opts ChunkerOpt // copied means the key space is covered from its smallest value; a // watermark at the largest value means nothing is left to cover. func startAfter(from Watermark) (lower int64, done bool) { + // INV: CO-4 if !from.Valid() { return math.MinInt64, false } @@ -155,6 +179,24 @@ func (c *Chunker) Rows() int64 { return c.rows } +// Cut reports the highest key inside any chunk Next has returned, or that +// the chunker resumed past; ok is false while no key has been cut. Every key +// above it lies in a chunk the copier has not started reading, so a change +// captured for such a key can be discarded (CO-4); every key at or below it +// lies in a chunk that is in flight or landed and must be applied or +// deferred, whatever the watermark says. +func (c *Chunker) Cut() (upper int64, ok bool) { + c.mu.Lock() + defer c.mu.Unlock() + if c.done { + return math.MaxInt64, true + } + if c.next == math.MinInt64 { + return 0, false + } + return c.next - 1, true +} + // Next cuts the next chunk from the live table with one bounded query on // db. ok is false once the key space is covered; the last chunk returned // before that has the largest int64 as its upper bound. @@ -170,8 +212,10 @@ func (c *Chunker) Next(ctx context.Context, db dbconn.RowQuerier) (chunk Chunk, } chunk, err = NewChunk(c.next, upper) if err != nil { - return Chunk{}, false, fmt.Errorf("%w: chunk boundary: %w", ErrInvariantViolation, err) + return Chunk{}, false, fmt.Errorf("%w (CO-4): chunk boundary: %w", ErrInvariantViolation, err) } + chunk.rows = c.rows + // INV: CO-4 if upper == math.MaxInt64 { c.done = true } else { @@ -187,15 +231,17 @@ func (c *Chunker) boundary(ctx context.Context, db dbconn.RowQuerier, lower, row var upper *int64 err := db.QueryRow(ctx, boundarySQL(c.target), lower, rows-1).Scan(&upper) if errors.Is(err, pgx.ErrNoRows) { + // INV: CO-4 return math.MaxInt64, nil } if err != nil { return 0, fmt.Errorf("cut chunk boundary on %s.%s from %d: %w", c.target.Schema(), c.target.Table(), lower, err) } if upper == nil { + // INV: ST-6 // The key column is a primary key and can never be NULL; a NULL here is // a table that is not the one the proof described. - return 0, fmt.Errorf("%w: chunk boundary on %s.%s from %d is NULL", ErrInvariantViolation, c.target.Schema(), c.target.Table(), lower) + return 0, fmt.Errorf("%w (ST-6): chunk boundary on %s.%s from %d is NULL", ErrInvariantViolation, c.target.Schema(), c.target.Table(), lower) } return *upper, nil } @@ -217,14 +263,22 @@ func boundarySQL(target preflight.CopySwapTarget) string { " OFFSET $2::bigint LIMIT 1" } -// Feedback reports how long the last chunk took to copy so the next one is -// sized toward the target time (D12). One step changes the size by at most -// a factor of two in either direction, within the configured floor and -// ceiling. -func (c *Chunker) Feedback(elapsed time.Duration) { +// Feedback reports how long chunk took to copy so the next chunk is sized +// toward the target time (D12). The new size is scaled from the row count +// chunk was cut to, not from the current size, so reports from workers +// copying concurrently each propose a size for the work they measured +// instead of compounding on one another. One step changes the size by at +// most a factor of two in either direction, within the configured floor and +// ceiling. A chunk this chunker did not cut is refused. +func (c *Chunker) Feedback(chunk Chunk, elapsed time.Duration) error { + // INV: ST-6 + if chunk.rows <= 0 { + return fmt.Errorf("%w (ST-6): feedback for chunk [%d, %d] that no chunker cut", ErrInvariantViolation, chunk.Lower(), chunk.Upper()) + } c.mu.Lock() defer c.mu.Unlock() - c.rows = nextRows(c.rows, elapsed, c.opts) + c.rows = nextRows(chunk.rows, elapsed, c.opts) + return nil } // nextRows scales rows by target/elapsed, clamped per step and to the diff --git a/pkg/copier/chunker_integration_test.go b/pkg/copier/chunker_integration_test.go index 4f20aec..6f9bf6a 100644 --- a/pkg/copier/chunker_integration_test.go +++ b/pkg/copier/chunker_integration_test.go @@ -66,7 +66,9 @@ func (f chunkerFixture) createSparse(t *testing.T, keyType string) preflight.Cop return f.prove(t, "sparse") } -// drain cuts chunks until the chunker reports the key space covered. +// drain cuts chunks until the chunker reports the key space covered, +// checking after every cut that the frontier Cut reports is the upper bound +// of the chunk just returned. func drain(t *testing.T, c *Chunker, db dbconn.RowQuerier) []Chunk { t.Helper() var chunks []Chunk @@ -76,6 +78,9 @@ func drain(t *testing.T, c *Chunker, db dbconn.RowQuerier) []Chunk { if !ok { return chunks } + cut, cutOK := c.Cut() + require.True(t, cutOK, "a returned chunk is a cut") + assert.Equal(t, chunk.Upper(), cut, "the frontier is the last returned chunk's upper bound") chunks = append(chunks, chunk) require.Less(t, len(chunks), 100, "chunking must terminate") } @@ -124,6 +129,8 @@ func TestChunkerCutsByRowCountNotKeyWidth(t *testing.T) { c, err := NewChunker(target, Watermark{}, ChunkerOptions{InitialRows: 4, MinRows: 4, MaxRows: 4}) require.NoError(t, err) + _, cutOK := c.Cut() + assert.False(t, cutOK, "nothing is cut before the first chunk") chunks := drain(t, c, f.pool) require.Len(t, chunks, 3) @@ -136,6 +143,11 @@ func TestChunkerCutsByRowCountNotKeyWidth(t *testing.T) { assert.Equal(t, int64(4), f.countIn(t, target, chunks[1])) assert.Equal(t, int64(3), f.countIn(t, target, chunks[2])) f.assertCoverage(t, target, chunks, 11) + // Every chunk records the size it was cut to, including the final one + // that holds fewer keys than that. + for i, chunk := range chunks { + assert.Equal(t, int64(4), chunk.Rows(), "chunk %d cut size", i) + } // Once covered, the chunker stays exhausted. _, ok, err := c.Next(t.Context(), f.pool) @@ -211,21 +223,28 @@ func TestChunkerFeedbackResizesTheNextChunk(t *testing.T) { // The first chunk copied in a quarter of the target time: the next one // doubles (the per-step cap), so it holds four keys {2,3,10,11}. - c.Feedback(DefaultTargetChunkTime / 4) + require.NoError(t, c.Feedback(first, DefaultTargetChunkTime/4)) assert.Equal(t, int64(4), c.Rows()) second, ok, err := c.Next(t.Context(), f.pool) require.NoError(t, err) require.True(t, ok) assertChunk(t, second, 2, 11) + assert.Equal(t, int64(4), second.Rows()) // The second chunk took three times the target: the next one shrinks by // half to two keys {12,13}. - c.Feedback(3 * DefaultTargetChunkTime) + require.NoError(t, c.Feedback(second, 3*DefaultTargetChunkTime)) assert.Equal(t, int64(2), c.Rows()) third, ok, err := c.Next(t.Context(), f.pool) require.NoError(t, err) require.True(t, ok) assertChunk(t, third, 12, 13) + + // A late report for the first chunk (two rows, a quarter of the target) + // proposes four rows from that chunk's size, not from the current two: + // feedback is about the chunk measured, whatever arrived since. + require.NoError(t, c.Feedback(first, DefaultTargetChunkTime/4)) + assert.Equal(t, int64(4), c.Rows()) } // TestChunkerSmallKeyTypes proves the bigint parameter discipline: bounds diff --git a/pkg/copier/chunker_test.go b/pkg/copier/chunker_test.go index 7f86fe5..70218a5 100644 --- a/pkg/copier/chunker_test.go +++ b/pkg/copier/chunker_test.go @@ -14,6 +14,7 @@ import ( func TestNewChunkerRejectsEmptyProof(t *testing.T) { _, err := NewChunker(preflight.CopySwapTarget{}, Watermark{}, ChunkerOptions{}) require.ErrorIs(t, err, ErrInvariantViolation) + assert.EqualError(t, err, "invariant violation (ST-6): copy-and-swap target proof is empty") } func TestChunkerOptionsDefaults(t *testing.T) { @@ -24,10 +25,53 @@ func TestChunkerOptionsDefaults(t *testing.T) { assert.Equal(t, DefaultMaxChunkRows, opts.MaxRows) require.NoError(t, opts.validate()) - // An explicit value is kept; only zero takes the default. - partial := ChunkerOptions{MaxRows: 7}.withDefaults() - assert.Equal(t, int64(7), partial.MaxRows) - assert.Equal(t, DefaultMinChunkRows, partial.MinRows) + // One bound set on its own pulls the other defaults inside it, so the + // filled options always validate. + cases := map[string]struct { + given ChunkerOptions + want ChunkerOptions + }{ + "ceiling below the default floor": { + ChunkerOptions{MaxRows: 7}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 7, MinRows: 7, MaxRows: 7}, + }, + "ceiling below the default initial size": { + ChunkerOptions{MaxRows: 500}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 500, MinRows: DefaultMinChunkRows, MaxRows: 500}, + }, + "floor above the default ceiling": { + ChunkerOptions{MinRows: 200_000}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 200_000, MinRows: 200_000, MaxRows: 200_000}, + }, + "floor above the default initial size": { + ChunkerOptions{MinRows: 5000}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 5000, MinRows: 5000, MaxRows: DefaultMaxChunkRows}, + }, + "initial size alone is kept as given": { + ChunkerOptions{InitialRows: 300}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 300, MinRows: DefaultMinChunkRows, MaxRows: DefaultMaxChunkRows}, + }, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + got := tc.given.withDefaults() + assert.Equal(t, tc.want, got) + require.NoError(t, got.validate()) + }) + } + + // An explicitly set value is never moved, even when it cannot validate: + // a floor the caller set above the ceiling the caller set is refused, + // not repaired. + contradictory := ChunkerOptions{MinRows: 50, MaxRows: 20}.withDefaults() + assert.Equal(t, int64(50), contradictory.MinRows) + assert.Equal(t, int64(20), contradictory.MaxRows) + require.ErrorIs(t, contradictory.validate(), ErrInvalidChunkerOptions) + + // A negative bound is refused as given rather than fitted around. + negative := ChunkerOptions{MaxRows: -5}.withDefaults() + assert.Equal(t, DefaultMinChunkRows, negative.MinRows) + assert.EqualError(t, negative.validate(), "invalid chunker options: maximum rows -5 must be positive") } func TestChunkerOptionsValidate(t *testing.T) { @@ -78,6 +122,67 @@ func TestStartAfter(t *testing.T) { assert.True(t, done) } +// TestChunkerCutBeforeAnyQuery covers the frontier states that need no +// table: nothing cut, resumed past a watermark, and already complete. The +// state after each Next is asserted by the integration tests. +func TestChunkerCutBeforeAnyQuery(t *testing.T) { + fresh := &Chunker{next: math.MinInt64} + _, ok := fresh.Cut() + assert.False(t, ok, "nothing has been cut before the first chunk") + + resumed := &Chunker{} + resumed.next, resumed.done = startAfter(NewWatermark(3)) + upper, ok := resumed.Cut() + assert.True(t, ok) + assert.Equal(t, int64(3), upper, "a resumed chunker has cut through its watermark") + + // A watermark at the smallest key resumes from the key after it, and + // that is a real cut, not the nothing-cut state. + lowest := &Chunker{} + lowest.next, lowest.done = startAfter(NewWatermark(math.MinInt64)) + upper, ok = lowest.Cut() + assert.True(t, ok) + assert.Equal(t, int64(math.MinInt64), upper) + + complete := &Chunker{} + complete.next, complete.done = startAfter(NewWatermark(math.MaxInt64)) + upper, ok = complete.Cut() + assert.True(t, ok) + assert.Equal(t, int64(math.MaxInt64), upper, "a complete chunker has cut the whole key space") +} + +// TestChunkerFeedbackSizesFromTheTimedChunk shows why Feedback scales the +// chunk it measured rather than the current size: several workers reporting +// the same measurement must agree on one next size, not multiply it. +func TestChunkerFeedbackSizesFromTheTimedChunk(t *testing.T) { + opts := ChunkerOptions{}.withDefaults() + c := &Chunker{opts: opts, rows: opts.InitialRows} + timed := Chunk{lower: 1, upper: 1000, rows: 1000} + + // Four workers each copied a 1000-row chunk in a fifth of the target + // time. Each report proposes 2000 (the step cap from 1000); had each + // scaled the current size the result would be 16000. + for range 4 { + require.NoError(t, c.Feedback(timed, DefaultTargetChunkTime/5)) + } + assert.Equal(t, int64(2000), c.Rows()) + + // Two slow reports on the same chunk likewise agree on 500, not 250. + for range 2 { + require.NoError(t, c.Feedback(timed, 3*DefaultTargetChunkTime)) + } + assert.Equal(t, int64(500), c.Rows()) + + // A chunk carrying no cut size was not produced by a chunker and cannot + // be sized from; the size is left alone. + foreign, err := NewChunk(1, 1000) + require.NoError(t, err) + err = c.Feedback(foreign, DefaultTargetChunkTime) + require.ErrorIs(t, err, ErrInvariantViolation) + assert.EqualError(t, err, "invariant violation (ST-6): feedback for chunk [1, 1000] that no chunker cut") + assert.Equal(t, int64(500), c.Rows()) +} + func TestNextRows(t *testing.T) { opts := ChunkerOptions{TargetChunkTime: time.Second, InitialRows: 1000, MinRows: 100, MaxRows: 5000} diff --git a/pkg/copier/types.go b/pkg/copier/types.go index b196b99..26acf8f 100644 --- a/pkg/copier/types.go +++ b/pkg/copier/types.go @@ -8,6 +8,9 @@ import "fmt" type Chunk struct { lower int64 upper int64 + // rows is the row count a Chunker cut this chunk to; zero for a chunk + // built by NewChunk, which sizes nothing. + rows int64 } // NewChunk validates and returns the closed range [lower, upper]. @@ -24,6 +27,11 @@ func (c Chunk) Lower() int64 { return c.lower } // Upper returns the inclusive upper bound. func (c Chunk) Upper() int64 { return c.upper } +// Rows returns the row count a Chunker cut the chunk to: the number of keys +// it holds, except for the final open-above chunk, which holds fewer. It is +// zero for a chunk not cut by a Chunker. +func (c Chunk) Rows() int64 { return c.rows } + // Watermark identifies the highest primary key below which every chunk was // copied. The zero value means no chunk has been copied yet; a valid // watermark is obtainable only from NewWatermark, so a value can never be diff --git a/pkg/copier/types_test.go b/pkg/copier/types_test.go index 09673d1..2816840 100644 --- a/pkg/copier/types_test.go +++ b/pkg/copier/types_test.go @@ -12,6 +12,7 @@ func TestNewChunk(t *testing.T) { require.NoError(t, err) assert.Equal(t, int64(-2), c.Lower()) assert.Equal(t, int64(4), c.Upper()) + assert.Equal(t, int64(0), c.Rows(), "a range built by hand carries no cut size") _, err = NewChunk(4, 3) assert.EqualError(t, err, "chunk lower bound 4 exceeds upper bound 3") } From 6179604b716fd6ebdc6c34c31b63ab1160389aa2 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Mon, 28 Sep 2026 20:27:59 +1000 Subject: [PATCH 3/3] copier: keep Cut and Feedback off the boundary query's lock Next held the one mutex across its database round trip, so every other worker's Cut, Rows, and Feedback queued behind it. Split the lock, fit the Min/Max defaults around an explicit InitialRows, clamp the growth before the int64 conversion, and record what a resumed copy owes CO-4. --- docs/invariants.md | 6 ++ pkg/copier/chunker.go | 85 ++++++++++++++++++++------- pkg/copier/chunker_test.go | 117 +++++++++++++++++++++++++++++++++++++ pkg/copier/doc.go | 4 +- 4 files changed, 189 insertions(+), 23 deletions(-) diff --git a/docs/invariants.md b/docs/invariants.md index 2a56bab..2b7bd73 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -97,6 +97,12 @@ is checkpointed and resumed from). Chunks land out of order, so a key can lie ab watermark yet inside a chunk already read; discarding a change for it would lose it. The applier therefore discards only for keys above the cut frontier, applies for keys in landed chunks, and defers for keys in in-flight chunks. With one worker the two positions coincide. +Resume reopens the gap between them: a run that stopped after checkpointing watermark W may have +landed chunks above W, and the resumed chunker's cut frontier starts at W, so the applier +discards changes for those keys while the copy's `ON CONFLICT DO NOTHING` would keep the stale +shadow rows. A resumed copy therefore owes one of two things before its first chunk is cut: +remove every shadow row with a key above W, or checkpoint the cut frontier alongside W and treat +(W, cut] as apply-not-discard. *Enforced today:* `pkg/copier` `Chunker` — chunks are consecutive closed ranges that tile the whole int64 key space (first open below, last open above), so every key a row can carry belongs to exactly one chunk, and `Cut` reports the frontier so that "not yet cut" always names a chunk diff --git a/pkg/copier/chunker.go b/pkg/copier/chunker.go index 5fd4129..aea0aed 100644 --- a/pkg/copier/chunker.go +++ b/pkg/copier/chunker.go @@ -40,15 +40,16 @@ var ( // bounded chunk: a non-positive target time or row count, or a floor // above the ceiling. ErrInvalidChunkerOptions = errors.New("invalid chunker options") - // ErrInvariantViolation is the sentinel a forged or empty proof wraps. - ErrInvariantViolation = errors.New("invariant violation") + // ErrInvariantViolation aliases dbconn's fail-closed error class so one + // errors.Is check covers a breach raised here or in the connection layer. + ErrInvariantViolation = dbconn.ErrInvariantViolation ) // ChunkerOptions bounds chunk sizing. Zero values take the defaults above, -// fitted to whatever bounds the caller did set: a floor or ceiling given on -// its own pulls the other defaults inside it, so setting one bound never -// makes the defaults contradict it. Explicitly set values are validated as -// given. +// fitted to whatever the caller did set: a floor, ceiling, or initial size +// given on its own pulls the other defaults around it, so setting one value +// never makes the defaults contradict it. Explicitly set values are +// validated as given. type ChunkerOptions struct { // TargetChunkTime is the copy duration each chunk is sized toward (D12). TargetChunkTime time.Duration @@ -64,17 +65,24 @@ func (o ChunkerOptions) withDefaults() ChunkerOptions { if o.TargetChunkTime == 0 { o.TargetChunkTime = DefaultTargetChunkTime } - // Bounds first, each fitted inside the other when only one was given, - // then the initial size fitted inside both. Only positive bounds are - // fitted to; a non-positive one is left for validate to refuse as given. + // Bounds first, each fitted inside the other and around an explicit + // initial size when only some were given, then the initial size fitted + // inside both. Only positive values are fitted to; a non-positive one is + // left for validate to refuse as given. if o.MinRows == 0 { o.MinRows = DefaultMinChunkRows if o.MaxRows > 0 { o.MinRows = min(o.MinRows, o.MaxRows) } + if o.InitialRows > 0 { + o.MinRows = min(o.MinRows, o.InitialRows) + } } if o.MaxRows == 0 { o.MaxRows = max(DefaultMaxChunkRows, o.MinRows) + if o.InitialRows > 0 { + o.MaxRows = max(o.MaxRows, o.InitialRows) + } } if o.InitialRows == 0 { o.InitialRows = DefaultInitialChunkRows @@ -131,11 +139,17 @@ func (o ChunkerOptions) validate() error { // yield chunks of the same row count rather than the same key width. // // A Chunker is safe for concurrent use: Next serializes callers so chunks -// stay consecutive, and Feedback may arrive from any worker. +// stay consecutive, while Cut, Rows, and Feedback never wait behind the +// boundary query a Next in progress is running. type Chunker struct { target preflight.CopySwapTarget opts ChunkerOptions + // nextMu serializes Next so consecutive chunks are cut from consecutive + // positions. It is the only lock held across the boundary query. + nextMu sync.Mutex + // mu guards the cursor state and is held only for reads and writes of it, + // never across a database round trip. mu sync.Mutex next int64 // lower bound of the next chunk; meaningful only while !done rows int64 // current chunk size in rows @@ -201,27 +215,45 @@ func (c *Chunker) Cut() (upper int64, ok bool) { // db. ok is false once the key space is covered; the last chunk returned // before that has the largest int64 as its upper bound. func (c *Chunker) Next(ctx context.Context, db dbconn.RowQuerier) (chunk Chunk, ok bool, err error) { - c.mu.Lock() - defer c.mu.Unlock() - if c.done { + c.nextMu.Lock() + defer c.nextMu.Unlock() + lower, rows, done := c.cursor() + if done { return Chunk{}, false, nil } - upper, err := c.boundary(ctx, db, c.next, c.rows) + upper, err := c.boundary(ctx, db, lower, rows) if err != nil { return Chunk{}, false, err } - chunk, err = NewChunk(c.next, upper) + chunk, err = NewChunk(lower, upper) if err != nil { return Chunk{}, false, fmt.Errorf("%w (CO-4): chunk boundary: %w", ErrInvariantViolation, err) } - chunk.rows = c.rows + chunk.rows = rows + c.advance(upper) + return chunk, true, nil +} + +// cursor snapshots the position and size the next chunk is cut with. Only +// Next moves the position, and Next is serialized, so the snapshot stays +// current for the boundary query even though mu is released while it runs. +func (c *Chunker) cursor() (lower, rows int64, done bool) { + c.mu.Lock() + defer c.mu.Unlock() + return c.next, c.rows, c.done +} + +// advance moves the position past a chunk closed at upper, or marks the key +// space covered when that chunk was open above. +func (c *Chunker) advance(upper int64) { + c.mu.Lock() + defer c.mu.Unlock() // INV: CO-4 if upper == math.MaxInt64 { c.done = true - } else { - c.next = upper + 1 + return } - return chunk, true, nil + c.next = upper + 1 } // boundary returns the key that closes a chunk of rows starting at lower: @@ -269,7 +301,8 @@ func boundarySQL(target preflight.CopySwapTarget) string { // copying concurrently each propose a size for the work they measured // instead of compounding on one another. One step changes the size by at // most a factor of two in either direction, within the configured floor and -// ceiling. A chunk this chunker did not cut is refused. +// ceiling. A chunk no chunker cut, such as a zero Chunk or one built with +// NewChunk, is refused. func (c *Chunker) Feedback(chunk Chunk, elapsed time.Duration) error { // INV: ST-6 if chunk.rows <= 0 { @@ -290,6 +323,14 @@ func nextRows(rows int64, elapsed time.Duration, opts ChunkerOptions) int64 { ratio = float64(opts.TargetChunkTime) / float64(elapsed) } ratio = math.Min(maxGrowthPerStep, math.Max(maxShrinkPerStep, ratio)) - scaled := int64(math.Round(float64(rows) * ratio)) - return min(opts.MaxRows, max(opts.MinRows, scaled)) + // Clamp before converting: a scaled value past the ceiling can also be + // past what int64 holds, and converting it first would wrap. + scaled := float64(rows) * ratio + if scaled >= float64(opts.MaxRows) { + return opts.MaxRows + } + if scaled <= float64(opts.MinRows) { + return opts.MinRows + } + return int64(math.Round(scaled)) } diff --git a/pkg/copier/chunker_test.go b/pkg/copier/chunker_test.go index 70218a5..b1b6f7f 100644 --- a/pkg/copier/chunker_test.go +++ b/pkg/copier/chunker_test.go @@ -1,10 +1,13 @@ package copier import ( + "context" "math" + "sync" "testing" "time" + "github.com/jackc/pgx/v5" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -51,6 +54,18 @@ func TestChunkerOptionsDefaults(t *testing.T) { ChunkerOptions{InitialRows: 300}, ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 300, MinRows: DefaultMinChunkRows, MaxRows: DefaultMaxChunkRows}, }, + "initial size below the default floor pulls the floor down": { + ChunkerOptions{InitialRows: 50}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 50, MinRows: 50, MaxRows: DefaultMaxChunkRows}, + }, + "initial size above the default ceiling pulls the ceiling up": { + ChunkerOptions{InitialRows: 200_000}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 200_000, MinRows: DefaultMinChunkRows, MaxRows: 200_000}, + }, + "initial size fits between an explicit ceiling and the fitted floor": { + ChunkerOptions{InitialRows: 50, MaxRows: 80}, + ChunkerOptions{TargetChunkTime: DefaultTargetChunkTime, InitialRows: 50, MinRows: 50, MaxRows: 80}, + }, } for name, tc := range cases { t.Run(name, func(t *testing.T) { @@ -183,6 +198,100 @@ func TestChunkerFeedbackSizesFromTheTimedChunk(t *testing.T) { assert.Equal(t, int64(500), c.Rows()) } +// TestChunkerCutAndFeedbackDoNotWaitForNext pins the lock split: while one +// worker's Next is out on its boundary query, the other workers' Cut, Rows, +// and Feedback calls are answered from the cursor state instead of queuing +// behind that round trip. The chunk Next finally returns still carries the +// size it was cut with, not the size Feedback moved to meanwhile. +func TestChunkerCutAndFeedbackDoNotWaitForNext(t *testing.T) { + const observeDeadline = 5 * time.Second + opts := ChunkerOptions{}.withDefaults() + c := &Chunker{opts: opts, rows: opts.InitialRows} + c.next, c.done = startAfter(NewWatermark(1000)) + db := &heldBoundary{started: make(chan struct{}), release: make(chan struct{}), upper: 2000} + releaseBoundary := sync.OnceFunc(func() { close(db.release) }) + t.Cleanup(releaseBoundary) + + var wg sync.WaitGroup + var chunk Chunk + var ok bool + var nextErr error + wg.Go(func() { chunk, ok, nextErr = c.Next(t.Context(), db) }) + select { + case <-db.started: + case <-time.After(observeDeadline): + t.Fatal("Next never reached the boundary query") + } + + observed := make(chan struct{}) + var cut int64 + var cutOK bool + var rowsDuringNext int64 + var feedbackErr error + wg.Go(func() { + defer close(observed) + cut, cutOK = c.Cut() + earlier := Chunk{lower: 1, upper: 1000, rows: opts.InitialRows} + feedbackErr = c.Feedback(earlier, DefaultTargetChunkTime/5) + rowsDuringNext = c.Rows() + }) + select { + case <-observed: + case <-time.After(observeDeadline): + t.Fatal("Cut, Rows, or Feedback waited behind Next's boundary query") + } + releaseBoundary() + wg.Wait() + + require.NoError(t, feedbackErr) + assert.True(t, cutOK) + assert.Equal(t, int64(1000), cut, "the frontier is the resumed watermark until Next lands its chunk") + assert.Equal(t, 2*opts.InitialRows, rowsDuringNext, "Feedback resized while Next was in flight") + + require.NoError(t, nextErr) + require.True(t, ok) + assert.Equal(t, int64(1001), chunk.Lower()) + assert.Equal(t, int64(2000), chunk.Upper()) + assert.Equal(t, opts.InitialRows, chunk.rows, "the chunk keeps the size it was cut with") + cut, cutOK = c.Cut() + assert.True(t, cutOK) + assert.Equal(t, int64(2000), cut) +} + +// heldBoundary is a RowQuerier whose boundary query blocks until released, +// so a test can observe the chunker while a Next is mid round trip. +type heldBoundary struct { + started chan struct{} + release chan struct{} + upper int64 +} + +func (h *heldBoundary) QueryRow(ctx context.Context, _ string, _ ...any) pgx.Row { + close(h.started) + select { + case <-h.release: + return heldRow{upper: h.upper} + case <-ctx.Done(): + return heldRow{err: ctx.Err()} + } +} + +// heldRow answers the boundary scan with one key, the way pgx does for a +// row whose single column is the chunk's upper bound. +type heldRow struct { + upper int64 + err error +} + +func (r heldRow) Scan(dest ...any) error { + if r.err != nil { + return r.err + } + upper := r.upper + *(dest[0].(**int64)) = &upper + return nil +} + func TestNextRows(t *testing.T) { opts := ChunkerOptions{TargetChunkTime: time.Second, InitialRows: 1000, MinRows: 100, MaxRows: 5000} @@ -207,4 +316,12 @@ func TestNextRows(t *testing.T) { // by zero or a shrink. assert.Equal(t, int64(2000), nextRows(1000, 0, opts)) assert.Equal(t, int64(2000), nextRows(1000, -time.Second, opts)) + + // A ceiling at the top of int64 is legal, and a fast chunk near it must + // hit the ceiling rather than wrap around to the floor. + unbounded := ChunkerOptions{TargetChunkTime: time.Second, InitialRows: 1000, MinRows: 100, MaxRows: math.MaxInt64} + assert.Equal(t, int64(math.MaxInt64), nextRows(math.MaxInt64/2+1, 500*time.Millisecond, unbounded)) + assert.Equal(t, int64(math.MaxInt64), nextRows(math.MaxInt64, time.Millisecond, unbounded)) + // float64 rounds the ceiling up to 2^63, so halving lands exactly on 2^62. + assert.Equal(t, int64(1)<<62, nextRows(math.MaxInt64, 2*time.Second, unbounded), "a slow chunk at the ceiling still halves") } diff --git a/pkg/copier/doc.go b/pkg/copier/doc.go index 1f648c1..a4b1a25 100644 --- a/pkg/copier/doc.go +++ b/pkg/copier/doc.go @@ -1,2 +1,4 @@ -// Package copier defines chunked shadow-table copy contracts enforcing CO-4 and LK-3. +// Package copier cuts a table's primary-key space into row-count chunks for +// the shadow-table copy, tiling the whole int64 range so every key belongs to +// exactly one chunk (CO-4) and copy work stays bounded per statement (LK-3). package copier