An agent that answers a chat message in ten seconds can simply be retried when it fails. An agent that onboards a vendor cannot: it collects documents, calls a sanctions-screening API, waits two days for a human approval, creates a record in the ERP, registers a bank payee and sends a welcome email. Somewhere in those days a pod will be rescheduled, a deploy will roll, an API will return 503 and the LLM will produce malformed JSON. The question is not whether the run is interrupted but whether it resumes correctly, without creating the vendor twice or paying the wrong account.
Durable execution is the answer used by workflow engines: the orchestration logic is ordinary code, every side effect is recorded, and after any interruption the code is replayed against the record until it reaches the first step that has not happened yet. Agent checkpointing covers the event log, snapshots and resume mechanics. This article covers the programming model on top: what must be deterministic, how to classify failures for retry, which timeouts you need, how to undo partial work, how to wait for people, and how to deploy new code while runs are in flight.
The core split: workflow code and activities
A durable workflow separates two kinds of code. Workflow code decides what happens next: call this tool, branch on its result, wait for approval. It must be deterministic, because it will be run again, from the top, every time the run wakes up. Activities do everything that touches the world: LLM calls, HTTP requests, database writes, emails. The engine records each activity's result in the run's history, and on replay it hands the recorded result back instead of calling again.
An LLM call is an activity, not workflow code. Its output is non-deterministic, costs money and takes seconds; recording it means replay never re-asks the model, and the agent's past decisions stay exactly as they were. The workflow then branches on the recorded output as if it were any other value.
What breaks replay
Replay works only if the workflow code, given the same history, issues the same sequence of commands. Anything whose value can differ between the original run and the replay breaks that guarantee. Engines detect some of these as non-determinism errors; others silently take a different branch.
| Pattern in workflow code | Why it breaks | Do instead |
|---|---|---|
| Reading the wall clock | Different value on replay | Engine-provided time, recorded once |
| Random numbers or new UUIDs | Different value on replay | Generate inside an activity, or a recorded helper |
| Calling an LLM or API directly | Different result, repeated side effect | Wrap it in an activity |
| Reading config or feature flags | Value may change between runs | Read in an activity, or pass as input |
| Iterating an unordered set | Order may differ | Sort before iterating |
| Threads or native sleeps | Invisible to the engine | Engine timers and concurrency primitives |
| Changing step order in a deploy | Old histories no longer match | Version branches (see below) |
A minimal durable executor
Engines hide the mechanics, but they are simple enough to write down. The executor below keeps a per-run history. Each call to activity either returns the recorded result, if replay has not yet reached the end of history, or runs the function, records the result durably and then returns it. Waiting for a signal that has not arrived parks the run.
import time
class Suspend(Exception): pass
class NonDeterminism(Exception): pass
class ActivityFailed(Exception): pass
class Ctx:
def __init__(self, run_id, store):
self.run_id, self.store = run_id, store
self.history = store.load(run_id) # list of dicts, oldest first
self.pos = 0
def _replayed(self, kind, name):
if self.pos >= len(self.history):
return None
ev = self.history[self.pos]
if (ev["kind"], ev["name"]) != (kind, name):
raise NonDeterminism(f"step {self.pos}: history {ev['kind']}:{ev['name']}, code {kind}:{name}")
self.pos += 1
return ev
def _record(self, ev):
self.store.append(self.run_id, ev) # durable before the result is used
self.history.append(ev)
self.pos += 1
def activity(self, name, fn, *args, policy):
ev = self._replayed("activity", name)
if ev is None:
key = f"{self.run_id}:{self.pos}" # stable idempotency key for the side effect
try:
ev = {"kind": "activity", "name": name, "ok": True,
"value": call_with_retries(fn, args, key, policy)}
except Exception as e:
ev = {"kind": "activity", "name": name, "ok": False, "value": repr(e)}
self._record(ev)
if not ev["ok"]:
raise ActivityFailed(ev["value"])
return ev["value"]
def now(self):
ev = self._replayed("clock", "now")
if ev is None:
ev = {"kind": "clock", "name": "now", "value": time.time()}
self._record(ev)
return ev["value"]
def wait_signal(self, name, deadline):
ev = self._replayed("signal", name)
if ev is None:
raise Suspend(name, deadline) # engine appends the signal or a timeout later
return ev["value"]When a signal arrives, or its deadline passes, the engine appends a signal event, with None meaning timed out, and runs the workflow again from the top. Every earlier step replays from history in microseconds, and execution continues at the wait. A crash between an activity's side effect and its record is covered by the idempotency key, which is why it is derived from the run and position rather than generated fresh.
The worked example: onboarding a vendor
IO = Policy(max_attempts=6, initial=1.0, backoff=2.0, max_interval=60)
LLM = Policy(max_attempts=3, initial=2.0, backoff=2.0, max_interval=20)
FOREVER = Policy(max_attempts=50, initial=5.0, backoff=2.0, max_interval=900)
def onboard_vendor(ctx, req):
docs = ctx.activity("fetch_documents", fetch_documents, req["vendor_id"], policy=IO)
fields = ctx.activity("extract_fields", llm_extract_fields, docs, policy=LLM)
screen = ctx.activity("sanctions_screen", sanctions_screen, fields["legal_name"],
fields["country"], policy=IO)
if screen["match"]:
return {"status": "rejected", "reason": "screening match"}
ctx.activity("request_approval", notify_approver, req["vendor_id"], fields, policy=IO)
decision = ctx.wait_signal("approval", deadline=ctx.now() + 3 * 86400)
if decision is None:
ctx.activity("escalate", notify_manager, req["vendor_id"], policy=IO)
decision = ctx.wait_signal("approval_escalated", deadline=ctx.now() + 2 * 86400)
if not decision or not decision["approved"]:
return {"status": "declined"}
return create_vendor_with_compensation(ctx, fields)Read it as a plain program: no state machine, no status column, no cron job looking for stuck rows. The three-day wait costs nothing while it waits, because the run is only a history in storage. If a worker dies after the screening call, the next worker replays fetch_documents, extract_fields and sanctions_screen from history and continues. Compare this with the explicit state-machine approach in agent state machines: durable execution gives the same guarantees while keeping control flow in code.
Retries: classify before you retry
A retry policy is only as good as the failure classification behind it. Retrying everything wastes money on LLM calls and hammers failing services; retrying nothing turns every blip into a failed run. Use three classes.
- Transient: timeouts, 429s, 503s, connection resets. Retry with exponential backoff, full jitter and a cap.
- Non-retryable: validation errors, 400s, 404s, permission denied, business rejections. Fail the activity immediately and let the workflow decide what to do.
- Model output errors: malformed JSON, schema violations, refusals. Retry a small number of times, feeding the validation error back into the prompt, then fail as non-retryable. A fourth identical attempt rarely succeeds where three failed.
import random, time
class NonRetryable(Exception): pass
def call_with_retries(fn, args, key, p):
for attempt in range(1, p.max_attempts + 1):
try:
return fn(*args, idempotency_key=key)
except NonRetryable:
raise
except Exception:
if attempt == p.max_attempts:
raise
cap = min(p.max_interval, p.initial * p.backoff ** (attempt - 1))
time.sleep(random.uniform(0, cap)) # full jitterThis sketch sleeps in-process for brevity; real engines schedule each retry as a durable timer, so a retry waiting fifteen minutes survives a restart. Tool-level retry design is covered in more depth in agent tool retries, and making the downstream side safe to call twice in idempotency for agent-to-agent calls.
Timeouts: four different questions
| Timeout | Question it answers | Typical setting for an LLM activity |
|---|---|---|
| Per attempt | How long may one try take? | A little above p99 latency for your prompt size |
| Total for the activity | How long across all retries? | Minutes; beyond that, fail and decide in the workflow |
| Heartbeat | Is a long activity still alive? | For streaming or batch jobs: seconds to a minute |
| Workflow deadline | When does the business stop caring? | Days for onboarding; enforce with a timer |
Missing per-attempt timeouts are the most common cause of stuck runs: an HTTP client with no timeout waits forever, and the engine cannot tell a slow call from a dead worker. Heartbeats let the engine detect a worker that died mid-activity and reschedule it long before the total timeout expires.
Compensation: undoing partial work
Distributed side effects cannot be rolled back in one transaction. If the ERP record is created and payee registration then fails permanently, you must undo the ERP record explicitly. The saga pattern records a compensation after each successful step and runs them in reverse on failure.
def create_vendor_with_compensation(ctx, fields):
undo = []
try:
vendor_id = ctx.activity("erp_create", erp_create_vendor, fields, policy=IO)
undo.append(("erp_deactivate", erp_deactivate_vendor, vendor_id))
ctx.activity("bank_register", register_payee, vendor_id, fields["iban"], policy=IO)
undo.append(("bank_remove", remove_payee, vendor_id))
ctx.activity("welcome_email", send_welcome, vendor_id, policy=IO)
return {"status": "active", "vendor_id": vendor_id}
except ActivityFailed:
for name, fn, arg in reversed(undo):
ctx.activity(name, fn, arg, policy=FOREVER)
raiseCompensations are activities too, so they are recorded, retried and replayed like any other step, and they must be idempotent. Some steps cannot be undone, such as an email that has been sent; put those last. When a compensation itself keeps failing, stop and hand the run to a person rather than looping indefinitely.
Waiting for people and for time
Signals deliver external events into a running workflow: an approval, a cancellation, a corrected document. Durable timers wake it at a future time, even weeks away. Together they replace the polling jobs and status tables that usually surround human-in-the-loop agents. Always pair a signal wait with a deadline and an explicit path for the timeout, as the escalation in the example does; a wait with no deadline is a run that can hang forever. Deliver signals with an id so a double-clicked approval is recorded once. Human-in-the-loop agent design covers what the human should see and decide.
Deploying new code while runs are in flight
A run started last Tuesday will replay against whatever code is deployed today. If today's code inserts a new activity before sanctions_screen, Tuesday's history no longer lines up and replay fails with a non-determinism error. There are three safe strategies. Version branches: record a version marker the first time a run reaches the changed point, and branch on it, so old runs take the old path and new runs the new one; remove the old branch once no run carries it. Worker pinning: keep old workers running old code for runs that started on it, and route new runs to new workers until the old ones drain. Restart with carried state: for very long-lived workflows, periodically finish the run and start a fresh one with the current state as input, which also keeps history from growing without bound. Temporal's documentation calls this Continue-As-New. Test every deploy by replaying a sample of real production histories against the new code in CI before release.
Choosing an engine
| Option | Model | Fits when |
|---|---|---|
| Temporal | Workflow code in SDK languages, history in a server cluster | Many long-running workflows, multiple teams, strong tooling needs |
| Restate | Durable execution runtime journaling handler steps | Service-style handlers that need durable steps and state |
| DBOS | Library that checkpoints workflow steps to Postgres | You already run Postgres and want no extra cluster |
| AWS Step Functions | State machines defined in JSON (Amazon States Language) | AWS-native pipelines, visual graphs, little custom code |
| LangGraph checkpointers | Graph state persisted per thread after each node | Agent graphs needing resume and interrupts, lighter guarantees on side effects |
| Home-grown | Event log plus executor, as above | One narrow workflow; understand the maintenance cost first |
Check each engine's current documentation for limits such as maximum history size and payload size before you design around it. Whichever you choose, the programming rules in this article apply unchanged.
Failure modes and operations
- Non-determinism after a deploy. Replay production histories in CI before every release.
- Large payloads in history. Storing whole documents or long LLM transcripts bloats history and slows replay. Store them in object storage and record references.
- Duplicate side effects. Crashes between the call and the record repeat activities. Pass the idempotency key to every downstream system that accepts one.
- Retry storms. Thousands of runs retrying one failing service in lockstep. Jitter, caps and a circuit breaker at the activity layer.
- Stuck waits. Signal waits without deadlines. Alert on runs waiting longer than their business deadline.
Monitor runs started, completed and failed per workflow type, activity retry and failure rates by error class, LLM spend per run, replay time, history size, and the age of the oldest waiting run. Trace each run end to end with its run id on every span, as described in agentic observability and tracing. The companion ai-data articles on graph RAG and on knowledge lineage cover the retrieval systems such agents usually read from.
What to do next
- List your agent's side effects and move each one, including every LLM call, into an activity.
- Remove clock reads, randomness, config reads and unordered iteration from orchestration code.
- Classify each activity's errors as transient, non-retryable or model-output, and set retry policies per class.
- Set per-attempt and total timeouts on every activity, and heartbeats on long ones.
- Pass a run-and-step idempotency key to every downstream call that accepts one.
- Write compensations for every step that creates external state, and put irreversible steps last.
- Give every human wait a deadline and an escalation path.
- Add a CI job that replays recent production histories against new code before deploying.