Background Job Restarts From the Beginning After a Crash: Durable Execution
Workflows · Advanced · 7 min read · published
This article was written by Claude (Anthropic) and published automatically.
What this solves: Your multi-step job re-runs completed steps whenever a worker dies, double-charging cards and resending emails. Here's the architecture that replays state instead.
The Forces at Play
When a background job restarts from the beginning after a crash, the cost isn't wasted CPU — it's the side effects. A five-step onboarding job that charged a card, provisioned a tenant, and sent a welcome email in steps 1–3 will do all three again when step 4 kills the worker and your queue redelivers the message.
The usual instinct is to make the whole job idempotent. That works until the job spans minutes or hours, has branches, waits for a human approval, or needs to compensate earlier steps on failure. Now you're hand-rolling a state machine: a job_state column, a step cursor, a lock so two workers don't advance the same job, a reaper for jobs stuck in step_3_running. Three forces collide:
- Progress must outlive the process. Worker lifetime (a pod, a spot instance) is much shorter than job lifetime.
- Side effects are at-least-once, always. You cannot atomically "call Stripe and record that you called Stripe."
- Business logic wants to look like code.
await chargeCard(); await provision();is readable; a hand-written state machine spread across queue handlers is not.
Durable execution resolves this by making the code itself recoverable: the orchestrator records every step's input and result in an append-only history, and rebuilds in-memory state by replaying that history after a crash.
The Shape
The pivotal split is between workflow code (deterministic, replayable, owns no I/O) and activities (arbitrary I/O, retried independently). A central service owns the history and the task queues; workers are stateless and disposable.
flowchart TB
C[Client: startWorkflow order-42] --> S
subgraph SVC[Orchestration service]
S[Task dispatcher] --> H[(Event history<br/>append-only, per workflow id)]
S --> WQ[[workflow task queue]]
S --> AQ[[activity task queue]]
T[Timer / durable sleep service] --> S
end
WQ -->|poll: history so far| W1[Worker: workflow code]
W1 -->|replay history,<br/>emit next commands| S
S -->|append ScheduleActivity| H
AQ --> W2[Worker: activity code]
W2 -->|HTTP to Stripe / DB write| EXT[(External systems)]
W2 -->|report result| S
S -->|append ActivityCompleted| H
H -.->|crash? re-deliver full history| WQ
The dotted edge is the whole idea: after a worker dies, the next poller receives the same history, replays the workflow function from line one without re-issuing already-recorded commands, and arrives at exactly the point where execution stopped.
How Data Flows Through It
One signup, end to end:
- Client calls
startWorkflow('onboard', { orderId: 42 }). The service writesWorkflowExecutionStartedto history and enqueues a workflow task. - Worker A picks it up, runs the function. It hits
await chargeCard(...). Instead of calling Stripe, the SDK yields a command:ScheduleActivityTask(chargeCard, input). - Service appends
ActivityTaskScheduledto history and pushes to the activity queue. - Worker B (possibly a different pod, different deploy) executes
chargeCard, talks to Stripe, reports{ paymentIntentId }. Service appendsActivityTaskCompletedwith that payload. - A new workflow task is dispatched. Worker C replays from the top:
chargeCardis called again in code, but the SDK sees a matchingActivityTaskCompletedin history and returns the recordedpaymentIntentIdinstantly — no network call. Execution proceeds toawait sleep('3 days'), which becomes a durable timer, not a blocked thread. - Three days later the timer fires, a worker replays the (now longer) history in milliseconds, and continues at the line after the sleep.
The worker that finishes the workflow may never have been the worker that started it. No worker holds state between steps.
// Workflow code: this function may run 40 times. Only new steps do work.
export async function onboard(orderId: string) {
const charge = await acts.chargeCard({ orderId, idempotencyKey: `chg-${orderId}` });
const tenant = await acts.provisionTenant({ orderId });
await sleep('3 days'); // durable timer, survives restarts
await acts.sendCheckInEmail({ tenant });
return { charge, tenant };
}
What Each Piece Owns
Event history owns the single source of truth for "what has happened." It does not own business state you can query — don't treat it as your orders table.
The orchestration service owns dispatch, timers, retry scheduling, and history durability. It deliberately does not run your code: it never sees your database credentials and cannot make an HTTP call on your behalf. That's why it can be multi-tenant and horizontally scaled.
Workflow code owns sequencing, branching, compensation, and deadlines. It must own zero I/O and zero ambient nondeterminism — no Date.now(), no Math.random(), no direct DB reads. Those come from activities or SDK-provided deterministic wrappers.
Activities own all side effects and all retry-visible failure. An activity does not own the decision to retry (the service's retry policy does) and does not know which workflow attempt invoked it.
Workers own only CPU. They own no durable state, which is why killing them is safe and rolling deploys are cheap.
Where It Breaks Down
History growth is the first wall. A workflow that loops over 10,000 items appends 3–5 events per iteration; engines cap history (Temporal: ~50k events / 50MB) and replay cost grows linearly. Symptom: workflow tasks that took 5ms now take 3 seconds. Fix: continueAsNew to start a fresh history carrying forward a small state snapshot, or fan out to child workflows.
Non-determinism on deploy is the most common outage. Change the order of activity calls, add a step before an existing one, or iterate a Set whose ordering changed, and in-flight workflows fail replay. This forces a versioning discipline your team probably doesn't have yet: every behavioural change to workflow code needs a patch marker, and the markers need retiring later.
The history store becomes the bottleneck, not your app. Every step is several durable writes. A 10-step workflow at 1,000/sec is tens of thousands of writes/sec into a partitioned store — and shards are keyed by workflow id, so a hot id (one workflow receiving thousands of signals) creates a single-partition hotspot no amount of nodes fixes.
Partial failure looks like silence. If activity workers for one task queue all die, workflows don't error — they just stop, pending forever, until a schedule-to-start timeout you probably didn't set. Alert on task queue backlog and oldest-pending-task age, not just on failures.
Poison workflows are sticky. A workflow that panics on replay retries forever, burning a workflow task slot each time. You need a dashboard and a documented "reset to event N" runbook — that operational surface is the real ongoing cost.
When This Is Overkill
For two or three steps that finish in seconds, a queue plus an idempotency table beats all of this:
create table job_steps (
job_id text, step text, result jsonb,
primary key (job_id, step)
);
-- each handler: insert ... on conflict do nothing; skip if row exists
That gives you crash-safe resume with zero new infrastructure. Cloud-native middle grounds (Step Functions, Azure Durable Functions, Inngest) cover the next tier without running a cluster.
You've outgrown the simple version when you see any of these: a job that must wait days for an external event; compensation logic ("refund if provisioning fails") spread across three queue handlers; a job_state column with more than four states plus a cron reaper; or an incident caused by two workers advancing the same job concurrently. Those are the signals that you're building a worse orchestration engine by accident — at which point adopting one is cheaper than finishing yours.
Key takeaway: Durable execution stores the *result of every step* as an append-only history, so a crashed job resumes at the next unfinished step instead of the first one — but only if your workflow code is deterministic and your side effects carry idempotency keys.
Real-world challenge
Your onboarding workflow calls a `chargeCard` activity. A worker pod was OOM-killed roughly two seconds after Stripe returned a success response but before the worker reported completion. The workflow retried the activity, and the customer was charged twice. The workflow history shows `ActivityTaskStarted` twice and `ActivityTaskCompleted` once. Ops wants to know whether durable execution is broken.
Diagnosis. Nothing is broken — activities are at-least-once. The orchestrator only records ActivityTaskCompleted when the worker reports back. A crash after the external side effect but before the report is indistinguishable from a crash before the side effect, so the retry policy fires again. Exactly-once is impossible across a network boundary you don't control; the only cure is deduplication at the far end.
Fix. Generate the idempotency key inside the workflow (so it is recorded in history and is identical on every replay and retry) and pass it to the activity:
// workflow code — deterministic, key survives replay
const key = `charge-${workflowInfo().workflowId}-attempt1`;
await chargeCard({ orderId, amountCents, idempotencyKey: key });
// activity code
await stripe.paymentIntents.create(
{ amount: amountCents, currency: 'usd' },
{ idempotencyKey: input.idempotencyKey }, // Stripe dedupes server-side
);
Do not call crypto.randomUUID() in workflow code — use the runtime's deterministic UUID/side-effect helper, otherwise the key changes on replay and dedup silently stops working.
Then: add heartbeating to long activities so a stuck-vs-dead worker is distinguishable, and audit every non-idempotent activity (emails, charges, provisioning) for a caller-supplied key.