AWS Step Functions and an agent framework both manage control flow, which makes it easy to build the same thing twice. Step Functions is a durable state machine: it records every transition, retries failed steps with backoff, waits up to a year for a callback, and shows an auditor exactly what happened. ADK for Java runs a language model with tools inside a conversation. The integration works when each does what it is good at. The state machine owns the business process and its pauses. The agent owns a bounded judgment step inside one task state, such as "classify this claim and draft a response", and returns a small, validated result.
This article covers that design end to end: choosing Standard or Express, running an ADK Runner inside a Lambda handler and on ECS for longer runs, the state machine definition with retries and catches, human approval through task tokens, the 256 KiB payload limit, and why Step Functions retries make tool idempotency mandatory. Service limits and behaviours were checked against the AWS Step Functions developer guide in October 2026. ADK class and method names are limited to ones verified on this site's other ADK Java pages. Confirm both against the versions you deploy.
Who owns what: the process versus the judgment
Put a step in the state machine when it needs durability, a timeout, a retry policy, a wait measured in hours or days, or an audit trail. Put it in the agent when it needs judgment over unstructured input. A loan workflow is a good example: collect documents, extract fields, check policy, approve, disburse. Extraction and the policy check may be agent steps. The sequencing, the human approval and the payment are workflow steps.
The rule that prevents most bugs is one owner per pause. ADK has its own pause mechanism, long-running tools with resumability, described in Task Resumability in ADK Java. If the workflow waits on a task token while the agent also waits on a long-running tool, there are two clocks, two places to resume from and two retry policies. Choose one. In this design the workflow owns every wait longer than a single agent turn, and the agent's tools return immediately.
Standard or Express
| Property (AWS docs, Oct 2026) | Standard | Express |
|---|---|---|
| Maximum duration | 1 year | 5 minutes |
| Execution semantics | Exactly-once (except your own Retry) | Async: at-least-once; sync: at-most-once |
| Integration patterns | Request Response, .sync, .waitForTaskToken | Request Response only |
| History | 25,000 events per execution; kept 90 days | CloudWatch Logs only, if enabled |
| Start by name | Idempotent for a running execution with the same name | Not managed |
| Redrive | Within 14 days | Not supported |
| Payload per state | 256 KiB | 256 KiB |
| Billing | Per state transition | Per execution, duration and memory |
The row that decides most agent designs is the third one. Express workflows support only Request Response integrations, so any human approval (.waitForTaskToken) or long ECS run (.sync) forces Standard. Express suits high-volume, short pipelines where each item is one agent call: classify a ticket, tag a document, all within five minutes. Because asynchronous Express is at-least-once, an item can be processed twice, so the idempotency rules below apply even more strictly.
The agent as a Lambda task
Build the Runner once per execution environment, in a static field, so warm invocations reuse it. Create a fresh session for each invocation: the workflow's state is the durable record, and the session lives only for one agent turn, so the in-memory session service is enough. The handler reads its input by pointer, runs one turn, demands JSON, validates it, and throws a typed exception when the agent misbehaves, so the state machine can route on the error name.
package com.example.agent;
import com.amazonaws.services.lambda.runtime.Context;
import com.amazonaws.services.lambda.runtime.RequestHandler;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.adk.runner.Runner;
import com.google.adk.sessions.InMemorySessionService;
import com.google.adk.sessions.Session;
import com.google.genai.types.Content;
import com.google.genai.types.Part;
import java.util.Map;
public final class TriageHandler implements RequestHandler<Map<String, Object>, Map<String, Object>> {
private static final String APP = "claims_triage";
private static final InMemorySessionService SESSIONS = new InMemorySessionService();
private static final Runner RUNNER = Runner.builder()
.agent(TriageAgent.ROOT_AGENT).appName(APP).sessionService(SESSIONS).build();
private static final ObjectMapper JSON = new ObjectMapper();
@Override
public Map<String, Object> handleRequest(Map<String, Object> in, Context ctx) {
String key = (String) in.get("idempotencyKey"); // execution name + state name
String claimText = ClaimStore.read((String) in.get("claimS3Uri"));
Session s = SESSIONS.createSession(APP, "workflow",
Map.<String, Object>of("idempotencyKey", key), null).blockingGet();
StringBuilder out = new StringBuilder();
try {
RUNNER.runAsync("workflow", s.id(), Content.fromParts(Part.fromText(claimText)))
.blockingForEach(ev -> { if (ev.finalResponse()) out.append(ev.stringifyContent()); });
} finally {
SESSIONS.deleteSession(APP, "workflow", s.id()).blockingAwait(); // warm containers would accumulate them
}
JsonNode r;
try { r = JSON.readTree(out.toString()); }
catch (Exception e) { throw new BadAgentOutput("not JSON"); }
if (!r.hasNonNull("category") || !r.hasNonNull("confidence")) throw new BadAgentOutput("missing fields");
String draftUri = ClaimStore.writeDraft(key, r.path("draft").asText("")); // keyed: retry-safe
return Map.of("category", r.get("category").asText(),
"confidence", r.get("confidence").asDouble(),
"draftS3Uri", draftUri);
}
}
/** Unhandled; Step Functions sees the exception class name as the error name. */
final class BadAgentOutput extends RuntimeException {
BadAgentOutput(String m) { super(m); }
}ClaimStore and TriageAgent are your code. The agent is an LlmAgent whose instruction asks for a JSON object with category, confidence and draft. The returned map stays far below 256 KiB because the draft goes to S3 and only its URI travels through the workflow. Java cold starts are worth managing for this handler: initialising the JVM, the ADK runner and the model client can dominate a short turn. The options, including SnapStart and provisioned concurrency, are compared in Lambda cold start architecture.
The state machine
The definition below uses JSONPath. The optimised lambda:invoke integration wraps the function output in a metadata object, with the result in Payload, so ResultSelector extracts it. The idempotency key is built from the execution name and the state name with the States.Format intrinsic. Both values are the same on every retry and every redrive of this state, which is exactly what the tools need.
{
"StartAt": "Triage",
"States": {
"Triage": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": {
"FunctionName": "arn:aws:lambda:REGION:ACCOUNT:function:claims-triage:live",
"Payload": {
"claimS3Uri.$": "$.claimS3Uri",
"idempotencyKey.$": "States.Format('{}:{}', $$.Execution.Name, $$.State.Name)"
}
},
"ResultSelector": { "result.$": "$.Payload" },
"ResultPath": "$.triage",
"TimeoutSeconds": 120,
"Retry": [
{ "ErrorEquals": ["Lambda.ServiceException", "Lambda.SdkClientException",
"Lambda.TooManyRequestsException"],
"IntervalSeconds": 2, "MaxAttempts": 4, "BackoffRate": 2.0,
"MaxDelaySeconds": 30, "JitterStrategy": "FULL" },
{ "ErrorEquals": ["com.example.agent.BadAgentOutput"], "MaxAttempts": 1 }
],
"Catch": [ { "ErrorEquals": ["States.ALL"], "ResultPath": "$.error", "Next": "HumanQueue" } ],
"Next": "Confident?"
},
"Confident?": {
"Type": "Choice",
"Choices": [ { "Variable": "$.triage.result.confidence", "NumericGreaterThanEquals": 0.85, "Next": "Act" } ],
"Default": "Approve"
},
"Approve": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke.waitForTaskToken",
"Parameters": {
"FunctionName": "arn:aws:lambda:REGION:ACCOUNT:function:request-review:live",
"Payload": { "taskToken.$": "$$.Task.Token", "draftS3Uri.$": "$.triage.result.draftS3Uri" }
},
"TimeoutSeconds": 604800,
"ResultPath": "$.review",
"Catch": [ { "ErrorEquals": ["States.ALL"], "ResultPath": "$.error", "Next": "HumanQueue" } ],
"Next": "Act"
},
"Act": { "Type": "Task", "Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "arn:aws:lambda:REGION:ACCOUNT:function:claims-act:live",
"Payload.$": "$" }, "End": true },
"HumanQueue": { "Type": "Task", "Resource": "arn:aws:states:::sqs:sendMessage",
"Parameters": { "QueueUrl": "https://sqs.REGION.amazonaws.com/ACCOUNT/claims-manual",
"MessageBody.$": "$" }, "End": true }
}
}Two retry details matter for agents. Transient Lambda service errors get exponential backoff with full jitter. Invalid agent output gets exactly one retry, because a second sample from the model often succeeds, but a loop of retries only burns tokens. Everything else is caught and routed to people rather than failing the execution, since an execution-level failure cannot be caught. A real definition should also set the Act state's own retries and its own idempotency key.
Human approval with task tokens
The Approve state passes its task token to a Lambda that records the token with the draft, notifies a reviewer and returns. The execution then waits. When the reviewer decides, your review service completes the task with the AWS SDK for Java v2:
import software.amazon.awssdk.services.sfn.SfnClient;
import software.amazon.awssdk.services.sfn.model.SendTaskFailureRequest;
import software.amazon.awssdk.services.sfn.model.SendTaskSuccessRequest;
import software.amazon.awssdk.services.sfn.model.TaskTimedOutException;
void complete(SfnClient sfn, String token, boolean approved, String reviewer) {
try {
if (approved) {
sfn.sendTaskSuccess(SendTaskSuccessRequest.builder().taskToken(token)
.output("{\"approved\":true,\"reviewer\":\"" + reviewer + "\"}").build());
} else {
sfn.sendTaskFailure(SendTaskFailureRequest.builder().taskToken(token)
.error("Rejected").cause("rejected by " + reviewer).build());
}
} catch (TaskTimedOutException e) {
// The task already timed out; the token is dead. Tell the reviewer, do not retry.
}
}Several behaviours of task tokens shape this code, according to the AWS documentation. A token must be returned by a principal in the same AWS account. If the task times out, Step Functions generates a new token, so a stored token can go stale. That is why Approve sets only TimeoutSeconds (a week): a human reviewer sends no heartbeats. HeartbeatSeconds belongs on tasks whose worker can call SendTaskHeartbeat, such as an ECS container holding a token. For a Lambda callback task, the heartbeat clock starts only after the function returns, and a silent wait fails with States.Timeout. In a real service, build the output JSON with a serializer rather than string concatenation.
Long agent runs and the history limit
A Lambda function can run for at most 15 minutes, and a single agent turn should take far less. Agent work that needs many tool calls over a longer period has two good homes. The first is ECS on Fargate, started with arn:aws:states:::ecs:runTask.sync: the workflow waits for the container to exit, and the container runs the same Runner code with the input pointer passed as an environment override. The second is to split the agent into several states, each one turn, with the transcript in S3 between them. Splitting costs more state transitions but gives each step its own timeout, retry policy and history entry.
Keep an eye on the 25,000-event history limit. A workflow that loops one state per agent step can reach it on long investigations. When a loop may run long, start a child execution per batch of iterations. Inputs and outputs above 256 KiB fail with States.DataLimitExceeded, which States.ALL does not catch. Keep transcripts and documents in S3 and pass only URIs.
Retries make idempotency mandatory
Standard workflows run each state exactly once, unless you configure retries, and agents need retries. A retry of the Triage state re-runs the whole agent turn, including any tool calls that completed before the failure. Redriving a failed execution does the same. So every tool with a side effect must be idempotent, keyed on something stable across attempts. Execution name plus state name is stable, while a UUID generated inside the handler is not. The tool reads the key from session state, where the handler put it:
public static Map<String, Object> createRefundTicket(
@Schema(name = "claimId") String claimId, ToolContext ctx) {
String key = "refund:" + ctx.state().get("idempotencyKey") + ":" + claimId;
return Tickets.findByKey(key) // durable lookup first
.map(t -> Map.<String, Object>of("ticket", t.id(), "replayed", true))
.orElseGet(() -> Map.of("ticket", Tickets.create(claimId, key).id(), "replayed", false));
}Lookup-then-create is still racy if two attempts overlap. Enforce the key with a unique constraint or a conditional write in the store, as described in the idempotency architecture. Start executions with a business-derived name, such as the claim id, so a duplicate trigger hits the idempotent StartExecution behaviour instead of launching a second workflow.
Failure modes
- Express with callbacks. A definition using
.waitForTaskTokencannot run as Express. Choose the type first, because it cannot be changed after creation. - Transcripts in the payload. A long conversation passes 256 KiB and fails with an error that
States.ALLdoes not catch. - Double side effects. A retry after a tool succeeded charges or emails twice. Fix it with keys derived from the execution name and state name.
- Two owners of a pause. An ADK long-running tool waiting inside a task-token state. Pick one mechanism.
- Stale tokens. A reviewer approves after the timeout and the call fails with
TaskTimedOut. Handle it and tell the reviewer. - Unvalidated agent output. A Choice state on a missing field fails at runtime. Validate in the handler and throw a typed error.
- History exhaustion. Per-step loops in one execution reach 25,000 events. Use child executions.
Trade-offs
| Option | Choose it when | Watch out for |
|---|---|---|
| Agent turn in Lambda | Turns finish in seconds to minutes | Cold starts; 15-minute ceiling |
| Agent on ECS with .sync | Long multi-tool investigations | Slower start; you own the container image |
| One state per agent step | You want per-step retries and audit | More transitions; history limit |
| Express workflow | High-volume single calls, no waits | At-least-once (async); no callbacks |
| ADK resumability instead | The agent, not the business process, owns the pause | Pause is outside the workflow's audit trail |
What to do next
- List the waits in your process. If any are human or longer than five minutes, use Standard.
- Draw the boundary: the workflow owns sequencing, waits and payments, and the agent owns judgment inside one state.
- Build the Runner once per container and create a fresh session per invocation. Return validated JSON with S3 pointers.
- Derive an idempotency key from the execution name and state name, and enforce it in every side-effecting tool.
- Add Lambda service-error retries with jitter, one retry for bad output, and a catch-all route to a human queue.
- Implement the review service with SendTaskSuccess and SendTaskFailure, handle timed-out tokens, set TimeoutSeconds on human waits, and use HeartbeatSeconds only where a worker heartbeats.
- Load-test a retry: kill the handler after a tool call and confirm nothing happens twice. For a non-AWS comparison, read Deploying ADK Java to Cloud Run and AWS Step Functions architecture.