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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion SAFETY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
7 changes: 5 additions & 2 deletions docs/copy-and-swap-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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 |
Expand Down
16 changes: 14 additions & 2 deletions docs/invariants.md
Original file line number Diff line number Diff line change
Expand Up @@ -90,8 +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).
*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.
*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
Expand Down
295 changes: 295 additions & 0 deletions pkg/copier/chunker.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,295 @@
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,
// 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
// 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 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 = 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
}

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, 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 (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.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 (CO-4): chunk boundary: %w", ErrInvariantViolation, err)
}
chunk.rows = c.rows
// INV: CO-4
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) {
// 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 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(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))
scaled := int64(math.Round(float64(rows) * ratio))
return min(opts.MaxRows, max(opts.MinRows, scaled))
}
Loading
Loading