Cancellation & recovery
Three engine methods handle the endings that aren’t “every step completed”: one you call deliberately, one you run on a timer, and one you reach for when a run failed for a reason you have since fixed.
Cancelling a workflow
Section titled “Cancelling a workflow”const cancelled = await engine.cancelWorkflow(workflowId);if (!cancelled.ok) throw new Error(cancelled.error.message); // workflow_not_foundEvery pending and waiting step is marked skipped, the workflow finishes
cancelled, and a workflow.cancelled event is emitted. Three things worth knowing:
- It interrupts a running step only if that step beats. There is no signal into
an in-flight handler on its own, so by default the step runs to completion and its
result lands on a workflow that has already finished. A step with a
heartbeat learns it was cancelled on its next beat:
the engine fires
ctx.signaland discards whatever the handler returns. Without one, cancellation remains a “stop scheduling more work” operation, not a kill. - Compensation does not run. Saga rollback is wired to the failure path only (see Saga compensation). If you need completed steps undone on a cancel, do it yourself after the call returns.
- Cancelling a sub-workflow child fails the parent step, which cascades into the parent workflow exactly as a child failure would.
Calling it on a workflow that already reached a terminal state is a no-op that still
returns ok — so a double-click on a cancel button is safe.
Recovering after a crash
Section titled “Recovering after a crash”A worker that dies mid-step leaves its step in running forever. Nothing detects that
on its own: the queue’s job expiry releases the job, but the row stays claimed, and
an atomic claim (markStepRunning) means a redelivered job can’t take it over.
recoverStuckWorkflows is the sweeper that clears them:
const { retriedSteps, recoveredSteps, recoveredWorkflows, expiredWorkflows } = await engine.recoverStuckWorkflows();It scans every running workflow in the partition and does three things:
- Re-queues a crashed step that has been silent longer than its liveness window, provided its attempt budget has room. A dead pod costs an attempt, not the run.
- Fails a crashed step whose budget is spent, cascading the workflow to
failedin the usual way (dependents skipped, compensation run). - Fails a run past its deadline, which nothing else would notice while the run is suspended or idle.
Nothing calls this for you. Run it on a schedule in one process — a cron job, a
setInterval, or a pg-boss schedule:
setInterval(() => { void engine.recoverStuckWorkflows().catch((e) => logger.error('sweep failed', e));}, 60_000);Once per partition is enough. Two sweepers racing on the same partition will not corrupt
state — the step transition is a plain UPDATE and the cascade is recomputed from
listSteps — but each one increments the workflow’s failedSteps counter, so a
double sweep can inflate that progress number. Run it from one process, or accept the
skew.
The stuck threshold
Section titled “The stuck threshold”stuckThreshold = stepExpirySeconds + stuckStepBufferSecondsThis is the default window, measured from when a step started. A step type that
declares heartbeatTimeoutMs overrides it with its own,
measured from when the step last reported in — which is what lets the window be short
without condemning work that is merely slow.
Both come from WorkflowEngineConfig, and both have defaults:
| Option | Default | Meaning |
|---|---|---|
stepExpirySeconds |
600 |
What you told the dispatcher a step may occupy a worker for |
stuckStepBufferSeconds |
300 |
Grace on top, so a step that is merely slow isn’t swept |
onStuckStep |
'retry' |
'retry' re-queues within the attempt budget; 'fail' always fails |
Per step type, heartbeatTimeoutMs replaces the first two entirely for that type.
const engine = createWorkflowEngine({ store, dispatcher, registry, partitionKey, config: { stepExpirySeconds: 900, stuckStepBufferSeconds: 300 }, // sweep at 20 min});What recovery does not cover
Section titled “What recovery does not cover”The sweeper only looks at steps stuck in running. A step stuck in pending with no
job behind it is invisible to it — that is the dual-write window between persisting a
transition and enqueueing the jobs it unlocks.
Closing that window is the job of
transactional dispatch: with
store-pg + dispatcher-pgboss on one Postgres, the write and its enqueues commit
together and the window doesn’t exist. On a queue that lives elsewhere (SQS, Redis) the
engine falls back to write-then-enqueue, and the repair path is the dispatcher
redelivering the completed step’s job, which re-drives readiness.
Retrying a failed run
Section titled “Retrying a failed run”A workflow that used up its retries is failed, and failed is terminal. Starting a
fresh run repeats every side effect the first one already committed — the charge, the
email, the file that was written. retryWorkflow is the alternative: resume the
existing run from where it stopped.
const retried = await engine.retryWorkflow(workflowId);if (!retried.ok) throw new Error(retried.error.message);retried.value.resetSteps; // ['fulfil', 'notify'] — what will run againEvery step that failed, was skipped in the fallout, or had its work
compensated away goes back to pending with a fresh attempt budget, its output,
error and timestamps cleared. Steps that completed are left exactly as they are —
they keep their output, and they do not run again. The workflow returns to running
with recomputed counters, and the frontier is dispatched.
→ examples/16-deadlines-and-retry.ts
Details worth knowing:
- A compensated step does re-run. Saga rollback undid its effects, so its work has
to happen again — that is why
compensatedis reset alongsidefailedandskipped. - A map parent’s children are dropped before it re-runs, so it fans out afresh rather than aggregating two generations of items.
- A guard is re-evaluated. A step skipped by a
whenguard is reset too, so the branch decision is made again against the current data. - It refuses a run that is not
failed. Acompleted,runningorcancelledworkflow returns aworkflow_not_retryableerror rather than being restarted. - It refuses a sub-workflow child, pointing you at the parent. Retrying a child cannot un-fail the parent step that was waiting on it; retry the parent, and it starts a fresh child.
- A
workflow.retriedevent is emitted for your run history.