progress: copy operation and engine-measured work source (v4) - #131
Conversation
The row copy has no pg_stat_progress view, so the tracker learns its counters from the engine: a WorkSource polled inside Progress under the same fence as the concurrent build, mutually exclusive with it.
The statement field, the pin claim, the docs index, and the core table still described the report as server-observed only; the second pollMu fence (StopWorkSource) was missing from the dependency-list rationale.
|
You have reached your Codex usage limits for code reviews. You can see your limits in the Codex usage dashboard. |
|
🤖 Adversarial review (1/2): 0 blocking, 1 non-blocking. I reviewed The locking holds:
I made 12 deliberate breaks to the code and 9 were caught. The caught ones:
The three survivors behave the same in every case the tracker can reach:
The PG 15 Non-blocking1. A poll that never returns wedges the copy step, because
The fix belongs in the contract. Test case: fails at
|
|
🤖 Adoption and integration review (2/2): 0 blocking, 3 non-blocking. OSS adoption 1. The package synopsis on pkg.go.dev still describes copy counters as future work. Integration ease for importers 2. The presence of
SchemaBot does not pin 3. A source that reads the catalog is a CO-9 read site. The docs name "a size read from the catalog" as a source of This review was generated by Claude Code (claude-opus-5-5). |
aparajon
left a comment
There was a problem hiding this comment.
🤖 Stamped with 0 blocking findings. The non-blocking notes are in the two review comments above.
This stamp was left by Claude Code (claude-opus-5-5).
morgo
left a comment
There was a problem hiding this comment.
🤖 Automated adversarial review, posted on Morgan Tocker's behalf.
The shape is right. An operation the server publishes no progress row for needs a way to report what it measured, and modelling it as a source the tracker polls — rather than counters the engine pushes — keeps the tracker's "one read per observation, lifetime is the caller's context" property intact. Landing the interface and the fence before the consumer, with the format version bumped in the same commit as the new operation, is the right order.
Verified rather than assumed:
- The new precedence in
Progressis behaviour-preserving for the build path. The old guard waspid == 0 || session == nil || s.Phase != PhaseRunning; the new code hoists the phase test and then tests source before build. BecauseSetWorkSourceclearssession/buildPIDandSetConcurrentBuildclearssource, the two are mutually exclusive, so the reorder cannot shadow a build that would previously have been polled. StopWorkSourcetakes the locks in the same orderProgressdoes (pollMuthenmu), matchingStopConcurrentBuildandCancelBuild, so the new fence introduces no lock-order inversion among the tracker's own methods.observeEngineWorkcannot alias a snapshot.workis a fresh value per call andSnapshotis taken by value, so two observers never share a*Work.- A source is polled only while the tracker is running, because the phase test returns before the source is consulted, and
Start/StartStep/Finishall null it.
SetWorkSource is StopConcurrentBuild without the fence, and its doc comment recommends it as a substitute
func (t *Tracker) SetWorkSource(source WorkSource) {
t.mu.Lock()
defer t.mu.Unlock()
t.source = source
t.session, t.buildPID = nil, 0
}
func (t *Tracker) StopConcurrentBuild() {
t.pollMu.Lock()
defer t.pollMu.Unlock()
t.mu.Lock()
defer t.mu.Unlock()
t.session, t.buildPID = nil, 0
}The clearing is character-for-character identical. The only difference is pollMu, and StopConcurrentBuild's own comment says what that difference buys:
the executor calls it before the build's own session can return to the pool, so no signal that read the build's PID completes after that backend could be running someone else's statement.
SetWorkSource's comment, meanwhile, advertises the clearing as a feature — "setting a source drops any concurrent build the step was polling" — and the type comment on Tracker states the rule as "Start, StartStep and Finish also clear those fields, but as resets under mu alone — by the time they run the step has already passed through its fence." SetWorkSource clears the same fields under mu alone but is called during the step, before any fence has run. It is the one clearing path that does not satisfy the precondition the comment gives for the others.
Failure scenario (observation). A step runs a concurrent index build; the executor has handed the tracker its reserved session. An operator calls Progress(ctx). The observer takes pollMu, reads session, pid, releases mu, and is inside dbconn.ConcurrentIndexProgress on that reserved pgx connection. The executor's step moves to the copy and calls SetWorkSource(copier), which returns immediately. Reading its doc, the executor concludes the tracker no longer holds the build and returns the reserved session to the pool. The pool hands that connection to another caller while the observer's query is still in flight on it — concurrent use of a single pgx connection, which is precisely the hazard the pollMu comment names.
Failure scenario (cancel), the worse one. Same start, but the operator calls CancelBuild(ctx). It takes pollMu, reads pid, and is inside the detached pg_cancel_backend send. SetWorkSource runs concurrently — it does not take pollMu, so it is not serialized against CancelBuild the way StopConcurrentBuild is — and the executor releases the build's backend. The signal lands on a recycled backend running an unrelated statement. This is the exact outcome StopConcurrentBuild's comment says the fence exists to prevent, reachable through a method whose documentation reads as the newer way to drop a build.
The mirror holds for SetConcurrentBuild. The new t.source = nil line drops the source under mu alone, so an observation already inside source.Work(ctx) keeps running. StopWorkSource's comment states the contract the engine needs — "the engine calls it before the source's state goes away" — but SetConcurrentBuild now offers a way to drop the source that does not honour it, and an engine transitioning copy → concurrent build will reach for it for the same reason.
TestWorkSourceAndBuildMutuallyExclusive (both directions) currently pins the un-fenced clearing as the intended behaviour, so the suite would not catch this being wrong.
Fix, in order of preference:
- Give both
Setmethods the fence. Four lines, and the asymmetry disappears:The cost is that a step transition waits behind an in-flight observation — the same costfunc (t *Tracker) SetWorkSource(source WorkSource) { t.pollMu.Lock() defer t.pollMu.Unlock() t.mu.Lock() defer t.mu.Unlock() t.source, t.session, t.buildPID = source, nil, 0 }
StopConcurrentBuildalready pays, at the same once-per-step frequency. TheTrackercomment's objection ("a reset that waited behind an observation would make polling a gate on execution") was written aboutStart/StartStep/Finish, which run after the fence; it does not apply to a mid-step handoff that has had no fence at all. - If the un-fenced behaviour is deliberate, say so at both call sites: that the clearing is bookkeeping only, that the corresponding
Stop*is still required before the dropped state is released, and that an engine must never treatSetWorkSourceas a replacement forStopConcurrentBuild. Then add a test that the fence is still needed, so the requirement is pinned somewhere other than a comment.
Option 1 is what I would do — the contract currently has three methods that clear the build's fields and only one of them is safe to clear them with, which is a distinction the next caller will not preserve.
Work runs arbitrary engine code under pollMu, and StopWorkSource waits for it
Progress holds pollMu across source.Work(ctx), and StopWorkSource acquires pollMu. So the engine's teardown blocks until an observer's callback into the engine returns. The WorkSource doc explicitly anticipates a source that takes locks — "a source sees one poll at a time and may read state that is not safe for concurrent observation" — which makes the inversion concrete rather than theoretical:
- Engine goroutine holds the copier's internal mutex and calls
StopWorkSource; it blocks onpollMu. - Observer holds
pollMuand is insideWork, blocking on the copier's internal mutex.
Deadlock, and neither participant's context breaks it, because pollMu and a sync.Mutex inside the source are not context-aware. The observer's context bounds nothing here: cancelling it does not release pollMu.
The same shape gives a milder problem even without a cycle. StopConcurrentBuild waits behind a database round-trip, which the server and the connection bound; StopWorkSource waits behind whatever the source does, on whatever context the observer supplied — context.Background() from a long-poll handler is enough to make "waits for an in-flight observation" unbounded.
No consumer lands in this PR, so this is a contract gap rather than a live bug, and the next PR in the stack is exactly where it becomes one. Two sentences on WorkSource would close it: that Work must not acquire a lock the engine may hold when it calls StopWorkSource, and that it must return promptly and honour ctx, because both the tracker's other observers and the engine's own teardown wait on it.
Notes
The doc comment's "engine-measured" framing is the useful part of the version bump. Splitting the counter families explicitly — blocks/tuples/lockers as server-observed, rows/bytes as engine-measured, neither fabricating the other's — gives a consumer a rule it can apply to operations that do not exist yet, which is more durable than the operation-by-operation table. Worth making sure the format_version 4 entry in the doc's version history says the counter families were split, not only that an operation was added, since that is the part a strict consumer's validation would key on.
ErrNoActiveBuild after SetWorkSource is a good deliberate choice and is easy to misread as a bug. A caller that cancels a build and gets ErrNoActiveBuild has no way to tell "the build finished" from "the step moved to a copy". That is correct — there is nothing to cancel either way — but it is the kind of thing an operator-facing caller renders as "no build running" when the truthful answer is "the step is copying now". If the engine ever surfaces that error to a human, the snapshot's operation is the field that disambiguates it.
Replacing a poll target under mu alone let the engine release the build session or the source state while an observation was still reading it; both setters now drain in-flight polls like the Stop* fences do, and the WorkSource contract requires Work to bound itself and hold no engine lock.
|
🤖 Adversarial review response — created by Kiran's code review agent (Amp, Claude Opus 4.6) — pull/131, follow-up commit Verdict: all eight findings across the three adversarial comments are fixed in the follow-up commit — one code change (both
No action: the confirmed items — Source: block/pg-sprite#131, review comments 5886048966 and 5886050903 and review 5355697487 at head |
Adds the
copyoperation and aprogress.WorkSourcethe tracker polls for engine-measured work, bumping the progress report toformat_version4. No consumer of the source lands here; the next PR in the stack wirescopier.Copierto it.Why
The copy-and-swap row copy (previous PR in this stack) is many short transactions across N workers, not one server statement, so PostgreSQL publishes no
pg_stat_progress_*row for it. The tracker's only work path today is the concurrent-build one — a reserved session and a backend PID queried insideProgress. The copier already knows its own counters (Position.RowsInserted, the shadow's size in the catalog); what was missing is a way to hand them to the tracker under the same polling and lifetime rules the build path has, so an operator'sProgress()during the copy reports rows and bytes instead of an empty step.What
pkg/progress:OperationCopy = "copy".WorkSourceinterface —Work(ctx) (Work, error)— andTracker.SetWorkSource/Tracker.StopWorkSource.Progresspolls the source inside the existingpollMucritical section, on the observer's context, so a source sees one poll at a time and may read state that is not safe for concurrent observation.StopWorkSourceis a fence likeStopConcurrentBuild: it waits for an in-flight observation before clearing the source, so the engine can release the source's state right after it returns. The two setters are fences too:SetWorkSourceandSetConcurrentBuildeach replace a poll target, so each takespollMuthenmu(the one order every taker uses) and drains an in-flight poll before the old target's owner is free to release it.WorkSourcecontract putsWorkon the engine's stop path: it must bound itself (memory reads, or catalog reads under a sessionstatement_timeout) and honour ctx, since the fences wait for it and an observer's context may never end; it must take no lock the engine holds while callingSet*/StopWorkSource; and a catalog-readingWorkis a CO-9 read site (pg_catalog.qualification, shadowing-search_pathtest).SetWorkSourcedrops the build session and PID (soCancelBuildrefuses withErrNoActiveBuildand stopping the source does not revive the build), andSetConcurrentBuilddrops the source.Start,StartStepandFinishreset the source like they reset the build fields; a source is polled only while the tracker is running.FormatVersion3 → 4. TheWorkdoc now distinguishes server-observed counters (blocks, tuples, lockers) from engine-measured ones (rows, bytes); neither operation fabricates the other's.docs/progress-report.md: version history (v4 splitsworkinto two counter families selected bydetail.operation),workpresence rule ("measured work only" — and consumers pick the family fromoperation, never fromworkbeing present), the two counter families, thecopyrow in the operations table, polling semantics including the four handoff fences and theWorkobligations, and a pinned copy-step example.SAFETY.md'spkg/progressdependency rationale names the same fences.Consumers that reject unknown
format_versionvalues need a pin bump to 4; nothing else in the shape changed, and the existing build example is unchanged apart from the version.Tests
pkg/progress/work_source_test.go(pure,-race): the copy-step JSON shape pinned as a literal (matching the doc example; rows and bytes carry distinct values so a swapped pair fails); the source polled once per observation; a source error returns the running snapshot withWorknil; no poll afterFinish, afterStartStep, afterStarton an abandoned run, on a pending tracker, or afterStopWorkSource; build and source mutually exclusive in both directions, includingCancelBuild→ErrNoActiveBuildonce a source owns the step; two concurrent pollers never overlap inside the source;StopWorkSource,SetWorkSourceandSetConcurrentBuildeach drain an in-flight observation (one table-driven test; removingpollMufrom either setter fails exactly its subtest);SetAttempt/StartStep/Finish/Startdo not wait behind an in-flight poll. Every test was checked against a mutant that removes the branch it covers.TestDocStatesCurrentFormatVersionpins the doc's stated version toFormatVersion.Before / after
References
copier.CopieraWorkSourceand wires it throughcopier.Options.docs/progress-report.md;docs/copy-and-swap-design.mdpackage map (updated in the next PR when the copier fills the counters).🤖 Drafted with Amp (Claude Opus 4.6); reviewed and edited by the author.