Skip to content

Skip paths archive sibling messages but never terminalize sibling step_tasks rows, leaving status='started' tasks on completed runs #638

Description

@cpursley

Summary

When a map step is skipped via when_exhausted: 'skip' / 'skip-cascade' (and via _cascade_force_skip_steps generally), the skip transaction updates step_states, decrements runs.remaining_steps, and archives the sibling tasks' pgmq messages — but it never updates the sibling step_tasks rows themselves. Those rows keep status = 'started' (with started_at set) forever, on a run that goes on to reach status = '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 on main as of 8ffc889 — schemas/0100_function_fail_task.sql and schemas/0100_function__cascade_force_skip_steps.sql both 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_exhausted in ('skip','skip-cascade')):
    • maybe_fail_step sets the step_state to 'skipped' (skip_reason = 'handler_failed', skipped_at = now(), remaining_tasks = NULL).

    • run_update decrements runs.remaining_steps, so the run can still complete.

    • The skip branch then archives sibling messages only:

      PERFORM pgmq.archive(r.flow_slug, ARRAY_AGG(st.message_id))
      FROM pgflow.step_tasks st ... WHERE st.status IN ('queued','started') ...

      There is no UPDATE pgflow.step_tasks SET status = ... anywhere in this branch.

  • _cascade_force_skip_steps has the same shape: its skipped CTE updates step_states only, and its archived_messages CTE archives messages only. Task rows are untouched.
  • maybe_complete_run then 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

  1. A flow with a map step (initialTasks > 1) configured whenExhausted: 'skip' and maxAttempts: 1.
  2. Start a run; let task 0 of the map step fail its only attempt while tasks 1..n are still queued/started.
  3. fail_task skips the step; the run completes.
  4. 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_tasks requires task.status = 'queued' AND run.status = 'started' AND a step_states row at 'started' — so no handler ever runs. The damage is to the data model, not execution:

  1. Phantom in-flight work on terminal runs. Any dashboard, health check, or metric of the shape WHERE status = 'started' over-counts forever.
  2. Recovery/requeue tooling built on status = 'started' churns. We hit this concretely: a stalled-task sweeper keyed on started + age repeatedly "recovered" these rows on completed runs — flipping them started → queued, incrementing requeue counters until a permanent-stall marker fired, calling pgmq.archive on 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, and set_vt on an archived id matches nothing), but the operational signal is permanently noisy and wrong for any deployment using whenExhausted: 'skip' on a map step.
  3. Invariant asymmetry. Every other terminal step outcome terminalizes its task rows; skip is the one path that doesn't, which future consumers won't expect.

Suggested fix

In the same transaction that skips the step (both the fail_task skip branch and _cascade_force_skip_steps), terminalize the sibling rows alongside the message archive — e.g.:

UPDATE pgflow.step_tasks st
SET status = 'skipped'          -- or whatever terminal marker fits the schema
FROM ...
WHERE st.run_id = ... AND st.step_slug = ... AND st.status IN ('queued', 'started');

If adding a task-level 'skipped' status is undesirable, any terminal status with a marker would resolve the phantom-row problem equally well.

Activity

  1. jumski commented on Aug 17, 2026

    @jumski
    Contributor

    Thanks for the report @cpursley !
    I'm gonna review it soon!

  2. jumski commented on Aug 21, 2026

    @jumski
    Contributor

    Maintainer 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, and whenExhausted: 'skip':

    1. Tasks 0 and 1 were started.
    2. Task 2 remained queued.
    3. Task 0 failed its only attempt.
    4. fail_task skipped 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:      3
    

    This 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: its skipped CTE updates step_states, and its archived_messages CTE archives active task messages, but it does not update those task rows.

    Root cause

    pgflow.fail_task updates only the task whose handler failed. When that task exhausts retries with when_exhausted IN ('skip', 'skip-cascade'), the function:

    • changes the parent step_states row to skipped;
    • clears remaining_tasks;
    • decrements runs.remaining_steps;
    • archives queued and started sibling messages;
    • emits the existing step:skipped event;
    • continues dependency or cascade handling;
    • may complete the run.

    It never changes queued or started sibling step_tasks rows.

    pgflow._cascade_force_skip_steps has the same state mismatch for any newly skipped step that already has task rows.

    Late callbacks do not repair the mismatch. complete_task and fail_task guard against callbacks for a step that is no longer started, 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_steps become skipped;
    • 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 skipped as a valid step_tasks.status value.

    Do not reuse another existing status:

    • failed is incorrect because sibling handlers did not fail and would inflate failure metrics;
    • completed is 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_at or skip_reason column in this fix. The parent step_states row already stores the authoritative skipped_at and skip_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:skipped event 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.sql
    

    Primary 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.sql
    

    Documentation that lists runtime statuses:

    NOMENCLATURE_GUIDE.md
    pkgs/website/src/content/docs/concepts/data-model.mdx
    

    Generated 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>.md
    

    database-types.ts currently represents the status column as string, 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 0 started and later failed;
    • task 1 started;
    • task 2 queued;
    • 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 reaches completed where appropriate.

    The new assertion must initially fail because the schema rejects skipped for task rows and the functions do not write it.

    Extend archives_task_messages_for_skipped_steps.test.sql to assert that a directly cascaded active map step changes both its started and queued tasks to skipped while preserving message archival.

    Extend the late callback tests to assert the sibling task remains skipped after:

    • a late complete_task callback;
    • a late fail_task callback.

    Extend the preexisting-skip and idempotency tests to assert:

    • _cascade_force_skip_steps changes 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, change step_tasks.valid_status from:

    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, so output_valid_only_for_completed remains satisfied after they become skipped;
    • timestamp constraints do not require a terminal timestamp for each status;
    • completed_at_or_failed_at remains satisfied;
    • no foreign key depends on the status value.

    3. Terminalize siblings in fail_task

    In the existing branch guarded by:

    IF v_task_exhausted AND v_step_skipped THEN

    keep 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 NULL in 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 failed and 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_steps

    Refactor the current archived_messages portion 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 IDs
    

    Add 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_messages in the final query so PostgreSQL cannot optimize away the side-effecting archive select.

    Using the skipped CTE's RETURNING rows 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 started state 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 core

    Apply 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_task locks the run and parent step state before it transitions the step. complete_task and sibling fail_task callbacks lock the same run/step state and guard against a parent step that is no longer started.

    The task update must filter only queued and started rows:

    • 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_tasks or 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_tasks

    Review the generated migration. It must:

    1. drop and recreate step_tasks.valid_status with skipped allowed;
    2. replace both changed functions;
    3. 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 task
    

    Expected post-migration statuses:

    failed, completed, skipped, skipped
    

    Documentation and release metadata

    Update NOMENCLATURE_GUIDE.md so task statuses include:

    skipped - Task was cancelled because its parent step was skipped
    

    While touching the same section, include the already-supported skipped step status if it is still omitted.

    Update pkgs/website/src/content/docs/concepts/data-model.mdx so 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 core

    Also 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 skipped requires 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:

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

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions