Pidoku

Orchestration and Durable Execution

Advanced 55 min Difficulty 4/5 Lesson 02 of 05

Prerequisites Anatomy of an Agent

The idea in one minute#

An agent task can run for an hour, wait a day for a human, and span a hundred model and tool calls. In that time the process running it will be restarted — by a deploy, an autoscaler, a crash. If the task’s progress lives in that process’s memory, it is lost, and so is the money already spent. Durable execution fixes this by recording every completed step in a persistent log, so that a task can be resumed anywhere, from exactly where it stopped, without repeating work or repeating side effects.

The idea is simple: the log is the truth; the process is disposable.

A picture#

flowchart LR
  API[":i-user: <b>Start task</b>"] --> Q[":natsdotio: <b>Task queue</b>"]
  Q --> W1[":i-bot: <b>Worker A</b><br/><small>runs the loop</small>"]
  W1 -->|"append each step"| LOG[(":postgresql: <b>Event history</b><br/><small>step 1: model call → result<br/>step 2: tool call → result<br/>step 3: waiting for approval</small>")]
  W1 -.->|"crashes"| X[":i-skull: <b>gone</b>"]
  LOG -->|"replay recorded results"| W2[":i-bot: <b>Worker B</b><br/><small>resumes at step 3</small>"]
  H[":i-hand: <b>Human approves</b>"] -->|"signal"| LOG
  W2 --> ACT[":i-wrench: <b>Step 4: tool</b><br/><small>with idempotency key</small>"]
  ACT --> LOG
  class API,H neutral
  class Q queue
  class W1,W2 compute
  class LOG memory
  class X,ACT warn

How it really works#

What can interrupt a task#

InterruptionWithout durabilityWith it
Worker crash or deployTask lost; restart from zeroAnother worker replays the log and continues
A tool is down for ten minutesTask failsThe step retries with backoff; the task waits
Human approval takes a dayA process sleeps for a day, if it survivesThe task is parked in storage; nothing runs
Rate limit from the model providerErrorThe step is rescheduled
The user closes the browserStream dies with the requestThe task continues; the client reconnects to its events

Two ways to get durability#

Checkpointing. After each step, write the agent’s state — messages, pending tool calls, counters — to a store. To resume, load the latest checkpoint and continue. This is what LangGraph’s checkpointers and most agent SDK session stores do. It is simple and enough for most agents.

Event-sourced replay. Record each step’s result in an append-only history. To resume, re-run the orchestration code from the start, but for every step already in the history, return the recorded result instead of executing it. The code races through the past and continues at the first step with no record. This is the model of Temporal, Restate and similar engines. It handles timers, signals, child tasks and versioning in a uniform way.

Replay imposes one rule: the orchestration code must be deterministic. Anything non-deterministic — a model call, a tool call, the current time, a random number — must be wrapped as a recorded step (an activity). Model calls are the obvious case: replaying must return the same completion that was recorded, not generate a new one.

Idempotency: the other half#

Durability guarantees a step is recorded after it finishes. A crash between “the email was sent” and “the result was recorded” means the step runs again. So every side-effecting tool must be safe to run twice:

  • Give each step a stable idempotency key (task ID + step number) and pass it to the tool.
  • The tool’s backend deduplicates on the key.
  • Where the backend cannot, use check-then-act: look up whether the effect already happened.

At-least-once execution plus idempotent tools gives effectively-once effects. There is no third way.

Human in the loop#

Approval is an interrupt: the task records “waiting for approval of action X”, releases its worker, and resumes when a signal arrives. Design details that matter:

  • Show the human the exact action and arguments, not the model’s summary of them.
  • The approval is bound to that action; if the agent changes the arguments, it must ask again.
  • A timeout has a defined outcome: cancel, escalate or proceed with the safe alternative.
  • Record who approved what and when. This is audit data.

Streaming and reconnection#

Users want to watch. Decouple the task from the connection: the worker appends events — tokens, tool calls, status — to a per-task stream, and clients subscribe to it with a cursor. A dropped connection resumes from the cursor. The same stream feeds the UI, the trace and the audit log.

Concurrency inside a task#

