Skip to content

copier: fill the copy step's progress counters from the copier - #132

Merged
Kiran01bm merged 2 commits into
mainfrom
kiran01bm/cs5c-copier-fillers
Oct 1, 2026
Merged

Kiran01bm merged 2 commits into
mainfrom
kiran01bm/cs5c-copier-fillers

Conversation

@Kiran01bm

@Kiran01bm Kiran01bm commented Sep 29, 2026 •

Copy link
Copy Markdown
Collaborator

Makes copier.Copier the progress tracker's WorkSource for the lifetime of Run, so a poll during the copy step reports rows and bytes instead of an empty step. Built on the progress contract from #131 (merged).

Why

The tracker can now poll an engine-measured WorkSource (#131), but nothing implements one. The copier already holds every fact the step's counters need — the ledger's inserted-row count, the two relation OIDs the shadow proof carries — and it is the only component that knows when the copy is running. Wiring it here, behind one Options field, keeps the orchestrator PR to orchestration.

What

  • pkg/copier:
    • Options.Tracker *progress.Tracker. When set, Run calls SetWorkSource(c) before the resume clear and the first chunk, and StopWorkSource() (the drain fence) before it returns; the caller owns the tracker's steps and statement, the copier fills only the current step's counters. A nil tracker changes nothing.
    • Copier.Work(ctx) — rows_copied is Position.RowsInserted (rows this run's committed chunks inserted; a resumed run counts only its own); rows_total is the source's pg_class.reltuples read once at the start of Run, clamped so a never-analyzed table (-1) reports 0 rather than an unsigned wraparound; bytes_copied / bytes_total are pg_table_size of the shadow and the source measured at each poll over the caller's pool. Outside Run the copier holds no connection, so the sizes read 0 and the row counters still report what the ledger knows.
    • The row-count read is a measurement, not a relation-identity check: a source that is gone reads as "no count" (COALESCE(…, -1)), so the chunk guard's ST-6 refusal stays the one answer for a replaced or dropped source (the existing TestCopierRefusesReplacedRelations fails without this).
    • mu now also guards the bound pool and the recorded total; report/its stop closure are the only writers.
    • Every catalog name the measurements read is pg_catalog-qualified (CO-9): pg_class, pg_table_size, the = operator, and the oid cast — an unqualified pg_table_size(oid) resolves by signature across the whole search_path, so a same-signature decoy in a user schema wins regardless of path order. The chunk guard's relation-identity query (relationOIDsSQL) is qualified the same way.
    • Each measurement runs in a read-only transaction of its own under the chunk budgets (SET LOCAL lock_timeout / statement_timeout, LK-2), so a poll on a caller-built pool with no session timeouts still ends at the copier's lock timeout instead of waiting out a queued lock. setBudgets is the one place the budgets are spelled; chunk transactions and measurements share it.
  • Docs in the same PR: progress-report.md gains a table defining the four copy counters, notes that the two sizes are of tables that differ in shape (a rate is derived by the consumer from two snapshots), states that on a resumed run rows_copied / rows_total is that run's share rather than completion, and that the size read is bounded by the copy's own budgets; copy-and-swap-design.md package map, architecture.md copier and progress rows, SAFETY.md copier and progress rows drop "progress fillers planned".

Tests

Real PostgreSQL (work_integration_test.go): a 4-worker copy with one chunk pinned mid-insert and every other chunk landed is polled through the tracker — rows_copied is exactly the source count minus the pinned chunk, rows_total equals the source count after ANALYZE, and bytes_copied / bytes_total equal pg_table_size of the shadow and the source as the test measures them (nothing writes either table while the pin holds); after the pin is released and Run returns, the tracker reports no work and Copier.Work reports the full row count with zero bytes. A table created WITH (autovacuum_enabled = false) and never analyzed reports rows_total 0. A resume from watermark 500 over a 1000-row source reports rows_copied 500 and rows_total 1000. A search_path that shadows the catalog — an impostor pg_class row carrying a fake reltuples for the source OID, an empty pg_namespace, and a pg_table_size(oid) decoy that records every call — leaves the counters real and the decoy uncalled, and the copy converges; unqualifying any one of pg_class, pg_table_size, or relationOIDsSQL fails it. A poll while the resume clear is fenced behind a straggling chunk transaction reports the copier's counters (the copier is the source before the clear, not only before the first chunk). A row count the copier cannot read (a closed pool) ends Run before the clear touches the shadow, as an ordinary error rather than an invariant violation, with the tracker left unregistered. A poll on a raw pgxpool.Pool — no session timeouts — queued behind an ACCESS EXCLUSIVE request on the source ends with the server's lock timeout (55P03) inside the copier's 500 ms, not the observer's patience; removing the budgets from the measurement transaction hangs it. Pure tests (work_test.go) pin the ledger-to-Work mapping and the zero state before Run. Seven mutants — clamp removed, tracker never registered, stop fence skipped, sizes swapped, pool kept after Run, rows dropped, missing relation failing Run — each fail at least one test.

Before / after

Before (#131 alone)                        After
┌─────────┐ Progress()                     ┌─────────┐ Progress()          ┌──────────────────┐
│ Tracker │──▶ source? ── none ──▶ work    │ Tracker │──▶ source ─────────▶│ Copier.Work(ctx) │
└─────────┘                        absent  └─────────┘                     │ rows: ledger     │
                                                ▲ SetWorkSource / Stop     │ total: reltuples │
┌────────┐                                 ┌────┴───┐                      │ bytes: pg_table_ │
│ Copier │ Run: copies, reports nothing    │ Copier │ Run ─────────────────│   size(shadow,   │
└────────┘                                 └────────┘                      │   source) / poll │
                                                                           └──────────────────┘

References

🤖 Drafted with Amp (Claude Opus 4.6); reviewed and edited by the author.

@Kiran01bm
Kiran01bm force-pushed the kiran01bm/cs5c-copier-fillers branch from 3849622 to dd72fb6 Compare September 29, 2026 07:25
Base automatically changed from kiran01bm/cs5c-progress-copy to main September 30, 2026 22:00
The copier is the tracker's work source for the lifetime of Run: rows
from committed chunks, the source's catalog row count read once, and
both tables' pg_table_size measured at each poll. Measured, not estimated.
@Kiran01bm
Kiran01bm force-pushed the kiran01bm/cs5c-copier-fillers branch from dd72fb6 to aac9ebe Compare September 30, 2026 22:04
@Kiran01bm
Kiran01bm marked this pull request as ready for review September 30, 2026 22:08
@chatgpt-codex-connector

Copy link
Copy Markdown

You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard.

@aparajon

Copy link
Copy Markdown
Collaborator

🤖 Adversarial review (1/2): 1 blocking, 1 non-blocking.

I reviewed aac9ebe, which makes the Copier the tracker's WorkSource for the length of Run. The new tests and pkg/copier pass locally, and CI is green on PostgreSQL 14 through 18.

The handoff holds:

  • report registers the source before the resume clear, and the deferred stop drains any poll in flight before Run returns.
  • Work snapshots db, rowsTotal and the ledger's rows under mu and does its database read outside the lock. It never takes a lock that the engine holds while it calls StopWorkSource.
  • The table lock session is an advisory lock, so a poll waiting on pg_table_size cannot deadlock against the chunk transactions.
  • On a pool built by dbconn, a poll queued behind an application's ACCESS EXCLUSIVE on the source ends at the session's lock_timeout. I checked this with a probe: it returned 55P03 after 3s and Run finished cleanly.

I made 11 deliberate breaks to the code and 8 were caught. The caught ones:

  • not clamping reltuples = -1
  • not registering the source, or dropping the stop fence
  • keeping the pool after Run
  • swapping the two sizes
  • dropping COALESCE
  • pg_total_relation_size in place of pg_table_size
  • reporting zero rows

Of the three survivors, moving measureRowsTotal after the clear changes nothing. The other two are below.

Blocking

1. Work's catalog reads are not pg_catalog-qualified, so the progress poll calls whatever pg_table_size(oid) the session's search_path offers.

tableSizesSQL calls pg_table_size($1::oid). PostgreSQL resolves a function by choosing the best signature match among every schema on the path, and it uses path order only to break ties between identical signatures. The catalog's function is pg_catalog.pg_table_size(regclass), and an oid argument needs a coercion to reach it. A pg_table_size(oid) in any schema on the path is an exact match, so it wins even though pg_catalog is searched first.

The dbconn session hook does not help here. It fixes path order, and order is not what decides this lookup.

With PostgreSQL 16's default search_path ("$user", public) on the fixture's dbconn pool, I planted public.pg_table_size(oid):

  • One tracker.Progress() call ran it twice, once per table, inside pg-sprite's session and as the copier's role.
  • bytes_copied and bytes_total reported the planted value.

The progress figure is the smaller problem. The larger one is that any role that can create a function in a schema on the copier's path gets its code run by pg-sprite during every poll of a live copy. On PostgreSQL 14, that includes public, where every role can create by default.

The WorkSource contract already states the obligation: "A Work that reads the catalog is a CO-9 read site: every catalog relation, function and operator is pg_catalog-qualified, and its test runs under a search_path that shadows the names it uses." CO-9 names progress probes explicitly, and the other pkg/progress read sites follow it (progress.go:252-254).

The fix is the two constants:

const rowsTotalSQL = `SELECT COALESCE((SELECT reltuples FROM pg_catalog.pg_class WHERE oid OPERATOR(pg_catalog.=) $1::pg_catalog.oid), -1)`
const tableSizesSQL = `SELECT COALESCE(pg_catalog.pg_table_size($1::pg_catalog.regclass), 0), COALESCE(pg_catalog.pg_table_size($2::pg_catalog.regclass), 0)`

With both constants changed, the test below passes, and so do the PR's copier tests and TestCopierRefusesReplacedRelations.

Test case: fails at aac9ebe
package copier_test

import (
	"math"
	"sync"
	"testing"
	"time"

	"github.com/jackc/pgx/v5/pgxpool"
	"github.com/stretchr/testify/assert"
	"github.com/stretchr/testify/require"

	"github.com/block/pg-sprite/pkg/copier"
)

// A schema on the session's search_path that defines pg_table_size(oid)
// must not be called by the copier's progress poll: the exact-signature
// match outranks pg_catalog.pg_table_size(regclass) whatever the path order.
func TestCopierWorkIgnoresAShadowingSizeFunction(t *testing.T) {
	f := newCopierFixture(t)
	target, lock, shadow := f.prepare(t, 2000)
	f.exec(t, "ANALYZE %s.orders")
	f.exec(t, `CREATE TABLE %s.decoy_calls (n int);
		CREATE FUNCTION %s.pg_table_size(oid) RETURNS bigint LANGUAGE sql
		AS 'INSERT INTO %s.decoy_calls VALUES (1); SELECT 424242::bigint'`)
	cfg := f.pool.Config().Copy()
	cfg.ConnConfig.RuntimeParams["search_path"] = f.schema + ", pg_catalog"
	pool, err := pgxpool.NewWithConfig(t.Context(), cfg)
	require.NoError(t, err)
	defer pool.Close()

	pin := f.pinShadowKeyAtFirstChunk(t, shadow, 1050, 0)
	tracker := copyTracker(t)
	c, err := copier.NewCopier(target, shadow, lock, copier.Watermark{}, copier.Options{
		Workers: 4, Chunker: copier.ChunkerOptions{InitialRows: 100, MaxRows: 100},
		Clock: pin.clock, Tracker: tracker,
	})
	require.NoError(t, err)
	results := make(chan error, 1)
	var wg sync.WaitGroup
	wg.Go(func() { results <- c.Run(t.Context(), pool) })
	t.Cleanup(wg.Wait)
	pin.wait(t)
	require.Eventually(t, func() bool {
		pos := c.Position()
		return pos.Cut == copier.NewWatermark(math.MaxInt64) && len(pos.InFlight) == 1
	}, 15*time.Second, 20*time.Millisecond)

	snap, err := tracker.Progress(t.Context())
	require.NoError(t, err)
	var calls int
	require.NoError(t, f.pool.QueryRow(t.Context(), "SELECT count(*) FROM "+f.schema+".decoy_calls").Scan(&calls))
	assert.Zero(t, calls, "the progress poll ran a function from the session's search_path")
	assert.NotEqual(t, uint64(424242), snap.Detail.Work.BytesTotal)
	pin.release(t)
	require.NoError(t, <-results)
}
--- FAIL: TestCopierWorkIgnoresAShadowingSizeFunction (1.82s)
        Error:      Should be zero, but was 2
        Messages:   the progress poll ran a function from the session's search_path
        Error:      Should not be: 0x67932

relationOIDsSQL, the ST-6 guard on main, is written in the same unqualified style. I did not probe it, and it is outside this diff, but it deserves the same audit.

Non-blocking

2. Nothing pins that the copier is the work source during the resume clear. Work's doc says the copier is the source "for the lifetime of Run", but moving c.report(pool) after clearAbove survives every test. On a large resumed shadow, the clear is the long quiet stretch, and that is when an operator most wants the step to show the copy. A resumed-run test that polls while the clear's first batch is pinned would cover it.

Also untested: a failed reltuples read. It aborts Run before any row moves, which is the right fail-closed answer, but no test covers it.

This review was generated by Claude Code (claude-opus-5-5).

@aparajon

Copy link
Copy Markdown
Collaborator

🤖 Adoption and integration review (2/2): 0 blocking, 2 non-blocking.

OSS adoption

1. Work is bounded only if the pool passed to Run is. Run accepts any *pgxpool.Pool.

  • The chunk transactions don't rely on the pool: they set their own lock_timeout and statement_timeout from Options with SET LOCAL (copy_chunk.go:75-76).
  • Work sends tableSizesSQL straight through the pool, so its only bound is whatever session settings the pool carries.

I tested this with the fixture's pool configuration, minus dbconn's session bounds and connect hook:

  • A poll queued behind an application's ACCESS EXCLUSIVE on the source waited out the whole lock hold (4.3s) with no error.
  • StopWorkSource drains that poll before Run returns.

The WorkSource contract asks Work to bound itself. An outside importer who builds their own pool, which is the ordinary way to use a library taking *pgxpool.Pool, loses that without any sign that anything is wrong. There are two ways to close this:

  • Run the size read in a short transaction with the same SET LOCAL budgets as a chunk.
  • Or state on Run that the pool must come from dbconn.

The first keeps the copier self-contained, as the chunk path already is.

Integration ease for importers

2. On a resumed copy, rows_copied / rows_total never reaches 1. rows_copied counts only this run's rows, while rows_total is the whole source's reltuples. A copy resumed at 60% therefore finishes at about 40%. progress-report.md documents both halves separately. But an importer rendering a percent, as SchemaBot's progress comments do for other engines, will divide the two, and nothing tells them the quotient is not a completion fraction on a resumed run. One clause in the counter table would close it: "on a resumed run, the ratio is this run's share, not the copy's completion". Reporting the rows already landed below the watermark as a separate counter would close it too.

This review was generated by Claude Code (claude-opus-5-5).

@aparajon aparajon left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 Approving with comments: 1 blocking finding in the review above. Work's catalog reads need pg_catalog. qualification, per CO-9 and the WorkSource contract.

This review was generated by Claude Code (claude-opus-5-5).

Every catalog name the copier's measurements read is pg_catalog-qualified
(CO-9): the row-count read, the size read, and the chunk guard's relation
identity check. An unqualified pg_table_size(oid) resolved by signature
across the whole search_path, so a same-signature function in a user
schema won regardless of path order.

Each measurement now runs in a read-only transaction of its own under the
chunk budgets (LK-2), so a poll on a caller-built pool without session
timeouts still ends at the copier's lock timeout instead of waiting out a
queued lock. setBudgets is the one place the budgets are spelled; chunk
transactions and measurements share it.

Tests: a shadowing search_path with an impostor pg_class row and a
pg_table_size(oid) decoy leaves the counters real and the decoy uncalled;
a poll during the resume clear reports the copier's counters; a failed
row-count read ends Run before the clear touches the shadow; a poll
queued behind an ACCESS EXCLUSIVE request on a raw pgxpool ends with the
server's lock timeout (55P03), not the observer's patience.

Docs: the progress report states that on a resumed run rows_copied /
rows_total is that run's share, not completion, and that the size read is
bounded by the copy's own budgets.
@Kiran01bm

Copy link
Copy Markdown
Collaborator Author

🤖 Adversarial review response — created by Kiran's code review agent (Amp, Claude Opus 4.6) — pull/132, follow-up commit

All four findings fixed in the follow-up commit; the relationOIDsSQL audit the blocking finding asked for is done in the same commit.

# Finding Status Explanation
C1-F1 rowsTotalSQL / tableSizesSQL read pg_class and pg_table_size(oid) unqualified; the WorkSource contract and CO-9 demand pg_catalog qualification and a shadowing-search_path test; audit relationOIDsSQL too fixed Every catalog name the copier's measurements read is now pg_catalog-qualified: pg_catalog.pg_class, pg_catalog.pg_table_size($n::pg_catalog.oid), and OPERATOR(pg_catalog.=). The pg_table_size case is the one the review called out — the server resolves an unqualified call by signature across the whole search_path, so a pg_table_size(oid) in a user schema beats pg_catalog.pg_table_size(regclass) regardless of path order — and the doc comment now says so. The chunk guard's relationOIDsSQL (pg_class ⋈ pg_namespace) is qualified the same way, with a CO-9 comment; it had the same hole. New TestCopierWorkResistsCatalogShadowing runs Run on testutil.NewCatalogShadowingPool: an impostor pg_class row carrying reltuples = 424242 for the source OID, an empty pg_namespace, and a pg_table_size(oid) decoy that records every call into a table. The counters are the real ones, the decoy is never called, and the copy converges; unqualifying any one of pg_class, pg_table_size, or relationOIDsSQL fails it.
C1-F2 No test pins that the copier is the work source during the resume clear; no test for a failed reltuples read fixed TestCopierReportsItsWorkDuringTheResumeClear resumes from watermark 150 with a straggling uncommitted insert (key 200) holding the clear's SHARE MODE fence, polls the tracker while pg_locks shows the fence waiting, and gets Work{RowsCopied: 0, RowsTotal: 300} plus the two sizes — the copier is the source before the clear, not only before the first chunk. TestCopierRefusesToStartWithoutARowCount runs Run on a closed pgxpool: Run returns the connection error (not ErrInvariantViolation), Position is zero, the shadow is untouched, and the tracker was never handed a source.
C2-F1 Work is bounded only if the pool is; a raw *pgxpool.Pool has no session timeouts — run the size read in a short transaction with the chunk SET LOCAL budgets fixed Each measurement (Work's size read and the run-start row-count read) now runs in a read-only transaction of its own through Copier.measure, which sets the chunk budgets first. setBudgets is extracted from setCopySession, so chunk transactions and measurements spell the budgets in one place (LK-2). TestCopierWorkIsBoundedOnACallerBuiltPool runs the copier on a raw pgxpool.New pool with LockTimeout 500 ms, parks the worker before its first chunk transaction, holds ACCESS SHARE on the source from the test, queues an application ACCESS EXCLUSIVE behind it (lock_timeout = 0), and polls: the poll ends with the server's 55P03 inside the copier's timeout, not the observer's patience; the copy then completes and converges once the application gives up. Removing the budgets from the measurement transaction hangs the poll. The chunker's boundary query (Chunker.Next, unchanged since the chunker landed) is the copier's one remaining statement that runs on the caller's pool without SET LOCAL budgets; it is outside this PR's diff and tracked as an internal follow-up, in its own small PR.
C2-F2 On a resumed run rows_copied / rows_total never reaches 1 — say so in the docs/progress-report.md counter table fixed docs/progress-report.md now states, below the counter table, that on a resumed run the ratio is this run's share of the source, not the copy's completion — the earlier run's rows below the watermark are in the shadow but not in this run's count, so the ratio ends short of 1 — and that completion is the watermark reaching the top of the key space, which the engine reports as the step ending. It also records that the size read runs under the copy's own budgets. Copier.Work's doc comment and the copier row in docs/copy-and-swap-design.md carry the same clause. A separate whole-copy counter was not added: the copier does not know how many rows the earlier run landed without scanning the shadow, and a scan is not a measurement the poll should pay for.
Review Approving review no action Confirmation only; no finding to address.

Source: block/pg-sprite#132, review comments 5921080138 and 5921080298 and review 5372845561 at head aac9ebe; fixes in the follow-up commit.

@Kiran01bm
Kiran01bm enabled auto-merge (squash) September 30, 2026 23:58
@Kiran01bm
Kiran01bm merged commit b2a44a2 into main Oct 1, 2026
16 checks passed
@Kiran01bm
Kiran01bm deleted the kiran01bm/cs5c-copier-fillers branch October 1, 2026 00:02
Kiran01bm added a commit that referenced this pull request Oct 1, 2026
Resolves the pkg/checksum rows in SAFETY.md, docs/architecture.md, and
docs/copy-and-swap-design.md against the copier progress-filler rows
from #132: this branch's checksum text, main's copier and progress text.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants