Workflow Engines
Learning Outcomes
- Decide which pipeline steps belong in deterministic code and which belong to an agent
- Wrap a non-deterministic agent call as a workflow activity with explicit timeouts and retries
- Configure retry policies and backoff for AI operations across Temporal, Prefect, and Airflow
- Maintain a durable audit trail of every agent decision through workflow event history
- Design a hybrid workflow combining validation, agent analysis, and human approval
Lesson Plan
| Segment | Duration | Topic |
|---|---|---|
| Intro | 3 min | Why deterministic engines wrap non-deterministic agents |
| Concept | 8 min | The determinism boundary and durable execution |
| Build | 9 min | Wrapping an agent call as a Temporal activity |
| Build | 10 min | Retry policies, timeouts, and idempotency for AI steps |
| Build | 8 min | Audit trails from event history |
| Compare | 7 min | Temporal vs. Prefect vs. Airflow for agent work |
| Build | 3 min | A hybrid validate → analyse → approve workflow |
| Wrap-up | 2 min | Key takeaways and what comes next |
Before You Begin
Pre-work:
- Complete Lesson 5: The Anthropic Agent SDK — you will reuse the agent built there
- Review Orchestration Patterns for the supervisor and pipeline patterns
- Recall agent safety from the Agentic AI course — workflow engines are where those guardrails get enforced operationally
Shopping List:
- Python 3.10+ and a virtual environment
pip install claude-agent-sdk(ships the Claude Code CLI bundled)- A Temporal dev server:
temporal server start-dev(download from temporal.io) - Optional:
pip install prefectandpip install apache-airflowfor the comparison - An
ANTHROPIC_API_KEYexported in your shell
The agent from Lesson 5 is non-deterministic by design: the model chooses each next step. That is the point of an agent — but it is poison for the parts of a pipeline that must behave identically every run: schema validation, idempotent writes, billing, compliance gates.
A workflow engine lets you draw a line. On one side, deterministic orchestration code the engine guarantees runs exactly once to completion, surviving crashes, deploys, and multi-day waits. On the other, the messy retryable calls to the outside world — HTTP, databases, and now LLM agents. In Temporal those outside-world calls are Activities:
Workflow (deterministic): validate() → call_agent() → persist() → notify()
(code) (Activity) (Activity) (Activity)
the agent loop (Claude + tools) runs inside call_agent(),
retried, timed-out, and audited by the engine
A Temporal Workflow Execution is a durable function execution with no imposed time limit — it runs effectively once and to completion whether it takes seconds or years. On every worker restart the engine replays the workflow code but only re-executes what is missing; activities that already succeeded are skipped and their recorded results returned.
A thin Activity wraps the Claude Agent SDK; the workflow only sees a plain function call. The SDK exposes query() for one-shot runs and ClaudeSDKClient for interactive sessions, both configured with ClaudeAgentOptions.
# activities.py
from dataclasses import dataclass
from temporalio import activity
from claude_agent_sdk import query, ClaudeAgentOptions, AssistantMessage, TextBlock
@dataclass
class AnalysisRequest:
document_id: str
text: str
@activity.defn
async def analyse_document(req: AnalysisRequest) -> str:
options = ClaudeAgentOptions(
system_prompt="You are a contract risk analyst. Return findings as JSON.",
max_turns=6,
allowed_tools=["Read", "Grep"],
)
output = ""
async for message in query(prompt=req.text, options=options):
if isinstance(message, AssistantMessage):
output += "".join(b.text for b in message.content
if isinstance(b, TextBlock))
activity.heartbeat(f"turn for {req.document_id}") # progress signal
return output
The workflow stays deterministic: it imports the Activity (inside workflow.unsafe.imports_passed_through()) and calls workflow.execute_activity(analyse_document, req, ...) with the timeouts from Step 3. The heartbeat() matters for agents — a turn can take minutes, and without heartbeats the engine cannot tell a slow-but-healthy run from a dead worker.
Agent calls fail in ways HTTP calls do not: rate limits, overloaded models, tool errors, responses that never converge. The engine bounds these with four timeout types — pick deliberately:
| Timeout | Bounds | Use for agents |
|---|---|---|
| Start-To-Close | One attempt's wall-clock time | Primary — caps a single agent run |
| Schedule-To-Close | Total time across all retries | Hard budget ceiling for the whole job |
| Heartbeat | Max gap between heartbeats | Detects a hung agent mid-loop |
| Schedule-To-Start | Time waiting in the task queue | Detects worker starvation |
These four types are defined in Temporal's activity-timeout docs. A RetryPolicy controls backoff; its fields are Initial Interval (default 1s), Backoff Coefficient (default 2.0), Maximum Interval (default 100× initial), Maximum Attempts (default unlimited), and Non-Retryable Error Types.
from datetime import timedelta
from temporalio.common import RetryPolicy
result = await workflow.execute_activity(
analyse_document, req,
start_to_close_timeout=timedelta(minutes=10),
schedule_to_close_timeout=timedelta(minutes=45), # hard budget ceiling
retry_policy=RetryPolicy(
initial_interval=timedelta(seconds=2),
backoff_coefficient=2.0,
maximum_interval=timedelta(seconds=60),
maximum_attempts=5,
non_retryable_error_types=["ContentPolicyError", "InvalidSchemaError"],
),
)
Two non-obvious rules for AI work:
- Mark deterministic failures non-retryable. A content-policy refusal or malformed-schema error fails identically every retry — retrying just burns tokens.
- Always cap the job. The default is unlimited retries; set
schedule_to_close_timeoutormaximum_attemptsso an agent stuck against an overloaded model cannot quietly run up a large bill. These are your spend ceiling; Lesson 7 adds circuit breakers and dead-letter queues on top.
Compliance and debugging ask the same question: what did the agent do, and why? The engine answers it for free. Temporal durably persists an Event History for every execution, and that history is your audit log — every activity scheduled, started, retried, completed, or failed is a recorded event with inputs, outputs, and timestamps.
Capture decisions as structured Activity results so they land in history — return a dataclass like AgentDecision(document_id, verdict, reasoning, model, input_tokens, output_tokens) rather than free text. Then inspect any run from the CLI:
temporal workflow show --workflow-id document-42 --output json
You get the full ordered event log: which model was called, how many attempts it took, the tokens consumed, the verdict returned — no separate logging pipeline that can drift from reality.
| Auditor's question | Answered by |
|---|---|
| Did the agent run on this document? | ActivityTaskScheduled event |
| How many retries, and why? | Retry events + failure details |
| What verdict, with what reasoning? | ActivityTaskCompleted payload |
| Was a human involved? | Signal / approval events |
Not every team needs Temporal. Match the engine to the shape of your agent workload.
| Engine | Best fit for agents | Retry primitives |
|---|---|---|
| Temporal | Long-running, stateful, human-in-the-loop | RetryPolicy, four timeouts, heartbeats |
| Prefect | Dynamic Python pipelines that change shape per run | retries, retry_delay_seconds, backoff, jitter |
| Airflow | Scheduled batch / nightly agent jobs | retries, retry_delay, execution_timeout |
Prefect attaches retries to the decorator. The @task and @flow decorators accept retries and retry_delay_seconds, plus an exponential_backoff utility and a retry_jitter_factor to avoid thundering-herd retries against a rate-limited model:
from prefect import task
from prefect.tasks import exponential_backoff
@task(retries=4, retry_delay_seconds=exponential_backoff(backoff_factor=2),
retry_jitter_factor=0.5, timeout_seconds=600)
def analyse_document(text: str) -> str: # body calls the Claude Agent SDK
...
Airflow defines retries in default_args — retries, retry_delay, execution_timeout, and retry_exponential_backoff=True to grow the delay — applied to each operator in a scheduled DAG.
Combine the pieces into the full pattern: deterministic validation, an AI analysis step, and a human approval gate. The deterministic steps protect the expensive agent step; the human gate enforces the guardrail philosophy from the Agentic AI course. A Worker polls the task queue and hosts this code — its max_concurrent_activities argument is your real cap against the model API, so set it deliberately to avoid tripping rate limits across every worker at once.
@workflow.defn
class ContractReview:
def __init__(self) -> None:
self._approved: bool | None = None
@workflow.signal
def submit_decision(self, approved: bool) -> None:
self._approved = approved # delivered by a human reviewer
@workflow.run
async def run(self, req: AnalysisRequest) -> str:
if not req.text or len(req.text) < 50: # 1. deterministic gate
return "rejected: document too short"
decision = await workflow.execute_activity( # 2. AI analysis Activity
analyse_document, req,
start_to_close_timeout=timedelta(minutes=10),
retry_policy=RetryPolicy(maximum_attempts=4),
)
if decision.verdict == "approve": # 3. auto-approve low risk
return "auto-approved by agent"
await workflow.execute_activity(notify_reviewer, decision,
start_to_close_timeout=timedelta(minutes=1))
await workflow.wait_condition( # 4. durable wait for human
lambda: self._approved is not None)
return "approved by human" if self._approved else "rejected by human"
The wait_condition on a signal is what durable execution buys you: the workflow pauses for days awaiting a human (a reviewer's UI fires temporal workflow signal --name submit_decision --input 'true'), the worker can be redeployed, and when the signal arrives the execution resumes exactly where it left off — no polling loop, no lost state.
if statement. Agents for judgement, code for rules.Questions & Answers
start_to_close_timeout and emit heartbeats with a much shorter heartbeat_timeout (say, 2 minutes). The long timeout permits the work; the heartbeat detects a genuinely dead worker quickly instead of waiting 30 minutes. For very long jobs, put checkpoint detail in the heartbeat so a retry can resume rather than restart.Key Takeaways
- Draw the determinism boundary. Deterministic rules and routing live in workflow code; anything that calls a model, network, or clock is an Activity. Agents are always Activities.
- Durable execution makes agents reliable. The engine runs the workflow exactly once across crashes and deploys, and replay never re-bills completed agent calls.
- Timeouts and retry policies are cost controls. Cap agent steps with a start-to-close timeout plus a maximum-attempts ceiling, and mark deterministic failures as non-retryable.
- Event history is a free audit trail. Structured decisions from agent Activities give you every attempt, token count, and verdict — queryable and replayable, no separate logging stack.
- Pick the engine to fit the workload. Airflow for scheduled batch, Prefect for dynamic Python pipelines, Temporal for long-lived, stateful, human-in-the-loop orchestration.
- Use agents for judgement, code for rules. A hybrid validate-analyse-approve workflow is cheaper, safer, and more observable than an all-agent pipeline.
Next Steps: Lesson 7: Error Handling & Recovery