Workflow Engines

50 min intermediate Lesson 6

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:

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 prefect and pip install apache-airflow for the comparison
  • An ANTHROPIC_API_KEY exported in your shell

1 The Determinism Boundary

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.

NOTE
Key Insight
If a step produces the same output for the same input every time, it is workflow code. If it talks to a model, a network, or a clock, it is an Activity. Agents are always Activities.
WARNING
Watch Out
Workflow code is replayed, so never call an agent, generate a UUID, read the clock, or do I/O directly inside it — replay would diverge and corrupt the execution. Push all of that into Activities.

2 Wrapping an Agent Call as an Activity

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.

TIP
Tip
Keep the Activity boundary at one logical agent job, not one model call. Let the SDK run its full tool loop inside the Activity; the engine then retries the whole job, far easier to reason about than retrying individual tool calls.

3 Timeouts and Retry Policies for AI Steps

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:

  1. Mark deterministic failures non-retryable. A content-policy refusal or malformed-schema error fails identically every retry — retrying just burns tokens.
  2. Always cap the job. The default is unlimited retries; set schedule_to_close_timeout or maximum_attempts so 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.
WARNING
Watch Out
Retries are only safe if the Activity is idempotent. If the agent writes to a database mid-loop, a retry double-writes. Make the side effect idempotent with a key derived from the input, or split the write into its own Activity after the agent returns.

4 Audit Trails from Event History

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
TIP
Tip
The richer and more structured the result payload from each agent Activity, the more your event history doubles as a queryable audit database — and the less you need a bolt-on logging system.

5 Choosing an Engine: Temporal vs. Prefect vs. Airflow

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.

NOTE
Decision heuristic
Airflow when the agent job is scheduled batch. Prefect when the pipeline is dynamic Python with minimal ceremony. Temporal when runs are long-lived, stateful, must survive deploys, or wait on humans — which describes most production multi-agent systems.

6 A Hybrid Workflow: Validate → Analyse → Approve

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.

TIP
Tip
Only step 2 is an agent. Validation and routing stay as deterministic code because they must behave identically every run; spending a model call on 'is this string longer than 50 chars' is slower, costlier, and less reliable than an if statement. Agents for judgement, code for rules.
NOTE
Pattern connection
This is the supervisor pattern from Lesson 2, but the supervisor is deterministic workflow code rather than another LLM. In production that trade is often right: predictable routing and a free audit trail, while agents do the genuinely fuzzy work.

Questions & Answers

Q: If the engine already retries on failure, why do I still need guardrails inside the agent?
Retries only help with transient failures. The engine cannot tell that an agent produced a confidently wrong answer, leaked PII, or took an unsafe action — it sees a successful Activity. Engine-level retries are the operational safety net; the agent guardrails from the Agentic AI course are the semantic one. You need both, and the workflow is where you wire semantic checks in as their own deterministic Activities.
Q: Won't replay re-run my agent and double my token bill?
No. Replay re-executes workflow code but not Activities that already completed — their recorded results are returned from event history. As long as the agent call lives inside an Activity (never in workflow code), a crash and replay costs zero extra tokens. That is why the determinism boundary in Step 1 is non-negotiable.
Q: How do I handle an agent step that legitimately takes 30 minutes?
Set a generous 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.
Q: My agent Activity writes to a database mid-loop. How do I make retries safe?
Either make the write idempotent — derive a deterministic key from the input and upsert, so a replayed write is a no-op — or split the write out of the agent Activity entirely: let the agent return a structured decision, then persist it in a separate idempotent Activity afterward. The second approach also keeps the audit trail clean, since the decision and the side effect become distinct events.

Key Takeaways

  1. 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.
  2. Durable execution makes agents reliable. The engine runs the workflow exactly once across crashes and deploys, and replay never re-bills completed agent calls.
  3. 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.
  4. 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.
  5. 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.
  6. 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