Composing durable work: steps, parallel fan-in, and sagas
The agent loop journals its own model turns and tool calls. But the same durable substrate is exposed directly, so you can wrap your own work in the same at-most-once, crash-safe, independently-provable guarantee. This is the "Option B" authoring model: write plain Go control flow (if / for / functions), and name the operations that must survive a crash. There is no graph DSL; the graph is a derived output (see DEBUGGING.md RenderMermaid).
Three primitives cover the common shapes, all built on the Durable port (EXTENSION-POINTS.md):
Step[T]: one named durable step.Parallel[T]/Task[T]: durable fan-out/fan-in.RunSaga+CompensatedFunc: transactional agents with reverse-order rollback.
Step[T]: one named durable operation
func Step[T any](ctx context.Context, d Durable, runID, name string,
fn func(context.Context) (T, error)) (T, error)Step runs fn as a named durable step keyed by (runID, name) and returns its typed result. On resume, a completed step returns its recorded result without re-running fn; if fn errors, nothing is recorded and the step re-runs on the next attempt. It is the building block the agent loop itself is made of, exposed for your own orchestration:
inv, err := agent.Step(ctx, store, runID, "fetch-invoice",
func(ctx context.Context) (Invoice, error) { return billing.Lookup(ctx, id) })Step is a package function, not a method, because Go methods cannot add type parameters. The result is journaled as a StepValue record, so it shows up in RenderMermaid as step: <name> and is independently provable via audit.ProveStep (see AUDIT.md). name must be unique within the run: a second Step with the same (runID, name) returns the first one's recorded result.
Step is the idempotency guard the MESSAGING.md webhook pattern uses to make a redelivered inbound event replay instead of re-fire.
Parallel[T] / Task[T]: durable fan-in
type Task[T any] struct {
Name string
Fn func(context.Context) (T, error)
}
func Parallel[T any](ctx context.Context, d Durable, runID string,
maxConcurrency int, tasks ...Task[T]) ([]T, error)Parallel runs each Task concurrently, each as its own durable Step, and returns the results in task order (not completion order). It is the durable, auditable fan-in that a compliance pipeline wants: run several independent checks at once (sanctions, credit, fraud), each crash-safe and at-most-once, each result committed to the journal and provable on its own, then aggregate.
- Each task's
Nameis its durable memoization key within the run, so it must be unique across the tasks in one call. - All tasks run even if some fail, so a failed check never hides the others. The returned error joins every task's error (
errors.Join) and isnilonly if all succeeded. A failed task was not journaled, so a later resume re-runs it while succeeded tasks are memoized. maxConcurrencycaps in-flight tasks;<= 0means one goroutine per task.
checks := []agent.Task[CheckResult]{
{Name: "sanctions_check", Fn: runSanctions},
{Name: "pep_check", Fn: runPEP},
{Name: "adverse_media_check", Fn: runAdverseMedia},
}
results, err := agent.Parallel(ctx, store, runID, 0, checks...) // 0 = unbounded concurrencyThis is deliberately a thin primitive over the journal, not a graph engine. Dynamic, model-driven routing stays in plain Go and sub-agents; Parallel covers the static fan-out/fan-in that a governed workflow's "parallel checks, then decide" stage is made of. The full worked flow (parallel durable checks, then a governed decision, then an offline proof) is examples/compliance.
Related but distinct: the agent loop already runs a single turn's tool calls concurrently (bounded by
SetMaxConcurrency).Parallelis for fan-out you author yourself outside a model turn. See the parallel-tool concurrency notes in KNOWN-LIMITATIONS.md.
Sagas: transactional agents with reverse-order compensation
A saga is the sequential/hierarchical transactional tier: "charged the card and booked the flight, then failed on the hotel, so cleanly refund and cancel." A tool declares how to undo its write with the Compensator port (EXTENSION-POINTS.md); CompensatedFunc builds a typed tool that carries both the forward action and its undo:
book := agent.CompensatedFunc("book_flight", "book a flight", agent.Safety{},
func(ctx context.Context, in BookArgs) (Booking, error) { return airline.Book(ctx, in) },
func(ctx context.Context, in BookArgs, out Booking) error { return airline.Cancel(ctx, out.PNR) })Run the agent with RunSaga instead of Run. If a step fails after earlier compensatable writes succeeded, RunSaga rolls those writes back in reverse order (recursing through sub-agent trees) and returns *SagaAborted:
_, err := a.RunSaga(ctx, runID, input)
var aborted *agent.SagaAborted
if errors.As(err, &aborted) {
// aborted.Cause is why it rolled back; the completed writes were compensated.
}RunSagaResult is the *Result-envelope counterpart (see DEBUGGING.md), and StreamSaga is the streaming counterpart of Stream. SagaAborted unwraps to its cause, so errors.Is against a sentinel still works.
How it survives a crash
The abort is derived from the journal: a saga step failure writes a durable StepSagaFail record, so a crash at any point resumes correctly. On re-entry a recorded failure sends the run straight to rollback, and each compensation is itself a durable memoized step, so it runs at-most-once if it completes. This split is the precise contract, proven adversarially in saga_dst_test.go (see TESTING.md):
- The forward non-idempotent effect is at-most-once (halt on unknown outcome).
- Compensators are at-least-once: memoized so they run once if they complete, but a crash mid-compensation re-runs them. This is why a
Compensatemust be idempotent.
The boundaries (read these before relying on it)
These are stated in full in KNOWN-LIMITATIONS.md; in brief:
- A saga step must be atomic. The failing step itself is not compensated (there is no recorded result to drive
Compensate), so a forward step must not leave a partial external side effect before returning an error. Make forward steps all-or-nothing or idempotent. - Unknown-outcome resume halts, it does not auto-roll-back. If a non-retriable step crashes after its attempt marker but before any result,
RunSagareturns*ResumeHalt: a human decides, because you cannot safely roll back a step that may have committed. - Compensation is hierarchical, not concurrent. Rollback recurses through a sub-agent tree (one causal order). Truly concurrent agents mutating shared state out of order need the provable convergence of the governance tier (GOVERNANCE.md), not a saga.
SubAgent(name, description, sub) composes agents into that durable tree: give the sub-agent the same Durable store as the parent for a unified journal, and a crash anywhere in the tree resumes the whole tree precisely (completed sub-agents reused, the in-flight one resumed, and ResumeHalt / PendingApproval / SagaAborted from deep in the tree propagating up). A ResumeHalt or PendingApproval raised inside a sub-agent surfaces from the parent's own Run (match it with errors.As); resolve it against the sub-run's ID and tool-use ID carried on the signal, then re-run the parent to resume down the path.
Clearing a ResumeHalt: ResolveHalt
A ResumeHalt is deliberately terminal until a human confirms the real outcome: the runtime cannot know whether the non-idempotent side effect (a charge, a send) actually committed. Once you have verified it out of band, ResolveHalt is the sanctioned escape. It injects the missing tool result under the halted tool-use ID (the same journal key the loop uses), so a re-run proceeds past the halt instead of halting again:
var halt *agent.ResumeHalt
if errors.As(err, &halt) {
// operator confirmed the charge did go through
_ = agent.ResolveHalt(ctx, store, halt.RunID, halt.ToolUseID, "charged (confirmed)", false)
msg, err = a.Run(ctx, halt.RunID, input) // resumes past the halt
}Pass the result value an actual call would have returned, and isError=true if the verified outcome was a failure the model should react to. It is idempotent (first result for a (runID, toolUseID) wins), so a retry or a racing driver injects it at most once.
Durable timers: Sleep / WaitUntil and the Waker
An agent often has to wait: for a deadline, a cool-off, a scheduled follow-up. Sleep(ctx, name, d) and WaitUntil(ctx, name, until) make that wait durable. Called from inside a retry-safe tool, they journal the wake time once (at-most-once by name) and pause the run with *Sleeping, the same durable pause as Interrupt. Because the wake time is journaled on the first call and memoized, a resumed or crash-recovered run waits to the same absolute instant rather than restarting the clock; no goroutine is held blocked across the wait.
wait := agent.Func("cooldown", "wait before retrying", agent.Safety{ReadOnly: true},
func(ctx context.Context, _ struct{}) (string, error) {
if err := agent.Sleep(ctx, "cooldown", time.Hour); err != nil {
return "", err // *Sleeping propagates out of Run; the run is paused durably
}
return "resumed", nil
})Re-invoking Run with the same runID at or after the wake time resumes past the Sleep. What re-invokes it is a Waker, the time-driven sibling of the inbound event trigger in MESSAGING.md: the SDK provides the durable, at-most-once timer and its resume safety, and the trigger is pluggable. Bind one with agent.WithWaker(ctx, w) and Sleep registers its wake automatically. MemWaker is the reference in-process implementation:
w := agent.NewMemWaker(func(ctx context.Context, runID string) error {
_, err := a.Run(agent.WithWaker(ctx, w), runID, savedInput) // resume; may sleep again
return err
})
w.Start(ctx, time.Second, nil) // tick: resume every run whose timer is dueBoundaries: MemWaker is a local-dev default, not a durable scheduler. Its in-memory timer set is lost on process exit, so the wake times must also live in the journal (they do), and a restarted deployment rebuilds pending wakes by scanning runs or hands the trigger to an external scheduler (cron, a queue). Tests inject a clock with agent.WithClock to advance time deterministically. This is the piece that makes an always-on ambient agent turnkey: a durable wait plus a trigger, with at-most-once and crash-resume intact across the wait.