diff --git a/SAFETY.md b/SAFETY.md index ab01189..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 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, 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 07b2f40..9971363 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 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 @@ -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..2b7bd73 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -90,8 +90,26 @@ 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 -in-flight chunk's key range (mutual exclusion, not tombstone retention). *Test obligation:* a +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. +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 +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 new file mode 100644 index 0000000..aea0aed --- /dev/null +++ b/pkg/copier/chunker.go @@ -0,0 +1,336 @@ +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 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 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 + // 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 + } + // 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 + if o.MinRows > 0 { + o.InitialRows = max(o.InitialRows, o.MinRows) + } + if o.MaxRows > 0 { + o.InitialRows = min(o.InitialRows, o.MaxRows) + } + } + 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 (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 +// 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, 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 + 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 (ST-6): 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) { + // INV: CO-4 + 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 +} + +// 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. +func (c *Chunker) Next(ctx context.Context, db dbconn.RowQuerier) (chunk Chunk, ok bool, err error) { + c.nextMu.Lock() + defer c.nextMu.Unlock() + lower, rows, done := c.cursor() + if done { + return Chunk{}, false, nil + } + upper, err := c.boundary(ctx, db, lower, rows) + if err != nil { + return Chunk{}, false, err + } + chunk, err = NewChunk(lower, upper) + if err != nil { + return Chunk{}, false, fmt.Errorf("%w (CO-4): chunk boundary: %w", ErrInvariantViolation, err) + } + 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 + return + } + c.next = upper + 1 +} + +// 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) { + // 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 (ST-6): 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 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 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 { + 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(chunk.rows, elapsed, c.opts) + return nil +} + +// 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)) + // 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_integration_test.go b/pkg/copier/chunker_integration_test.go new file mode 100644 index 0000000..6f9bf6a --- /dev/null +++ b/pkg/copier/chunker_integration_test.go @@ -0,0 +1,318 @@ +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, +// 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 + for { + chunk, ok, err := c.Next(t.Context(), db) + require.NoError(t, err) + 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") + } +} + +// 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) + _, 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) + // 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) + // 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) + 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}. + 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}. + 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 +// 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..b1b6f7f --- /dev/null +++ b/pkg/copier/chunker_test.go @@ -0,0 +1,327 @@ +package copier + +import ( + "context" + "math" + "sync" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "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) + assert.EqualError(t, err, "invariant violation (ST-6): copy-and-swap target proof is empty") +} + +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()) + + // 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}, + }, + "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) { + 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) { + 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) +} + +// 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()) +} + +// 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} + + // 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)) + + // 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 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") }