Repository navigation
Skip paths archive sibling messages but never terminalize sibling step_tasks rows, leaving status='started' tasks on completed runs #638
Description
Activity
Thanks for the report @cpursley !
I'm gonna review it soon!Reacted by Chase PursleyMaintainer verification and implementation plan
Status
The core report is confirmed on current
main.A direct reproduction used a root map step with three tasks,
maxAttempts: 1, andwhenExhausted: 'skip':- Tasks
0and1were started. - Task
2remained queued. - Task
0failed its only attempt. fail_taskskipped the map step and completed the run.
The committed state was:
run.status: completed run.remaining_steps: 0 step_states.status: skipped step_tasks[0].status: failed step_tasks[1].status: started step_tasks[2].status: queued active queue messages: 0 archived messages: 3This proves that message handling and orchestration completion work, while the persisted sibling task states remain non-terminal.
The same omission exists in
_cascade_force_skip_steps: itsskippedCTE updatesstep_states, and itsarchived_messagesCTE archives active task messages, but it does not update those task rows.Root cause
pgflow.fail_taskupdates only the task whose handler failed. When that task exhausts retries withwhen_exhausted IN ('skip', 'skip-cascade'), the function:- changes the parent
step_statesrow toskipped; - clears
remaining_tasks; - decrements
runs.remaining_steps; - archives queued and started sibling messages;
- emits the existing
step:skippedevent; - continues dependency or cascade handling;
- may complete the run.
It never changes queued or started sibling
step_tasksrows.pgflow._cascade_force_skip_stepshas the same state mismatch for any newly skipped step that already has task rows.Late callbacks do not repair the mismatch.
complete_taskandfail_taskguard against callbacks for a step that is no longerstarted, archive the callback message if needed, and return without changing the sibling task row.Acceptance invariant
After any transaction commits a step as skipped:
step_states.status = 'skipped' implies no step_tasks row for that step has status IN ('queued', 'started')More specifically:
- the task that exhausted its attempts remains
failed; - previously completed tasks remain
completed; - previously failed tasks remain
failed; - queued or started siblings become
skipped; - queued or started tasks for steps newly skipped by
_cascade_force_skip_stepsbecomeskipped; - all corresponding active PGMQ messages remain archived;
- run counters, step counters, outputs, skip reasons, event order, and return values remain unchanged;
- replayed and late callbacks do not revive or rewrite skipped task rows;
- existing affected rows are repaired during migration.
Chosen data model
Add
skippedas a validstep_tasks.statusvalue.Do not reuse another existing status:
failedis incorrect because sibling handlers did not fail and would inflate failure metrics;completedis incorrect because the handlers produced no accepted result;- deleting the rows would destroy task history and observability;
- leaving them queued or started preserves the bug.
Do not add a task-level
skipped_atorskip_reasoncolumn in this fix. The parentstep_statesrow already stores the authoritativeskipped_atandskip_reason, and every task references that row by(run_id, step_slug).Do not clear
started_at,last_worker_id,attempts_count, or requeue history. Those fields describe what happened before the orchestration cancelled the task and remain useful history.Do not add a partial index for skipped tasks. No current core query filters task rows by
status = 'skipped', and run/step lookups already use existing keys and indexes.Do not emit task-level Realtime events. The existing
step:skippedevent remains the public orchestration event.Files to change
Schema source of truth:
pkgs/core/schemas/0060_tables_runtime.sql pkgs/core/schemas/0100_function_fail_task.sql pkgs/core/schemas/0100_function__cascade_force_skip_steps.sqlPrimary regression tests:
pkgs/core/supabase/tests/fail_task_when_exhausted/skip_archives_sibling_messages.test.sql pkgs/core/supabase/tests/_cascade_force_skip_steps/archives_task_messages_for_skipped_steps.test.sql pkgs/core/supabase/tests/complete_task/late_complete_after_skip_does_not_mutate_step_or_run.test.sql pkgs/core/supabase/tests/fail_task_when_exhausted/late_fail_after_skip_does_not_double_decrement_remaining_steps.test.sql pkgs/core/supabase/tests/_cascade_force_skip_steps/does_not_archive_preexisting_skipped_step_messages.test.sql pkgs/core/supabase/tests/_cascade_force_skip_steps/idempotent_second_call.test.sqlDocumentation that lists runtime statuses:
NOMENCLATURE_GUIDE.md pkgs/website/src/content/docs/concepts/data-model.mdxGenerated artifacts:
pkgs/core/supabase/migrations/<timestamp>_pgflow_terminalize_skipped_tasks.sql pkgs/core/supabase/migrations/atlas.sum pkgs/core/src/database-types.ts and vendored/generated copies if gen-types changes them .changeset/<generated-name>.mddatabase-types.tscurrently represents the status column asstring, so a generated type diff is not expected. The generation and verification commands are still required.Test-first implementation sequence
1. Add failing task-status assertions
Extend
skip_archives_sibling_messages.test.sql, which already creates the exact required shape:- a three-task map step;
- task
0started and later failed; - task
1started; - task
2queued; - all three messages archived after the skip.
Add an ordered assertion equivalent to:
(task_index, status) (0, failed) (1, skipped) (2, skipped)Also assert that the skipped step has zero task rows with
status IN ('queued', 'started')and that the run reachescompletedwhere appropriate.The new assertion must initially fail because the schema rejects
skippedfor task rows and the functions do not write it.Extend
archives_task_messages_for_skipped_steps.test.sqlto assert that a directly cascaded active map step changes both its started and queued tasks toskippedwhile preserving message archival.Extend the late callback tests to assert the sibling task remains
skippedafter:- a late
complete_taskcallback; - a late
fail_taskcallback.
Extend the preexisting-skip and idempotency tests to assert:
_cascade_force_skip_stepschanges task rows only for step states newly changed by that invocation;- task rows under a step that was already skipped before the call remain untouched;
- a second call returns zero newly skipped steps and does not change task statuses or emit duplicate events.
Update each pgTAP
plan(N)count exactly.Run each changed test and confirm it fails for the missing task-state transition, not for malformed setup or SQL syntax.
2. Expand the task status constraint
In
0060_tables_runtime.sql, changestep_tasks.valid_statusfrom:status in ('queued', 'started', 'completed', 'failed')
to:
status in ('queued', 'started', 'completed', 'failed', 'skipped')
No other table constraint needs expansion:
- active task rows already have
output IS NULL, sooutput_valid_only_for_completedremains satisfied after they become skipped; - timestamp constraints do not require a terminal timestamp for each status;
completed_at_or_failed_atremains satisfied;- no foreign key depends on the status value.
3. Terminalize siblings in
fail_taskIn the existing branch guarded by:
IF v_task_exhausted AND v_step_skipped THENkeep the current message archive first, because its query filters task rows by
status IN ('queued', 'started').Immediately after that archive, update all still-active rows for the skipped step:
UPDATE pgflow.step_tasks AS task SET status = 'skipped' WHERE task.run_id = fail_task.run_id AND task.step_slug = fail_task.step_slug AND task.status IN ('queued', 'started');
Important boundaries:
- do not include
message_id IS NOT NULLin the update; an active task row must become terminal even if its message ID is absent; - do not update the exhausted task, because it is already
failedand the status predicate excludes it; - do not overwrite completed or previously failed siblings;
- do not change the existing step event, dependency propagation, cascade call, run completion, retry handling, or function return query;
- keep the archive and task update in the same PL/pgSQL function transaction so either both commit or both roll back.
4. Terminalize tasks in
_cascade_force_skip_stepsRefactor the current
archived_messagesportion of the CTE chain so message archival consumes the rows changed by a new task update CTE.The intended dependency is:
steps_to_skip -> skipped step_states -> skipped active step_tasks -> archived message IDsAdd a CTE with this behavior:
skipped_tasks AS ( UPDATE pgflow.step_tasks AS task SET status = 'skipped' WHERE task.run_id = _cascade_force_skip_steps.run_id AND task.step_slug IN ( SELECT skipped_step.step_slug FROM skipped AS skipped_step ) AND task.status IN ('queued', 'started') RETURNING task.message_id )
Then archive only non-null message IDs returned by
skipped_tasks:archived_messages AS ( SELECT pgmq.archive(v_flow_slug, array_agg(task.message_id)) AS result FROM skipped_tasks AS task WHERE task.message_id IS NOT NULL HAVING count(task.message_id) > 0 )
Retain the existing reference to
archived_messagesin the final query so PostgreSQL cannot optimize away the side-effecting archive select.Using the
skippedCTE'sRETURNINGrows is required. It preserves the helper's current idempotency and prevents a later cascade call from touching messages or tasks belonging to a step that was already terminal before this invocation.Created downstream steps normally have no task rows, so the task update is a no-op for them. It still correctly handles the helper's supported
startedstate and direct/internal invocations.5. Apply the schema locally and make focused tests pass
Follow the repository schema workflow:
cd pkgs/core pnpm nx supabase:start core pnpm nx supabase:status core pnpm nx fix-sql coreApply the changed schema files to the reported local database URL with
psql, then run the focused pgTAP files through./scripts/run-test-with-colors.Do not generate the migration until the schema files and focused tests are correct.
Concurrency and callback reasoning
fail_tasklocks the run and parent step state before it transitions the step.complete_taskand siblingfail_taskcallbacks lock the same run/step state and guard against a parent step that is no longer started.The task update must filter only
queuedandstartedrows:- if a sibling completes before the skip wins the step lock, it remains
completed; - if the skip commits first, the sibling becomes
skipped, and its late callback returns without rewriting it; - if
start_tasksor stalled-task recovery races on a task row, PostgreSQL row locking and the active-status predicate leave the committed row skipped after the skip transaction wins; - replayed skip calls do not rewrite already terminal tasks.
This fix does not attempt to stop handler code already executing in a worker. It makes the persisted orchestration state truthful and keeps late callbacks harmless. Cooperative worker cancellation is tracked separately in #646.
Migration and existing-data repair
Migrations must be generated from
pkgs/core/schemas/*.sql; do not write the migration from scratch.After schema behavior and tests pass:
cd pkgs/core pnpm nx fix-sql core ./scripts/atlas-migrate-diff terminalize_skipped_tasksReview the generated migration. It must:
- drop and recreate
step_tasks.valid_statuswithskippedallowed; - replace both changed functions;
- run a one-time data repair only after the expanded constraint is active.
Add the data repair to the generated migration using the repository's existing data-migration precedent:
UPDATE pgflow.step_tasks AS task SET status = 'skipped' FROM pgflow.step_states AS step WHERE step.run_id = task.run_id AND step.step_slug = task.step_slug AND step.status = 'skipped' AND task.status IN ('queued', 'started');
The backfill deliberately:
- repairs skipped steps whether the containing run is still started or already completed;
- leaves completed and failed task outcomes unchanged;
- leaves attempts, worker identity, timestamps, errors, and requeue history unchanged;
- does not attempt to archive messages again, because the affected historical skip paths already archived them;
- does not touch analogous rows on failed runs, which belong to Failed runs leave unfinished sibling tasks queued or started instead of cancelled #645.
After any manual data-migration addition to the generated Atlas migration, refresh the migration hash:
./scripts/atlas-migrate-hash --yes
Review migration ordering and run an upgrade rehearsal with fixture rows representing:
skipped parent + failed task skipped parent + completed task skipped parent + started task skipped parent + queued taskExpected post-migration statuses:
failed, completed, skipped, skippedDocumentation and release metadata
Update
NOMENCLATURE_GUIDE.mdso task statuses include:skipped - Task was cancelled because its parent step was skippedWhile touching the same section, include the already-supported
skippedstep status if it is still omitted.Update
pkgs/website/src/content/docs/concepts/data-model.mdxso the runtime table descriptions list current step and task statuses. Keep the wording about logical orchestration state; do not promise hard termination of running JavaScript handlers.Add a patch changeset for
@pgflow/core. The repository's fixed changeset group will handle synchronized package versions.Suggested changeset summary:
Terminalize queued and started task rows when their parent step is skipped.Required verification
Run focused tests first, followed by the repository checks:
cd pkgs/core ./scripts/run-test-with-colors supabase/tests/fail_task_when_exhausted/skip_archives_sibling_messages.test.sql ./scripts/run-test-with-colors supabase/tests/_cascade_force_skip_steps/archives_task_messages_for_skipped_steps.test.sql ./scripts/run-test-with-colors supabase/tests/complete_task/late_complete_after_skip_does_not_mutate_step_or_run.test.sql ./scripts/run-test-with-colors supabase/tests/fail_task_when_exhausted/late_fail_after_skip_does_not_double_decrement_remaining_steps.test.sql ./scripts/run-test-with-colors supabase/tests/_cascade_force_skip_steps/does_not_archive_preexisting_skipped_step_messages.test.sql ./scripts/run-test-with-colors supabase/tests/_cascade_force_skip_steps/idempotent_second_call.test.sql pnpm nx verify-migrations core pnpm nx gen-types core pnpm nx verify-gen-types core --skip-nx-cache pnpm nx test:pgtap coreAlso inspect the final diff and generated migration to ensure:
- only active sibling rows transition to skipped;
- archive behavior remains present in both functions;
- no event or counter logic changed accidentally;
- no generated vendored/type artifact is stale;
- the migration contains the historical backfill after the constraint expansion;
- no temporary migration remains.
Stop conditions
The implementation is complete when:
- every focused and full pgTAP test passes;
- migration and generated-type verification pass;
- a fresh database built from migrations matches the schema source;
- an upgrade fixture repairs old active rows under skipped steps;
- the exact three-task reproduction ends with
failed,skipped,skipped; - no queued or started task row remains under a skipped step;
- all three messages in that reproduction remain archived;
- late callbacks and replayed cascade calls remain idempotent.
If implementation reveals that adding
skippedrequires a public typed status union outside generated database types, stop and update the plan before broadening the API change.Explicitly out of scope
This issue must not absorb adjacent cancellation work:
- unfinished task rows on failed runs and the proposed
cancelledtask status: Failed runs leave unfinished sibling tasks queued or started instead of cancelled #645; - notifying workers to cooperatively abort already-running handlers: Use terminal events to cooperatively abort in-flight worker tasks #646;
- a user-facing cancel-run API;
- hard process termination;
- configurable worker control channels;
- changing stalled-task recovery beyond tests needed to prove skipped rows are no longer eligible.
- Tasks
- added 3 commits that reference this issue
on Aug 27, 2026 - added a commit that references this issue
on Sep 2, 2026 - added a parent issue
on Sep 2, 2026
Summary
When a map step is skipped via
when_exhausted: 'skip'/'skip-cascade'(and via_cascade_force_skip_stepsgenerally), the skip transaction updatesstep_states, decrementsruns.remaining_steps, and archives the sibling tasks' pgmq messages — but it never updates the siblingstep_tasksrows themselves. Those rows keepstatus = 'started'(withstarted_atset) forever, on a run that goes on to reachstatus = 'completed'.Observed in the SQL at commit
5c132f3(we vendor pgflow's core SQL into an Elixir port and verified this with integration tests there), and confirmed still present onmainas of8ffc889—schemas/0100_function_fail_task.sqlandschemas/0100_function__cascade_force_skip_steps.sqlboth archive sibling messages without touching the sibling task rows.Where in the code
All references are to the compiled core SQL (function bodies as installed):
fail_task, skip branch (when_exhaustedin('skip','skip-cascade')):maybe_fail_stepsets the step_state to'skipped'(skip_reason = 'handler_failed',skipped_at = now(),remaining_tasks = NULL).run_updatedecrementsruns.remaining_steps, so the run can still complete.The skip branch then archives sibling messages only:
There is no
UPDATE pgflow.step_tasks SET status = ...anywhere in this branch._cascade_force_skip_stepshas the same shape: itsskippedCTE updatesstep_statesonly, and itsarchived_messagesCTE archives messages only. Task rows are untouched.maybe_complete_runthen fires, so the run reaches'completed'with'started'task rows still attached to the skipped step.The message handling is correct and race-free (done in-transaction under
FOR UPDATE); it's specifically the task rows that are left behind.Reproduction sketch
initialTasks > 1) configuredwhenExhausted: 'skip'andmaxAttempts: 1.queued/started.fail_taskskips the step; the run completes.SELECT status, count(*) FROM pgflow.step_tasks WHERE run_id = $1 GROUP BY 1— the skipped step's sibling tasks are still'started'(or'queued'), on a'completed'run.Why it matters
The rows are permanently undispatchable by construction —
start_tasksrequirestask.status = 'queued'ANDrun.status = 'started'AND astep_statesrow at'started'— so no handler ever runs. The damage is to the data model, not execution:WHERE status = 'started'over-counts forever.status = 'started'churns. We hit this concretely: a stalled-task sweeper keyed onstarted+ age repeatedly "recovered" these rows on completed runs — flipping themstarted → queued, incrementing requeue counters until a permanent-stall marker fired, callingpgmq.archiveon already-archived message ids, and logging recovery work that never happened. The requeue is pure no-op churn (the messages were archived by the skip, andset_vton an archived id matches nothing), but the operational signal is permanently noisy and wrong for any deployment usingwhenExhausted: 'skip'on a map step.Suggested fix
In the same transaction that skips the step (both the
fail_taskskip branch and_cascade_force_skip_steps), terminalize the sibling rows alongside the message archive — e.g.:If adding a task-level
'skipped'status is undesirable, any terminal status with a marker would resolve the phantom-row problem equally well.