From 329a0e9033690234f3fb7c27a560b17a12099051 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Tue, 29 Sep 2026 16:19:14 +1000 Subject: [PATCH 1/3] progress: copy operation and engine-measured work source (v4) 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. --- docs/progress-report.md | 72 +++++-- pkg/progress/progress.go | 77 ++++--- pkg/progress/progress_test.go | 6 +- pkg/progress/work_source.go | 50 +++++ pkg/progress/work_source_test.go | 355 +++++++++++++++++++++++++++++++ 5 files changed, 515 insertions(+), 45 deletions(-) create mode 100644 pkg/progress/work_source.go create mode 100644 pkg/progress/work_source_test.go diff --git a/docs/progress-report.md b/docs/progress-report.md index bc7ef2f..722f1b2 100644 --- a/docs/progress-report.md +++ b/docs/progress-report.md @@ -16,8 +16,10 @@ phase or operation value is a contract change and bumps `format_version`, even i is added or renamed. Adding a field bumps `format_version` so a strict consumer can detect the new shape from the -version. The current version is **3**: version 2 added `detail.statement`; version 3 added -`detail.current_locker_pid`, `work.lockers_total`, and `work.lockers_done`. +version. The current version is **4**: version 2 added `detail.statement`; version 3 added +`detail.current_locker_pid`, `work.lockers_total`, and `work.lockers_done`; version 4 added +the `copy` operation, whose `work` is measured by the engine rather than read from a server +progress view. The [plan report](plan-report.md), [lint report](lint-report.md), and [suggest report](suggest-report.md) are separate contracts with their own `format_version`; @@ -52,18 +54,27 @@ licenses a consumer to intervene in the change itself. | `active` | bool | always | Whether an operation is executing now. `false` with `phase: "running"` means a concurrent build's progress row has left the server view. | | `attempt` | int | bounded retries only | The current attempt number when the executor is inside its bounded retry loop. | | `current_locker_pid` | int | while waiting on a locker | PostgreSQL backend PID currently blocking the concurrent build; omitted when none is published. | -| `work` | object | server-observed work only | Present exactly when the server published a progress row; then **every** counter below is present, so a fresh build reports honest zeros rather than an empty object. | +| `work` | object | measured work only | Present exactly when something measured the step's work — the server published a progress row for a concurrent build, or the engine reported the copy step's counters; then **every** counter below is present, so a fresh build or an empty copy reports honest zeros rather than an empty object. | `statement` is the submitter's statement after qualification and canonicalization, so a consumer rendering it into a shared surface must clamp and escape it. ### Work counters -`blocks_done` / `blocks_total`, `tuples_done` / `tuples_total`, and -`lockers_done` / `lockers_total` come from -`pg_stat_progress_create_index` during a concurrent index build. `rows_copied` / -`rows_total` and `bytes_copied` / `bytes_total` are reserved for copy-and-swap and are `0` -on every native operation — the engine never fabricates copy counters. +Two operations publish `work`, and each measures only its own counters; the other +operation's counters are `0`, never estimated. + +**Server-observed** (`concurrent-index-build`): `blocks_done` / `blocks_total`, +`tuples_done` / `tuples_total`, and `lockers_done` / `lockers_total` come from +`pg_stat_progress_create_index`, read over the executor's reserved session. + +**Engine-measured** (`copy`): `rows_copied` / `rows_total` and `bytes_copied` / +`bytes_total` come from the copy-and-swap row copy itself, which knows the rows it has landed +from the chunks it committed. The copy step is not a single server statement, so PostgreSQL +publishes no progress view for it; the tracker instead polls the engine's work source for +the step. Every counter is something the engine measured — a row count from committed chunks, +a size read from the catalog — never a projection; a counter the engine cannot measure stays +`0`. Native operations report no rows or bytes. ## Phases @@ -85,15 +96,19 @@ returns the identical snapshot, elapsed values included. | `optimistic` | One bounded direct native attempt. | | `brief` | A brief transactional sequence step. | | `validate-constraint` | A constraint-validation scan. | -| `concurrent-index-build` | A concurrent index build (the one operation with server-observed `work`). | +| `concurrent-index-build` | A concurrent index build (`work` is server-observed). | +| `copy` | The copy-and-swap row copy from the source table into its shadow (`work` is engine-measured). | ## Polling semantics The tracker is caller-owned and has no goroutines or timers: polling lifetime is exactly the caller's context. A poll during an active concurrent index build performs one read of the -server's progress view over the executor's reserved session; every other poll is pure -memory. On a query error the returned snapshot still carries the last-known tracker state — -`phase` is never empty — with the error returned alongside for the caller to classify. +server's progress view over the executor's reserved session; a poll during a `copy` step +asks the engine's work source once; every other poll is pure memory. On a query or source +error the returned snapshot still carries the last-known tracker state — `phase` is never +empty — with the error returned alongside for the caller to classify. Pollers serialize +against each other, so the reserved session and the work source each see one observation +at a time. The tracker is also the operator's stop path for a running concurrent index build: `Tracker.CancelBuild` signals the build's backend over the same reserved session, and only @@ -122,7 +137,7 @@ A poll during step 2 of a 3-step sequence, mid concurrent index build: ```json { - "format_version": 3, + "format_version": 4, "phase": "running", "step": 2, "total_steps": 3, @@ -150,3 +165,34 @@ A poll during step 2 of a 3-step sequence, mid concurrent index build: } } ``` + +A poll during step 2 of a 4-step copy-and-swap, mid row copy (`TestSnapshotJSONShapeForACopyStep` +pins it): + +```json +{ + "format_version": 4, + "phase": "running", + "step": 2, + "total_steps": 4, + "elapsed_ns": 2750000000, + "step_elapsed_ns": 750000000, + "detail": { + "operation": "copy", + "statement": "INSERT INTO public.t_shadow (id) SELECT id FROM public.t WHERE id BETWEEN $1::bigint AND $2::bigint ON CONFLICT (id) DO NOTHING", + "active": true, + "work": { + "rows_copied": 1200, + "rows_total": 5000, + "bytes_copied": 98304, + "bytes_total": 409600, + "blocks_done": 0, + "blocks_total": 0, + "tuples_done": 0, + "tuples_total": 0, + "lockers_total": 0, + "lockers_done": 0 + } + } +} +``` diff --git a/pkg/progress/progress.go b/pkg/progress/progress.go index c3fc618..eb7caf4 100644 --- a/pkg/progress/progress.go +++ b/pkg/progress/progress.go @@ -46,7 +46,7 @@ const ( // field semantics. Adding a phase or operation value is a contract change // and bumps this version, even when no field is added or renamed. Adding a // field also bumps this version so strict consumers can detect the new shape. -const FormatVersion = 3 +const FormatVersion = 4 // Operation is the current operation's execution class. type Operation string @@ -63,13 +63,18 @@ const ( OperationValidate Operation = "validate-constraint" // OperationConcurrentIndex is a concurrent index build. OperationConcurrentIndex Operation = "concurrent-index-build" + // OperationCopy is the copy-and-swap row copy from the source table into + // its shadow. + OperationCopy Operation = "copy" ) -// Work reports server-observed work. It is present only when the server -// published a progress row, and then every counter marshals explicitly — a -// fresh build reports honest zeros, never an empty object a consumer must -// guess at. Rows and bytes are reserved for copy-and-swap; native operations -// do not fabricate them. +// Work reports observed work. It is present only when something measured +// it — the server published a progress row for a concurrent build, or the +// engine's own WorkSource reported the step's counters — and then every +// counter marshals explicitly — a fresh build reports honest zeros, never an +// empty object a consumer must guess at. Rows and bytes belong to the +// copy-and-swap row copy; blocks, tuples and lockers to a concurrent index +// build. Neither operation fabricates the other's counters. type Work struct { RowsCopied uint64 `json:"rows_copied"` RowsTotal uint64 `json:"rows_total"` @@ -112,27 +117,32 @@ type Snapshot struct { // Tracker is a concurrency-safe progress source. The caller owns it; it has // no goroutines. Progress performs the one read needed for an active index -// build, making polling lifetime identical to the caller's context. +// build, or the one WorkSource call for an engine-measured step, making +// polling lifetime identical to the caller's context. // // Two locks split the tracker's concerns: mu guards the state fields and is // held only for memory access, so the executor's own updates never wait for // a database read; pollMu serializes observers, so the reserved session — // a single pgx connection that is not safe for concurrent use — only ever -// carries one progress query at a time. +// carries one progress query at a time, and a WorkSource sees one poll at a +// time. // -// StopConcurrentBuild is the one state change that takes pollMu, because it -// is the fence between the build and the pool: the executor calls it before -// the build's session can be released, so an observation or cancel signal -// still in flight completes against a backend the build still owns. Start, -// StartStep and Finish also clear the build fields, but as resets under mu -// alone — by the time they run the build's step has already passed through -// StopConcurrentBuild, and a reset that waited behind an observation would -// make polling a gate on execution. +// StopConcurrentBuild and StopWorkSource are the two state changes that +// take pollMu, because each is the fence between the step's work and its +// release: the executor calls StopConcurrentBuild before the build's +// session can be released, so an observation or cancel signal still in +// flight completes against a backend the build still owns, and the engine +// calls StopWorkSource before the source's state goes away. 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, and a reset that waited behind an observation would make polling a +// gate on execution. type Tracker struct { mu sync.RWMutex pollMu sync.Mutex clock Clock session dbconn.RowQuerier + source WorkSource phase Phase started time.Time stepStart time.Time @@ -163,18 +173,18 @@ func (t *Tracker) Start(total int, operation Operation) { defer t.mu.Unlock() t.phase, t.started, t.stepStart, t.ended = PhaseRunning, now, now, time.Time{} t.step, t.total, t.detail = 0, total, Detail{Operation: operation, Active: true} - t.session, t.buildPID = nil, 0 + t.session, t.buildPID, t.source = nil, 0, nil } // StartStep advances a sequence to a 1-based step, records the exact SQL the -// executor will run, and drops any build session from a prior step, so a later -// step can never poll a stale build. +// executor will run, and drops any build session or work source from a prior +// step, so a later step can never poll a stale build or a finished source. func (t *Tracker) StartStep(step int, operation Operation, statement string) { t.mu.Lock() defer t.mu.Unlock() t.step, t.stepStart = step, t.clock.Now() t.detail = Detail{Operation: operation, Statement: statement, Active: true} - t.session, t.buildPID = nil, 0 + t.session, t.buildPID, t.source = nil, 0, nil } // SetAttempt records the current bounded retry attempt. @@ -186,11 +196,13 @@ func (t *Tracker) SetAttempt(attempt int) { // SetConcurrentBuild enables on-demand server progress for pid. The executor // supplies its reserved verdict session so polling cannot starve behind the -// build session even when the pool has only two connections. +// build session even when the pool has only two connections. A step's work +// comes from one place, so any WorkSource the step had is dropped. func (t *Tracker) SetConcurrentBuild(session dbconn.RowQuerier, pid uint32) { t.mu.Lock() defer t.mu.Unlock() t.session, t.buildPID = session, pid + t.source = nil } var ( @@ -374,15 +386,16 @@ func (t *Tracker) Finish(err error) { } t.ended = now t.detail.Active = false - t.session, t.buildPID = nil, 0 + t.session, t.buildPID, t.source = nil, 0, nil } // Progress returns a snapshot and, for an active concurrent index build, -// queries PostgreSQL's progress view by the executor-owned backend PID. On a -// query error the snapshot still carries the last-known tracker state. The -// state lock is released before the query, so concurrent pollers serialize -// only against each other (and StopConcurrentBuild), never against the -// executor's own state updates. +// queries PostgreSQL's progress view by the executor-owned backend PID; for +// a step with a WorkSource it asks the source for the step's counters. On a +// query or source error the snapshot still carries the last-known tracker +// state. The state lock is released before the query, so concurrent pollers +// serialize only against each other (and the two stop fences), never against +// the executor's own state updates. func (t *Tracker) Progress(ctx context.Context) (Snapshot, error) { t.pollMu.Lock() defer t.pollMu.Unlock() @@ -396,9 +409,15 @@ func (t *Tracker) Progress(ctx context.Context) (Snapshot, error) { s.Elapsed = now.Sub(t.started) s.StepElapsed = now.Sub(t.stepStart) } - session, pid := t.session, t.buildPID + session, pid, source := t.session, t.buildPID, t.source t.mu.RUnlock() - if pid == 0 || session == nil || s.Phase != PhaseRunning { + if s.Phase != PhaseRunning { + return s, nil + } + if source != nil { + return observeEngineWork(ctx, s, source) + } + if pid == 0 || session == nil { return s, nil } p, active, err := dbconn.ConcurrentIndexProgress(ctx, session, pid) diff --git a/pkg/progress/progress_test.go b/pkg/progress/progress_test.go index a19d054..c8de4e5 100644 --- a/pkg/progress/progress_test.go +++ b/pkg/progress/progress_test.go @@ -176,7 +176,7 @@ func TestStartResetsPriorRunState(t *testing.T) { // The JSON shape is the adapter-facing contract: exact keys, exact // omissions, driven through a real poll so the test pins what a consumer -// actually receives. A consumer pins format_version 3 against this test. +// actually receives. A consumer pins format_version 4 against this test. func TestSnapshotJSONShape(t *testing.T) { session := fakeSession{query: func(context.Context, string, ...any) pgx.Row { return fakeRow{scan: func(dest ...any) error { @@ -207,7 +207,7 @@ func TestSnapshotJSONShape(t *testing.T) { raw, err := json.Marshal(snapshot) require.NoError(t, err) assert.JSONEq(t, `{ - "format_version": 3, + "format_version": 4, "phase": "running", "step": 2, "total_steps": 3, @@ -266,7 +266,7 @@ func TestSnapshotJSONOmitsUnsetOptionalFields(t *testing.T) { raw, err := json.Marshal(snapshot) require.NoError(t, err) assert.JSONEq(t, `{ - "format_version": 3, + "format_version": 4, "phase": "pending", "elapsed_ns": 0, "step_elapsed_ns": 0, diff --git a/pkg/progress/work_source.go b/pkg/progress/work_source.go new file mode 100644 index 0000000..5163673 --- /dev/null +++ b/pkg/progress/work_source.go @@ -0,0 +1,50 @@ +package progress + +import "context" + +// WorkSource reports engine-derived work for the current step: the counters +// of an operation whose progress PostgreSQL publishes no view for, such as +// the copy-and-swap row copy, which knows its own rows copied from the +// chunks it has landed. The tracker polls it inside Progress, on the +// observer's context, so the source must be safe to call from any goroutine +// while the step runs. Every counter the source does not measure stays zero +// — a source reports what it observed, never an estimate dressed as a count. +type WorkSource interface { + Work(ctx context.Context) (Work, error) +} + +// SetWorkSource makes source the current step's work for every later poll. +// A step's work comes from one place: setting a source drops any concurrent +// build the step was polling, and SetConcurrentBuild drops the source. The +// engine calls StopWorkSource before the source's state goes away. +func (t *Tracker) SetWorkSource(source WorkSource) { + t.mu.Lock() + defer t.mu.Unlock() + t.source = source + t.session, t.buildPID = nil, 0 +} + +// StopWorkSource waits for an in-flight observation and then stops polling +// the source. It is the fence between the source and its owner's return: +// the engine calls it before the state the source reads is released, so +// no poll that began while the step ran completes against a source whose +// owner has gone. +func (t *Tracker) StopWorkSource() { + t.pollMu.Lock() + defer t.pollMu.Unlock() + t.mu.Lock() + defer t.mu.Unlock() + t.source = nil +} + +// observeEngineWork merges one source observation into s. On a source error +// the snapshot still carries the last-known tracker state, as the concurrent +// build path does on a query error. +func observeEngineWork(ctx context.Context, s Snapshot, source WorkSource) (Snapshot, error) { + work, err := source.Work(ctx) + if err != nil { + return s, err + } + s.Detail.Work = &work + return s, nil +} diff --git a/pkg/progress/work_source_test.go b/pkg/progress/work_source_test.go new file mode 100644 index 0000000..567ceb9 --- /dev/null +++ b/pkg/progress/work_source_test.go @@ -0,0 +1,355 @@ +package progress_test + +import ( + "context" + "encoding/json" + "errors" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/block/pg-sprite/pkg/progress" +) + +// fakeSource satisfies progress.WorkSource with a caller-supplied +// observation, standing in for the engine's copy step. +type fakeSource struct { + work func(ctx context.Context) (progress.Work, error) +} + +func (s fakeSource) Work(ctx context.Context) (progress.Work, error) { return s.work(ctx) } + +// copyCounters are distinct per field so a swapped pair cannot pass. +var copyCounters = progress.Work{RowsCopied: 1200, RowsTotal: 5000, BytesCopied: 98304, BytesTotal: 409600} + +// countingSource reports copyCounters and counts its polls. +func countingSource(polls *atomic.Int32) fakeSource { + return fakeSource{work: func(context.Context) (progress.Work, error) { + polls.Add(1) + return copyCounters, nil + }} +} + +// forbiddenSource fails the test if it is ever polled. +func forbiddenSource(t *testing.T, why string) fakeSource { + return fakeSource{work: func(context.Context) (progress.Work, error) { + t.Errorf("work source polled: %s", why) + return progress.Work{}, nil + }} +} + +// runningTrackerWithSource returns a tracker mid copy step, polling source. +func runningTrackerWithSource(t *testing.T, source progress.WorkSource) *progress.Tracker { + t.Helper() + tracker, err := progress.NewTracker(&fakeClock{now: time.Unix(100, 0)}) + require.NoError(t, err) + tracker.Start(1, progress.OperationCopy) + tracker.StartStep(1, progress.OperationCopy, "INSERT INTO public.t_shadow (id) SELECT id FROM public.t WHERE id BETWEEN $1::bigint AND $2::bigint ON CONFLICT (id) DO NOTHING") + tracker.SetWorkSource(source) + return tracker +} + +// The copy step's JSON is the adapter-facing contract for format_version 4: +// the copy operation value, and work carrying the engine's rows and bytes +// with the build counters at honest zero. +func TestSnapshotJSONShapeForACopyStep(t *testing.T) { + var polls atomic.Int32 + clock := &fakeClock{now: time.Unix(100, 0)} + tracker, err := progress.NewTracker(clock) + require.NoError(t, err) + tracker.Start(4, progress.OperationAdmitting) + clock.now = clock.now.Add(2 * time.Second) + const copySQL = "INSERT INTO public.t_shadow (id) SELECT id FROM public.t WHERE id BETWEEN $1::bigint AND $2::bigint ON CONFLICT (id) DO NOTHING" + tracker.StartStep(2, progress.OperationCopy, copySQL) + tracker.SetWorkSource(countingSource(&polls)) + clock.now = clock.now.Add(750 * time.Millisecond) + + snapshot, err := tracker.Progress(t.Context()) + require.NoError(t, err) + raw, err := json.Marshal(snapshot) + require.NoError(t, err) + assert.JSONEq(t, `{ + "format_version": 4, + "phase": "running", + "step": 2, + "total_steps": 4, + "elapsed_ns": 2750000000, + "step_elapsed_ns": 750000000, + "detail": { + "operation": "copy", + "statement": "INSERT INTO public.t_shadow (id) SELECT id FROM public.t WHERE id BETWEEN $1::bigint AND $2::bigint ON CONFLICT (id) DO NOTHING", + "active": true, + "work": { + "rows_copied": 1200, + "rows_total": 5000, + "bytes_copied": 98304, + "bytes_total": 409600, + "blocks_done": 0, + "blocks_total": 0, + "tuples_done": 0, + "tuples_total": 0, + "lockers_total": 0, + "lockers_done": 0 + } + } + }`, string(raw)) + assert.Equal(t, int32(1), polls.Load(), "one poll asks the source once") +} + +// Each poll asks the source afresh, so the counters a consumer sees are the +// source's current observation, not a value captured at registration. +func TestProgressPollsTheWorkSourceOnEveryObservation(t *testing.T) { + var polls atomic.Int32 + source := fakeSource{work: func(context.Context) (progress.Work, error) { + return progress.Work{RowsCopied: uint64(100 * polls.Add(1))}, nil + }} + tracker := runningTrackerWithSource(t, source) + + first, err := tracker.Progress(t.Context()) + require.NoError(t, err) + second, err := tracker.Progress(t.Context()) + require.NoError(t, err) + require.NotNil(t, first.Detail.Work) + require.NotNil(t, second.Detail.Work) + assert.Equal(t, uint64(100), first.Detail.Work.RowsCopied) + assert.Equal(t, uint64(200), second.Detail.Work.RowsCopied) +} + +// A source error leaves the snapshot with the last-known tracker state and +// no work, and the error reaches the poller — the same shape as a failed +// server progress query. +func TestProgressReturnsSnapshotAlongsideWorkSourceError(t *testing.T) { + sourceErr := errors.New("shadow size lookup failed") + tracker := runningTrackerWithSource(t, fakeSource{work: func(context.Context) (progress.Work, error) { + return progress.Work{}, sourceErr + }}) + + s, err := tracker.Progress(t.Context()) + require.ErrorIs(t, err, sourceErr) + assert.Equal(t, progress.PhaseRunning, s.Phase, "the snapshot must keep the last-known state on error") + assert.Equal(t, 1, s.Step) + assert.Equal(t, progress.OperationCopy, s.Detail.Operation) + assert.Nil(t, s.Detail.Work, "a failed observation reports no counters rather than zeros") +} + +// The source is polled only while the step runs: a terminal tracker reports +// its frozen snapshot without asking the source. +func TestFinishedTrackerDoesNotPollTheWorkSource(t *testing.T) { + tracker := runningTrackerWithSource(t, forbiddenSource(t, "the run has finished")) + tracker.Finish(nil) + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.PhaseFinished, s.Phase) + assert.Nil(t, s.Detail.Work) +} + +// A later step never polls the earlier step's source: the copy's source +// must not report copy counters against the checksum step. +func TestStartStepDropsThePriorStepsWorkSource(t *testing.T) { + tracker := runningTrackerWithSource(t, forbiddenSource(t, "the step it belonged to has ended")) + tracker.StartStep(2, progress.OperationBrief, "ALTER TABLE public.t_shadow VALIDATE CONSTRAINT c") + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Equal(t, 2, s.Step) + assert.Nil(t, s.Detail.Work) +} + +// A new run never polls the prior run's source, even when the prior run +// was abandoned without Finish or StopWorkSource: Start alone resets it. +func TestStartDropsThePriorRunsWorkSource(t *testing.T) { + tracker := runningTrackerWithSource(t, forbiddenSource(t, "the run it belonged to was abandoned")) + tracker.Start(1, progress.OperationCopy) + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.PhaseRunning, s.Phase) + assert.Nil(t, s.Detail.Work) +} + +// The source is polled only while the tracker is running: a source wired +// before Start is not consulted for a pending snapshot. +func TestPendingTrackerDoesNotPollTheWorkSource(t *testing.T) { + tracker, err := progress.NewTracker(&fakeClock{now: time.Unix(100, 0)}) + require.NoError(t, err) + tracker.SetWorkSource(forbiddenSource(t, "the tracker has not started")) + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.PhasePending, s.Phase) + assert.Nil(t, s.Detail.Work) +} + +// After StopWorkSource a poll reports the step without counters: the +// engine has released the state the source read. +func TestStopWorkSourceEndsPolling(t *testing.T) { + var polls atomic.Int32 + tracker := runningTrackerWithSource(t, countingSource(&polls)) + _, err := tracker.Progress(t.Context()) + require.NoError(t, err) + tracker.StopWorkSource() + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Equal(t, progress.PhaseRunning, s.Phase) + assert.Nil(t, s.Detail.Work) + assert.Equal(t, int32(1), polls.Load(), "the poll after StopWorkSource must not reach the source") +} + +// A step's work comes from one place: registering a concurrent build drops +// the source, and registering a source drops the build. +func TestWorkSourceAndConcurrentBuildAreMutuallyExclusive(t *testing.T) { + t.Run("build replaces source", func(t *testing.T) { + tracker := runningTrackerWithSource(t, forbiddenSource(t, "a concurrent build replaced it")) + tracker.SetConcurrentBuild(fakeSession{query: func(context.Context, string, ...any) pgx.Row { + return fakeRow{scan: func(...any) error { return pgx.ErrNoRows }} + }}, 4242) + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.False(t, s.Detail.Active, "the build path answered: its row is gone") + assert.Nil(t, s.Detail.Work) + }) + t.Run("source replaces build", func(t *testing.T) { + var polls atomic.Int32 + tracker := runningTrackerWithBuild(t, fakeSession{query: func(context.Context, string, ...any) pgx.Row { + return fakeRow{scan: func(...any) error { + t.Error("the reserved session was queried after a work source replaced the build") + return pgx.ErrNoRows + }} + }}) + tracker.SetWorkSource(countingSource(&polls)) + + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + require.NotNil(t, s.Detail.Work) + assert.Equal(t, copyCounters, *s.Detail.Work) + assert.Equal(t, int32(1), polls.Load()) + }) + t.Run("source ends the build's ownership of its backend", func(t *testing.T) { + tracker := runningTrackerWithBuild(t, neverQueried(t)) + tracker.SetWorkSource(countingSource(new(atomic.Int32))) + require.ErrorIs(t, tracker.CancelBuild(t.Context()), progress.ErrNoActiveBuild, + "a stale build PID must not be signalled once a source owns the step's work") + + tracker.StopWorkSource() + s, err := tracker.Progress(t.Context()) + require.NoError(t, err) + assert.Nil(t, s.Detail.Work, "stopping the source must not revive the replaced build") + }) +} + +// Two concurrent pollers never reach the source at the same time, so a +// source may read state that is not safe for concurrent observation. +func TestProgressSerializesConcurrentWorkSourcePolls(t *testing.T) { + firstEntered := make(chan struct{}) + overlap := make(chan struct{}) + release := make(chan struct{}) + var entries atomic.Int32 + tracker := runningTrackerWithSource(t, fakeSource{work: func(context.Context) (progress.Work, error) { + switch entries.Add(1) { + case 1: + close(firstEntered) + case 2: + close(overlap) + } + <-release + return copyCounters, nil + }}) + + var pollers sync.WaitGroup + for range 2 { + pollers.Go(func() { + _, err := tracker.Progress(t.Context()) + assert.NoError(t, err) + }) + } + <-firstEntered + secondPollerMustStillWait := time.After(100 * time.Millisecond) + select { + case <-overlap: + t.Fatal("two pollers reached the work source concurrently") + case <-secondPollerMustStillWait: + } + close(release) + pollers.Wait() + assert.Equal(t, int32(2), entries.Load(), "both pollers must complete, one after the other") +} + +// StopWorkSource must drain an in-flight observation before returning: the +// engine releases the state the source reads as soon as StopWorkSource +// returns. +func TestStopWorkSourceDrainsInFlightObservation(t *testing.T) { + entered := make(chan struct{}) + release := make(chan struct{}) + var observationFinished atomic.Bool + tracker := runningTrackerWithSource(t, fakeSource{work: func(context.Context) (progress.Work, error) { + close(entered) + <-release + observationFinished.Store(true) + return copyCounters, nil + }}) + + var workers sync.WaitGroup + workers.Go(func() { + _, err := tracker.Progress(t.Context()) + assert.NoError(t, err) + }) + <-entered + + stopReturned := make(chan struct{}) + workers.Go(func() { + tracker.StopWorkSource() + assert.True(t, observationFinished.Load(), + "StopWorkSource must not return while an observation still reads the source") + close(stopReturned) + }) + stopMustStillBlock := time.After(100 * time.Millisecond) + select { + case <-stopReturned: + t.Fatal("StopWorkSource returned while an observation was in flight") + case <-stopMustStillBlock: + } + close(release) + workers.Wait() +} + +// The engine's own state updates never wait behind a slow source poll: +// polling is observability, not a gate on the copy. +func TestStateMutatorsDoNotWaitForInFlightWorkSourcePoll(t *testing.T) { + entered := make(chan struct{}) + release := make(chan struct{}) + tracker := runningTrackerWithSource(t, fakeSource{work: func(context.Context) (progress.Work, error) { + close(entered) + <-release + return copyCounters, nil + }}) + + var workers sync.WaitGroup + defer workers.Wait() + defer close(release) + workers.Go(func() { + _, err := tracker.Progress(t.Context()) + assert.NoError(t, err) + }) + <-entered + + mutated := make(chan struct{}) + workers.Go(func() { + tracker.SetAttempt(2) + tracker.SetWorkSource(countingSource(new(atomic.Int32))) + close(mutated) + }) + mutatorDeadline := time.After(5 * time.Second) + select { + case <-mutated: + case <-mutatorDeadline: + t.Fatal("a state mutator waited behind an in-flight work source poll") + } +} From 1f502161d7481dc74a106fce36d6df6008c5f411 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Tue, 29 Sep 2026 17:20:37 +1000 Subject: [PATCH 2/3] docs: describe the copy step in the progress contract and SAFETY.md 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. --- SAFETY.md | 8 +++++--- docs/README.md | 2 +- docs/progress-report.md | 12 ++++++++---- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/SAFETY.md b/SAFETY.md index e55c738..64ec7bd 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -33,7 +33,7 @@ The invariant registry (invariant IDs referenced below) lives in | `pkg/diffplan` — desired schema → routed convergence plan, the declarative front door as a library (the CLI `diff` and embedding orchestrators share it) | ❌ periphery | exists | — | | `pkg/migrate` — one gated statement → resolve, classify, route, execute → one verdict; the imperative front door as a library (the CLI `migrate` and embedding orchestrators share it), plus the desired-state execution loop (`RunDesired`: derive the convergence plan, admit it as a whole, run each planned statement back through the same pipeline) | ❌ periphery² | exists | — | | `internal/cli` — CLI, flags, help, prompts | ❌ periphery | `migrate`, `pull`, `diff`, `fmt`, `lint`, `suggest`, `capabilities`, and `status` exist | — | -| `pkg/progress` — strategy-wide progress snapshots; the executors' observation seam (core imports it, so its locking discipline is core-critical); copy counters reserved for later | ✅ core | native progress exists | — | +| `pkg/progress` — strategy-wide progress snapshots; the executors' observation seam (core imports it, so its locking discipline is core-critical); the `WorkSource` seam for engine-measured steps such as the copy | ✅ core | native progress and the work-source seam exist | — | | orchestrator adapter | ❌ periphery | planned (Phase 11) | OC-* hold *at* the boundary | | `internal/testutil` | ❌ test-only | exists | — | @@ -93,9 +93,11 @@ The short version — the full rules live in [docs/tcb-model.md](docs/tcb-model. `pkg/progress` (the executors' progress-observation seam: they write state into a caller-owned tracker whose mutators take only a memory lock, and its polling reads ride the reserved verdict session behind a separate poll lock — the executor's own state - updates never wait for a database read, but the verdict handoff *is* observer-gated: + updates never wait for a database read, but the two handoffs *are* observer-gated: `StopConcurrentBuild` deliberately drains an in-flight poll before the executor reclaims - the session, a wait bounded by the poller's context and the session's `statement_timeout`), + the session, a wait bounded by the poller's context and the session's `statement_timeout`, + and `StopWorkSource` drains one before an engine step releases the state its `WorkSource` + reads, a wait bounded by the poller's context), stdlib. Adding one requires a recorded decision (see the rubric in [docs/tcb-model.md](docs/tcb-model.md) — copy small things, take pinned dependencies only for load-bearing expertise). diff --git a/docs/README.md b/docs/README.md index 4cddbc2..227c88b 100644 --- a/docs/README.md +++ b/docs/README.md @@ -49,7 +49,7 @@ Aurora-only. Why that combination is the product is [vision.md](vision.md); star | [limitations.md](limitations.md) | The **current limitations** — schema changes pg-sprite refuses today, why they are unsafe or unsupported, and where an operator must act outside the engine. | | [lint-report.md](lint-report.md) | The **lint report contract** — the versioned JSON shape `pg-sprite lint` emits for offline CI gating: finding fields (verbatim SQL, line/column), the codes table, severities and exit behavior, the offline-conservatism rules, and how the contract versions relative to the plan report. | | [suggest-report.md](suggest-report.md) | The **suggest report contract** — the versioned JSON shape `pg-sprite suggest` emits for offline advice: the typed caveat vocabulary (what changes about how you must run a safer form, and what a failed step leaves behind), the typed guidance codes for rewrites the planner cannot construct, and the operation → safer form → caveats table (pinned by test). | -| [progress-report.md](progress-report.md) | The **progress report contract** — the versioned JSON snapshot a caller receives when polling a running change through the `*WithProgress` entry points: phases and operations vocabularies, the terminal-freeze rule, server-observed work counters, and polling semantics (pinned by test). | +| [progress-report.md](progress-report.md) | The **progress report contract** — the versioned JSON snapshot a caller receives when polling a running change through the `*WithProgress` entry points: phases and operations vocabularies, the terminal-freeze rule, server-observed and engine-measured work counters, and polling semantics (pinned by test). | | [engine-role.md](engine-role.md) | The **engine-role provisioning contract** — the tiered minimum access a PostgreSQL user needs to run schema changes against tables it does not own: role membership for owner-gated DDL, schema `CREATE` for index builds and shadow objects, `SET ROLE` for owner-correct shadow creation, replication access for CDC, and the explicit list of powers the engine role must *not* have. Preflight refusals name the missing `GRANT` and point here. | | [invalid-index-recovery.md](invalid-index-recovery.md) | The **operator runbook** for the native path's invalid-index outcomes — an invalid index the executor found or left. How the typed states are told apart, which ones `RebuildAbandonedIndex` recovers on its own (and the lock-and-identity proof that makes its drop safe), when the entry may be another actor's healthy in-flight build, and what to check when the executor could prove nothing. | | [testing.md](testing.md) | The **test-suite guide** — how to run the suite (unit, per-major, all supported majors, compose database), current coverage, the remaining executor-phase test obligations, and the vanilla-PostgreSQL-matrix vs real-Aurora validation boundary. | diff --git a/docs/progress-report.md b/docs/progress-report.md index 722f1b2..b3ebe59 100644 --- a/docs/progress-report.md +++ b/docs/progress-report.md @@ -4,8 +4,9 @@ The progress snapshot is the machine-readable observation a caller receives when running schema change through the `*WithProgress` executor entry points. It is the one JSON shape an operator or orchestrator consumes to display or act on execution progress. This document is the contract: the fields, the closed vocabularies, and the behavior required of -a consumer. The Go source of truth is `pkg/progress`; `TestSnapshotJSONShape` pins the exact -keys, including the example at the end of this page. +a consumer. The Go source of truth is `pkg/progress`; `TestSnapshotJSONShape` and +`TestSnapshotJSONShapeForACopyStep` pin the exact keys, including the two examples at the end +of this page. ## Versioning: `format_version` @@ -56,8 +57,11 @@ licenses a consumer to intervene in the change itself. | `current_locker_pid` | int | while waiting on a locker | PostgreSQL backend PID currently blocking the concurrent build; omitted when none is published. | | `work` | object | measured work only | Present exactly when something measured the step's work — the server published a progress row for a concurrent build, or the engine reported the copy step's counters; then **every** counter below is present, so a fresh build or an empty copy reports honest zeros rather than an empty object. | -`statement` is the submitter's statement after qualification and canonicalization, so a -consumer rendering it into a shared surface must clamp and escape it. +`statement` is the SQL the engine is running for the step: for a native operation the +submitter's statement after qualification and canonicalization, for a `copy` step the +engine's own frozen chunk insert (the template with its `$1`/`$2` key bounds, not a chunk's +rendered values). Either way it is real SQL that reached the server, so a consumer rendering +it into a shared surface must clamp and escape it. ### Work counters From d75d38b0659ce8da2ee6c69a34c2e6408db168d3 Mon Sep 17 00:00:00 2001 From: Kiran Muddukrishna Date: Wed, 30 Sep 2026 18:30:03 +1000 Subject: [PATCH 3/3] progress: fence SetWorkSource/SetConcurrentBuild and bound the source 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. --- SAFETY.md | 12 +++++--- docs/progress-report.md | 23 +++++++++++--- pkg/progress/progress.go | 42 ++++++++++++++++---------- pkg/progress/work_source.go | 32 ++++++++++++++++++-- pkg/progress/work_source_test.go | 52 +++++++++++++++++++++++--------- 5 files changed, 119 insertions(+), 42 deletions(-) diff --git a/SAFETY.md b/SAFETY.md index 64ec7bd..9c8e7a3 100644 --- a/SAFETY.md +++ b/SAFETY.md @@ -93,11 +93,13 @@ The short version — the full rules live in [docs/tcb-model.md](docs/tcb-model. `pkg/progress` (the executors' progress-observation seam: they write state into a caller-owned tracker whose mutators take only a memory lock, and its polling reads ride the reserved verdict session behind a separate poll lock — the executor's own state - updates never wait for a database read, but the two handoffs *are* observer-gated: - `StopConcurrentBuild` deliberately drains an in-flight poll before the executor reclaims - the session, a wait bounded by the poller's context and the session's `statement_timeout`, - and `StopWorkSource` drains one before an engine step releases the state its `WorkSource` - reads, a wait bounded by the poller's context), + updates never wait for a database read, but the handoffs that end a poll target's + ownership *are* observer-gated: `StopConcurrentBuild` and `SetWorkSource` drain an + in-flight poll before the executor reclaims the build's session, a wait bounded by the + poller's context and the session's `statement_timeout`, and `StopWorkSource` and + `SetConcurrentBuild` drain one before an engine step releases the state its `WorkSource` + reads, a wait the `WorkSource` contract requires `Work` to bound itself — memory reads or + catalog reads under a session `statement_timeout`, never the observer's context alone), stdlib. Adding one requires a recorded decision (see the rubric in [docs/tcb-model.md](docs/tcb-model.md) — copy small things, take pinned dependencies only for load-bearing expertise). diff --git a/docs/progress-report.md b/docs/progress-report.md index b3ebe59..202f0f9 100644 --- a/docs/progress-report.md +++ b/docs/progress-report.md @@ -19,8 +19,9 @@ is added or renamed. Adding a field bumps `format_version` so a strict consumer can detect the new shape from the version. The current version is **4**: version 2 added `detail.statement`; version 3 added `detail.current_locker_pid`, `work.lockers_total`, and `work.lockers_done`; version 4 added -the `copy` operation, whose `work` is measured by the engine rather than read from a server -progress view. +the `copy` operation and split `work` into two counter families — the server-observed build +counters and the engine-measured copy counters (`rows_*`, `bytes_*`) — selected by +`detail.operation`, so `work` is no longer a signal that a concurrent index build is running. The [plan report](plan-report.md), [lint report](lint-report.md), and [suggest report](suggest-report.md) are separate contracts with their own `format_version`; @@ -55,7 +56,7 @@ licenses a consumer to intervene in the change itself. | `active` | bool | always | Whether an operation is executing now. `false` with `phase: "running"` means a concurrent build's progress row has left the server view. | | `attempt` | int | bounded retries only | The current attempt number when the executor is inside its bounded retry loop. | | `current_locker_pid` | int | while waiting on a locker | PostgreSQL backend PID currently blocking the concurrent build; omitted when none is published. | -| `work` | object | measured work only | Present exactly when something measured the step's work — the server published a progress row for a concurrent build, or the engine reported the copy step's counters; then **every** counter below is present, so a fresh build or an empty copy reports honest zeros rather than an empty object. | +| `work` | object | measured work only | Present exactly when something measured the step's work — the server published a progress row for a concurrent build, or the engine reported the copy step's counters; then **every** counter below is present, so a fresh build or an empty copy reports honest zeros rather than an empty object. Which counters mean anything is decided by `operation`, not by `work` being present — see [Work counters](#work-counters). | `statement` is the SQL the engine is running for the step: for a native operation the submitter's statement after qualification and canonicalization, for a `copy` step the @@ -66,7 +67,12 @@ it into a shared surface must clamp and escape it. ### Work counters Two operations publish `work`, and each measures only its own counters; the other -operation's counters are `0`, never estimated. +operation's counters are `0`, never estimated. A consumer selects the counter family from +`detail.operation`, never from the presence of `work`: a present `work` says only that +something measured the step, and a consumer that reads its presence as "a concurrent index +build is running" will show a `copy` step as a build with zero blocks. Render the build +counters for `concurrent-index-build`, the copy counters for `copy`, and nothing from `work` +for an operation you do not recognize. **Server-observed** (`concurrent-index-build`): `blocks_done` / `blocks_total`, `tuples_done` / `tuples_total`, and `lockers_done` / `lockers_total` come from @@ -114,6 +120,15 @@ empty — with the error returned alongside for the caller to classify. Pollers against each other, so the reserved session and the work source each see one observation at a time. +The handoffs that end a poll target's ownership of the step — stopping a build or a work +source, or replacing one with the other — wait for an observation in flight before they +return, so the engine never releases a session or the state a source reads while a poll is +still using it. That makes the work source part of the engine's stop path: its `Work` must +bound itself (memory reads, or catalog reads on a session with `statement_timeout` set) and +honour the poll's context, because a poll that returns only when its observer gives up +would stall the engine behind an observer that never does. A `Work` that reads the catalog +is a CO-9 read site — `pg_catalog`-qualified, tested under a shadowing `search_path`. + The tracker is also the operator's stop path for a running concurrent index build: `Tracker.CancelBuild` signals the build's backend over the same reserved session, and only while the build is active — the tracker never hands out the backend PID, so a caller cannot diff --git a/pkg/progress/progress.go b/pkg/progress/progress.go index eb7caf4..26a7bae 100644 --- a/pkg/progress/progress.go +++ b/pkg/progress/progress.go @@ -1,6 +1,8 @@ // Package progress defines the strategy-wide, machine-readable execution -// progress contract. It deliberately contains copy counters that native -// operations leave empty so copy-and-swap can implement the same contract. +// progress contract. One snapshot shape serves every operation: a concurrent +// index build's counters are read from the server's progress view, the +// copy-and-swap row copy's counters come from the engine's own WorkSource, +// and each operation leaves the other's counters at zero. package progress import ( @@ -127,16 +129,18 @@ type Snapshot struct { // carries one progress query at a time, and a WorkSource sees one poll at a // time. // -// StopConcurrentBuild and StopWorkSource are the two state changes that -// take pollMu, because each is the fence between the step's work and its -// release: the executor calls StopConcurrentBuild before the build's -// session can be released, so an observation or cancel signal still in -// flight completes against a backend the build still owns, and the engine -// calls StopWorkSource before the source's state goes away. 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, and a reset that waited behind an observation would make polling a -// gate on execution. +// Four state changes take pollMu, because each ends a poll target's +// ownership of the step's work and must not do so under a poll still in +// flight: StopConcurrentBuild and SetWorkSource release the build's +// session, which the executor calls before that session can return to the +// pool, so an observation or cancel signal in flight completes against a +// backend the build still owns; StopWorkSource and SetConcurrentBuild +// release the source, which the engine calls before the state the source +// reads goes away. Each takes pollMu before mu, the one order every taker +// uses. 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, and a reset that waited behind an observation would make +// polling a gate on execution. type Tracker struct { mu sync.RWMutex pollMu sync.Mutex @@ -197,8 +201,12 @@ func (t *Tracker) SetAttempt(attempt int) { // SetConcurrentBuild enables on-demand server progress for pid. The executor // supplies its reserved verdict session so polling cannot starve behind the // build session even when the pool has only two connections. A step's work -// comes from one place, so any WorkSource the step had is dropped. +// comes from one place, so any WorkSource the step had is dropped — after +// waiting, as StopWorkSource does, for a poll still reading it, so the +// source's owner never finds its state observed after the handoff. func (t *Tracker) SetConcurrentBuild(session dbconn.RowQuerier, pid uint32) { + t.pollMu.Lock() + defer t.pollMu.Unlock() t.mu.Lock() defer t.mu.Unlock() t.session, t.buildPID = session, pid @@ -252,7 +260,9 @@ const cancelBuildSQL = `SELECT state, // StopConcurrentBuild, which the executor calls before the build's session // can return to the pool, so the PID it signals still belongs to the build // — never to an unrelated statement that reused the same pooled backend. -// The signal is sent only to a backend the server reports active in the +// A step whose work a WorkSource owns has no build to signal and reports +// ErrNoActiveBuild; the snapshot's operation tells a caller which case it +// is in. The signal is sent only to a backend the server reports active in the // same statement as the read, which rules out the common way a cancel is // lost — a signal landing on an idle backend — without making the signal // itself observable. @@ -394,8 +404,8 @@ func (t *Tracker) Finish(err error) { // a step with a WorkSource it asks the source for the step's counters. On a // query or source error the snapshot still carries the last-known tracker // state. The state lock is released before the query, so concurrent pollers -// serialize only against each other (and the two stop fences), never against -// the executor's own state updates. +// serialize only against each other (and the four handoff fences), never +// against the executor's own state updates. func (t *Tracker) Progress(ctx context.Context) (Snapshot, error) { t.pollMu.Lock() defer t.pollMu.Unlock() diff --git a/pkg/progress/work_source.go b/pkg/progress/work_source.go index 5163673..18e4418 100644 --- a/pkg/progress/work_source.go +++ b/pkg/progress/work_source.go @@ -7,17 +7,43 @@ import "context" // the copy-and-swap row copy, which knows its own rows copied from the // chunks it has landed. The tracker polls it inside Progress, on the // observer's context, so the source must be safe to call from any goroutine -// while the step runs. Every counter the source does not measure stays zero +// while the step runs; the tracker itself never calls it from more than one +// goroutine at a time. Every counter the source does not measure stays zero // — a source reports what it observed, never an estimate dressed as a count. +// +// Work is on the engine's own stop path, so it has three obligations beyond +// returning counters: +// +// - Work bounds itself. SetWorkSource, SetConcurrentBuild and +// StopWorkSource wait for a poll in flight, and an observer may hand +// the tracker a context that never ends, so a Work that only returns +// when its caller gives up can stall the engine indefinitely. Work +// reads memory, or reads the catalog on a session whose +// statement_timeout is set, and honours ctx as well — a cancelled +// observer must not keep the poll alive. +// - Work takes no lock the engine holds while it calls SetWorkSource, +// SetConcurrentBuild or StopWorkSource. Those calls wait for Work to +// return; a Work that waits for the engine's lock would deadlock the +// step with its observer. +// - 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. type WorkSource interface { Work(ctx context.Context) (Work, error) } // SetWorkSource makes source the current step's work for every later poll. // A step's work comes from one place: setting a source drops any concurrent -// build the step was polling, and SetConcurrentBuild drops the source. The -// engine calls StopWorkSource before the source's state goes away. +// build the step was polling, and SetConcurrentBuild drops the source. +// +// Like StopWorkSource, it waits for an observation or cancel signal in +// flight, so the build it replaces is never released while a poll still +// reads its session, and a source it replaces has left its last poll +// before the engine lets the state behind it go. The engine calls +// StopWorkSource before the source's state goes away. func (t *Tracker) SetWorkSource(source WorkSource) { + t.pollMu.Lock() + defer t.pollMu.Unlock() t.mu.Lock() defer t.mu.Unlock() t.source = source diff --git a/pkg/progress/work_source_test.go b/pkg/progress/work_source_test.go index 567ceb9..155118c 100644 --- a/pkg/progress/work_source_test.go +++ b/pkg/progress/work_source_test.go @@ -282,10 +282,11 @@ func TestProgressSerializesConcurrentWorkSourcePolls(t *testing.T) { assert.Equal(t, int32(2), entries.Load(), "both pollers must complete, one after the other") } -// StopWorkSource must drain an in-flight observation before returning: the -// engine releases the state the source reads as soon as StopWorkSource -// returns. -func TestStopWorkSourceDrainsInFlightObservation(t *testing.T) { +// assertHandoffDrainsInFlightObservation starts a poll that blocks inside the +// source, then runs handoff on another goroutine and requires that it does +// not return until the observation has left the source. +func assertHandoffDrainsInFlightObservation(t *testing.T, handoff func(*progress.Tracker)) { + t.Helper() entered := make(chan struct{}) release := make(chan struct{}) var observationFinished atomic.Bool @@ -303,25 +304,46 @@ func TestStopWorkSourceDrainsInFlightObservation(t *testing.T) { }) <-entered - stopReturned := make(chan struct{}) + handoffReturned := make(chan struct{}) workers.Go(func() { - tracker.StopWorkSource() + handoff(tracker) assert.True(t, observationFinished.Load(), - "StopWorkSource must not return while an observation still reads the source") - close(stopReturned) + "the handoff must not return while an observation still reads the source") + close(handoffReturned) }) - stopMustStillBlock := time.After(100 * time.Millisecond) + handoffMustStillBlock := time.After(100 * time.Millisecond) select { - case <-stopReturned: - t.Fatal("StopWorkSource returned while an observation was in flight") - case <-stopMustStillBlock: + case <-handoffReturned: + t.Fatal("the handoff returned while an observation was in flight") + case <-handoffMustStillBlock: } close(release) workers.Wait() } +// Every call that ends the source's ownership of the step's work drains an +// in-flight observation before returning: the engine releases the state the +// source reads as soon as the call returns, whether it stopped the source or +// replaced it with another source or a concurrent build. +func TestSourceHandoffsDrainInFlightObservation(t *testing.T) { + t.Run("StopWorkSource", func(t *testing.T) { + assertHandoffDrainsInFlightObservation(t, (*progress.Tracker).StopWorkSource) + }) + t.Run("SetWorkSource", func(t *testing.T) { + assertHandoffDrainsInFlightObservation(t, func(tracker *progress.Tracker) { + tracker.SetWorkSource(countingSource(new(atomic.Int32))) + }) + }) + t.Run("SetConcurrentBuild", func(t *testing.T) { + assertHandoffDrainsInFlightObservation(t, func(tracker *progress.Tracker) { + tracker.SetConcurrentBuild(neverQueried(t), 4242) + }) + }) +} + // The engine's own state updates never wait behind a slow source poll: -// polling is observability, not a gate on the copy. +// polling is observability, not a gate on the copy. Only the handoffs that +// end a poll target's ownership wait; the resets that advance the run do not. func TestStateMutatorsDoNotWaitForInFlightWorkSourcePoll(t *testing.T) { entered := make(chan struct{}) release := make(chan struct{}) @@ -343,7 +365,9 @@ func TestStateMutatorsDoNotWaitForInFlightWorkSourcePoll(t *testing.T) { mutated := make(chan struct{}) workers.Go(func() { tracker.SetAttempt(2) - tracker.SetWorkSource(countingSource(new(atomic.Int32))) + tracker.StartStep(2, progress.OperationBrief, "ALTER TABLE public.t ADD COLUMN c integer") + tracker.Finish(nil) + tracker.Start(1, progress.OperationCopy) close(mutated) }) mutatorDeadline := time.After(5 * time.Second)