Skip to content

Failed runs leave unfinished sibling tasks queued or started instead of cancelled #645

Description

@jumski

Summary

When pgflow fails a run, it archives queued and started PGMQ messages but does not terminalize every corresponding step_tasks row. The task that directly caused the failure becomes failed; unfinished sibling and parallel tasks can remain queued or started on a permanently failed run.

This is the failed-run counterpart to #638, which covers active task rows stranded under skipped steps.

Confirmed reproduction

A root map step with three tasks was run with one allowed attempt:

  1. Tasks 0 and 1 were started.
  2. Task 2 remained queued.
  3. Task 0 exhausted its attempts with when_exhausted = 'fail'.

The resulting state on current main was:

run.status:             failed
step_states.status:     failed

step_tasks[0].status:   failed
step_tasks[1].status:   started
step_tasks[2].status:   queued

active queue messages:  0
archived messages:      3

The active messages were correctly archived, so the sibling rows were no longer dispatchable. Their persisted statuses did not describe that terminal state.

After backdating task 1, the current built-in requeue_stalled_tasks() selected it despite the failed run, changed it from started to queued, and incremented requeued_count. Its archived message could not become visible again, so the recovery attempt performed no executable work.

Known affected paths

1. pgflow.fail_task with when_exhausted = 'fail'

fail_or_retry_task marks only the exhausted task as failed. The function then marks the step and run failed and archives every queued or started message in the run. It does not change the other active task rows.

File:

pkgs/core/schemas/0100_function_fail_task.sql

2. pgflow.complete_task type violation

When a single step returns a non-array value required by a map step, complete_task fails the current task, step, and run and archives all active messages. Other queued or started task rows remain active in the database.

File:

pkgs/core/schemas/0100_function_complete_task.sql

3. pgflow.cascade_resolve_conditions with when_unmet = 'fail'

A failed condition can terminalize the run while independent branches already have queued or started tasks. The function archives all active messages but does not terminalize those task rows.

File:

pkgs/core/schemas/0100_function_cascade_resolve_conditions.sql

4. Late callback guards do not provide eventual cleanup

complete_task returns without task mutation when the run is already failed. fail_task can mark a particular late failing task as failed, but tasks whose workers never call back remain queued or started indefinitely.

Why this matters

The queue and run state say the work is terminal, while step_tasks says it is active.

Consequences include:

  • failed runs appearing to contain live work;
  • stalled-task recovery attempting to recover undispatchable tasks;
  • misleading queued, started, recovered, and permanently stalled metrics;
  • dashboards and health checks reporting phantom work;
  • future cleanup or recovery code making decisions from an invalid status invariant.

This normally does not redispatch the handler because start_tasks requires both runs.status = 'started' and step_states.status = 'started', and the messages were archived. The defect is persisted state and operational behavior.

Desired invariant

After a run commits as failed:

runs.status = 'failed'
    implies
no step_tasks row for that run has status IN ('queued', 'started')

Task outcomes must preserve their meaning:

  • the task that actually failed remains failed;
  • tasks that completed before the run failure remain completed;
  • unfinished queued or started tasks become terminal without being mislabeled as failures.

Proposed data model

Add a task-level cancelled status for unfinished work invalidated by a terminal run.

Preferred transition:

queued  -> cancelled
started -> cancelled

Do not mark these rows failed: their handlers did not necessarily fail, and doing so would inflate task-failure metrics.

Do not mark them skipped: #638 uses skipped when the parent step itself has the explicit skipped outcome. A task cancelled because another step failed belongs to a different terminal reason.

Do not delete the rows: they retain task history, attempts, worker identity, and timing information.

Chosen schema semantics

Keep the task representation minimal:

  • add status = 'cancelled';
  • do not add cancelled_at; use the parent run's failed_at as the cancellation time;
  • do not add a task-level cancellation reason; derive run_failed from the parent run;
  • preserve existing attempts, worker identity, timestamps, errors, and requeue history;
  • let orchestration cancellation win over every late callback.

A future cancellation source may add its own proven data requirements. This fix does not reserve columns for it.

Proposed implementation surfaces

Schema and functions:

pkgs/core/schemas/0060_tables_runtime.sql
pkgs/core/schemas/0100_function_fail_task.sql
pkgs/core/schemas/0100_function_complete_task.sql
pkgs/core/schemas/0100_function_cascade_resolve_conditions.sql
pkgs/core/schemas/0062_function_requeue_stalled_tasks.sql

Likely behavior:

  1. Expand step_tasks.valid_status with cancelled.
  2. In each run-failure transaction, preserve completed and failed rows and set every queued or started row in the failed run to cancelled.
  3. Archive messages in the same transaction, using ordering or UPDATE ... RETURNING so status changes do not hide message IDs from the archive query.
  4. Backfill queued or started task rows attached to existing failed runs.
  5. Harden requeue_stalled_tasks so recovery eligibility mirrors dispatch eligibility.

Stalled-task recovery hardening

requeue_stalled_tasks currently selects task rows from status and age alone. It joins the run but does not require a started run, and it does not join the parent step state.

The recovery predicate should include:

run.status = 'started'
AND step_state.status = 'started'

This mirrors start_tasks: recovering a task is useful only if pgflow could dispatch it again.

This guard is valuable even after the cancellation backfill. It protects recovery from future terminal paths that accidentally miss task cleanup.

The implementation must preserve genuine stalled-task behavior:

  • started task on a started step and started run remains recoverable;
  • failed, completed, skipped, or otherwise terminal runs and steps are ignored;
  • the existing max-requeue and permanent-stall behavior remains unchanged.

Migration and backfill

Generate the migration from schema source through Atlas. Do not write it from scratch.

After the new status is allowed, repair historical rows with the equivalent of:

UPDATE pgflow.step_tasks AS task
SET status = 'cancelled'
FROM pgflow.runs AS run
WHERE run.run_id = task.run_id
  AND run.status = 'failed'
  AND task.status IN ('queued', 'started');

The migration must preserve completed and failed outcomes and retain task history columns. It adds no cancellation timestamp or reason backfill; runs.failed_at is authoritative.

Required tests

Exhausted task failure

Create a multi-task map step with:

task 0: started, then failed
 task 1: started
 task 2: queued

Expected final statuses:

failed, cancelled, cancelled

Assert the run and step are failed, no active messages remain, and no task remains queued or started.

Type violation

Keep an independent map or single branch active while another step triggers a single-to-map type violation. Assert the directly invalid task remains failed and unrelated unfinished tasks become cancelled.

Condition failure

Keep an independent branch active while a newly ready conditional step fails with when_unmet = 'fail'. Assert active task rows across the failed run become cancelled.

Late callbacks

After cancellation:

  • late complete_task must not change cancelled to completed;
  • late fail_task must not change cancelled to failed;
  • run and step counters and events must remain unchanged;
  • repeated callbacks remain idempotent.

Orchestration cancellation wins. A late physical result is ignored just as current terminal-state guards ignore it.

Recovery guard

Prove that:

  • a genuinely stalled task on a started run and started step is requeued;
  • a stale started row on a failed run is ignored;
  • a stale started row under a terminal step is ignored;
  • existing requeue counts and permanent-stall behavior still work.

Upgrade fixture

Before applying the migration, create failed-run fixtures with completed, failed, started, and queued task rows. After migration, expect:

completed, failed, cancelled, cancelled

Concurrency requirements

Run/step failure and task callbacks can race. The solution must preserve these outcomes:

  • a task committed as completed before run failure remains completed;
  • a task still queued or started when run failure commits becomes cancelled;
  • a late callback cannot revive a cancelled row;
  • message archival and task terminalization commit atomically;
  • replayed failure paths do not emit duplicate events or rewrite terminal task outcomes.

Review lock ordering across fail_task, complete_task, condition resolution, and stalled-task recovery before finalizing statement order.

Public API and documentation

Adding cancelled is a public persisted status value. Update every status list in schema documentation and generated/public types if any typed union exists.

Generated Supabase table types currently use string, but generation and verification remain required.