Steps that do not depend on each other should run in parallel: several searches, several file reads, several sub-agents. The orchestrator fans out, waits for all (or the first, or a quorum), and continues. Parallel steps are recorded individually, so a crash mid-fan-out resumes only the unfinished ones. Cap the fan-out width; a model that decides to spawn 200 parallel searches is a budget problem.

Versioning running tasks#

A task started yesterday is running yesterday’s orchestration code and prompt. When you deploy a change, in-flight tasks either finish on the old version or must be explicitly migrated. Durable engines provide version markers for this; with plain checkpoints, store the code and prompt version in the checkpoint and route resumed tasks to a compatible worker. Changing a system prompt under a half-finished task also discards its cached prefix and can confuse the model about what it already decided.

Choosing the machinery#

SituationUse
Tasks of seconds to minutes; losing one is tolerableIn-memory loop with a session store
Minutes to hours; approvals; must resumeAgent framework checkpointing on Postgres
Hours to days; many timers and signals; strict guarantees; many task typesA durable execution engine: Temporal, Restate, or a cloud equivalent

The agent framework and the durable engine combine: the framework defines the loop, and each model or tool call becomes an activity of the engine. Temporal ships integrations of this kind for several agent SDKs.

Code#

Event-sourced replay in miniature: a journal of step results, a crash, and a resume that does not repeat the completed steps.

Go
// durable.go — journal each step's result; on resume, replay results instead of re-running.
package main

import (
	"errors"
	"fmt"
)

// journal is the durable event history. In production it is a database table.
type journal struct{ results []string }

type run struct {
	j        *journal
	cursor   int
	executed int // side effects actually performed in this process
}

var errCrash = errors.New("worker crashed")

// step returns the recorded result if the step already ran, else executes and records it.
func (r *run) step(name string, do func() (string, error)) (string, error) {
	if r.cursor < len(r.j.results) {
		res := r.j.results[r.cursor]
		r.cursor++
		fmt.Printf("  replay   %-16s → %s\n", name, res)
		return res, nil
	}
	res, err := do()
	if err != nil {
		return "", err
	}
	r.j.results = append(r.j.results, res) // durable write happens before moving on
	r.cursor++
	r.executed++
	fmt.Printf("  execute  %-16s → %s\n", name, res)
	return res, nil
}

// task is deterministic orchestration: the same steps in the same order on every replay.
func task(r *run, crashAt int) error {
	steps := []string{"model: plan", "tool: fetch data", "model: analyse", "tool: write file", "model: summarise"}
	for i, name := range steps {
		_, err := r.step(name, func() (string, error) {
			if i == crashAt {
				return "", errCrash
			}
			return fmt.Sprintf("ok#%d", i+1), nil
		})
		if err != nil {
			return err
		}
	}
	return nil
}

func main() {
	j := &journal{}

	fmt.Println("worker A:")
	a := &run{j: j}
	fmt.Println("  result:", task(a, 3)) // crashes while executing step 4

	fmt.Println("worker B, same journal:")
	b := &run{j: j}
	fmt.Println("  result:", task(b, -1))

	fmt.Printf("\nsteps executed: A=%d, B=%d; total %d for a 5-step task — nothing ran twice.\n", a.executed, b.executed, a.executed+b.executed)
}

Step 4 did run twice in attempts — once failing in A, once succeeding in B. If A had finished the write and died before recording it, B would write again. That is the gap an idempotency key closes.

Remember this#

  • The log is the truth; the process is disposable.
  • Checkpointing is enough for most agents; event-sourced replay for long, complex, guaranteed tasks.
  • Orchestration code is deterministic; model and tool calls are recorded activities.
  • At-least-once execution plus idempotent tools equals effectively-once effects.
  • Approvals are interrupts, bound to the exact action, with a timeout and an audit record.

Try it#

  1. Run durable.go. Move the crash to step 1. How much work does worker B repeat?
  2. Make step 4 non-idempotent in your head — an email send. Describe the exact crash timing that sends two emails, and the fix.
  3. Design the approval record for “merge a pull request”: which fields must it contain?

Check yourself#

  1. Why must a model call be a recorded activity under replay?
  2. What does an idempotency key protect against that durable execution alone does not?
  3. What should happen to tasks in flight when you deploy a new prompt?

↑↓ navigate↵ openesc close

drag to pan · scroll to zoom