Document that database cancellation describes orchestration state. It does not guarantee that JavaScript already executing in a worker stopped before producing external side effects.

Worker-side cooperative abort is explored separately in #646.

Acceptance criteria

  • cancelled has a documented, unambiguous task-level meaning.
  • Every known failed-run path terminalizes all queued and started task rows.
  • Completed and genuinely failed task outcomes remain unchanged.
  • All active PGMQ messages remain archived on run failure.
  • No failed run retains queued or started task rows after migration.
  • Stalled-task recovery only considers dispatchable runs and steps.
  • Late callbacks cannot revive cancelled tasks.
  • Existing failure events, counters, and function return contracts remain stable.
  • Historical rows are backfilled safely.
  • Focused pgTAP, full pgTAP, migration, and generated-type checks pass.

Out of scope

Related and sequencing

Activity

  1. added
    priority:p1Next batch: current correctness, user blocker, or operational safety
    on Aug 28, 2026
  2. jumski commented on Aug 28, 2026

    @jumski
    ContributorAuthor

    Follow-up: queue archival alone is not a sufficient invariant

    I considered a smaller version of this fix: keep sibling task rows as queued / started, rely on archiving their PGMQ messages, and only add active run/step guards to stalled recovery. After tracing the current two-phase poller and callback locking, I do not think that is a safe long-term boundary.

    The worker path is:

    pgmq.read_with_poll() commits
    → worker holds message + msg_id
    → pgflow.start_tasks()
    → handler executes
    → complete_task() / fail_task()
    

    Archiving prevents future reads, but it cannot revoke a message that a worker already holds. start_tasks() does not establish eligibility from queue presence; its candidate CTE checks run/step status under the statement snapshot and its guarded update rechecks only step_tasks.status = 'queued'.

    One relevant interleaving on current main is:

    1. A worker reads sibling message M.
    2. start_tasks() takes a snapshot while the run still appears started.
    3. Another transaction fails the run and archives M, but leaves the sibling task row queued.
    4. The guarded task update can still change that row to started and return it to the worker because the failure transaction never changed or locked that row.
    5. Updating visibility for the now-archived message affects zero rows; current start_tasks() can still return the task.

    #656 should make visibility extension structurally required and reject a missing message, which closes that particular post-archive claim window. It does not remove the broader problem of making every mutation path infer task liveness from parent state.

    There is also a callback race in complete_task(): the failed-run guard runs before the FOR UPDATE lock. If another task fails the run between that guard and lock acquisition, the code reads the updated run row but does not recheck its status under lock. A sibling row left as started can still enter the completion CTE. Terminalizing the task row makes the existing status = 'started' update guard no-op after the failure wins.

    Smallest robust boundary

    On every run-failure path:

    1. update all remaining queued / started sibling task rows to a terminal status;
    2. capture their message IDs with RETURNING;
    3. archive those messages afterward in the same transaction;
    4. preserve the task-before-queue lock order established by Skip paths archive sibling messages but never terminalize sibling step_tasks rows, leaving status='started' tasks on completed runs #638;
    5. also require run.status = 'started' and step_state.status = 'started' in stalled recovery as defense in depth.

    A terminal task-row update provides the row lock and current-state recheck that start_tasks(), complete_task(), and fail_task() already use. If the claim wins first, the handler may already execute and cancellation cannot undo its external side effects. If failure wins first, the later guarded claim/callback cannot revive the task.

    cancelled remains the accurate terminal value: the culprit stays failed, completed work stays completed, and unfinished work invalidated by the failed run becomes cancelled. Reusing failed corrupts failure metrics; reusing skipped conflates run failure with explicit step-skip policy. No cancelled_at or reason column is needed initially; runs.failed_at supplies the timestamp and reason context.

    Conclusion: the recovery guards are necessary, but archive-only plus guards should not replace task terminalization. This issue should remain the failed-run counterpart to #638 and retain its concurrency tests and historical backfill.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething isn't workingpkgs/corepriority:p1Next batch: current correctness, user blocker, or operational safety

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions