AWS Step Functions
Pause a Step Functions execution while Convoy processes your batch, then let Convoy resume it directly — no webhook endpoint, no polling loop. Your state machine waits at a .waitForTaskToken Task state (zero compute cost, up to a year), and Convoy calls states:SendTaskSuccess / SendTaskFailure when the cargo finishes.
This is the Step Functions analog of Direct AWS Durable Callbacks. Both use the same AWS connection (IAM role + ExternalId) — a connection set up for one can be extended to the other by adding one IAM statement.
How it works
Step Functions execution
│
├─ Task state (Resource: arn:aws:states:::lambda:invoke.waitForTaskToken)
│ invokes a small "submit" Lambda with $$.Task.Token in the payload
│
├─ submit Lambda → POST /cargo/load
│ callback: { type: "aws_sfn_task_token",
│ connection_id: "awsconn_…",
│ task_token: "<$$.Task.Token>" }
│
├─ …execution suspended, zero compute cost, up to 1 year…
│
└─ Convoy finishes the cargo →
AssumeRole(your role, ExternalId) →
states:SendTaskSuccess(taskToken, output=<result envelope>)
(or SendTaskFailure on cargo failure)
→ execution resumes with the envelope as the state's output.waitForTaskToken works with any task type that can carry the token (Lambda invoke, SQS, ECS, …); the examples here use Lambda invoke since that’s the common case.
Setup
Register (or reuse) an AWS connection
If you already have a connection for durable callbacks, reuse it — skip to the next step. Otherwise, register one in the dashboard: open your project’s AWS tab, launch the Connect AWS wizard, and choose the Step Functions task token capability (see the full wizard walkthrough in the AWS connections guide).
Save the external_id (cnvyext_…) the wizard shows — it is shown exactly
once — along with the connection_id and the connection-linked API key.
Grant the SendTask actions
Add this statement to your connection role’s inline policy (or deploy the CloudFormation template with EnableSfnTaskCallbacks=true):
{
"Effect": "Allow",
"Action": [
"states:SendTaskSuccess",
"states:SendTaskFailure",
"states:SendTaskHeartbeat"
],
"Resource": "*"
}SendTask* cannot be resource-scoped to a state machine — the TaskToken itself is the capability. Possession of a valid token is required to affect any execution, which is why Resource: "*" here is standard for every SFN task-token integration.
Verify the connection
Once the role exists, open the connection in the dashboard and run the Step Functions task token capability check (Verify).
Convoy assumes your role and runs a dry-run probe (a SendTaskHeartbeat against a fake token — the expected InvalidToken error proves the permission grant without touching any real execution). You can verify the durable capability at the same time if the role also carries the callback grant.
The ARN-addressed capabilities (lambda_invoke, eventbridge, sqs, kinesis, sfn_start, s3) need a concrete target ARN for their probes — the dashboard’s verify dialog collects it. For example, verifying the sfn_start capability probes against your state machine’s ARN.
Add the wait state to your state machine
"WaitForConvoy": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke.waitForTaskToken",
"Parameters": {
"FunctionName": "submit-to-convoy",
"Payload": {
"task_token.$": "$$.Task.Token",
"prompt.$": "$.prompt"
}
},
"TimeoutSeconds": 86400,
"Catch": [{ "ErrorEquals": ["ConvoyCargoFailed"], "Next": "HandleFailure" }],
"Next": "UseResult"
}TimeoutSecondsbounds how long you’re willing to wait for the cargo.- On cargo failure Convoy calls
SendTaskFailurewitherror: "ConvoyCargoFailed"— yourCatchclause matches on it; thecausefield carries the error detail.
Submit from the submit Lambda
import json
import os
import urllib.request
def handler(event, context):
body = {
"params": {
"model": "claude-haiku-4-5",
"max_tokens": 1024,
"messages": [{"role": "user", "content": event["prompt"]}],
},
"callback": {
"type": "aws_sfn_task_token",
"connection_id": "awsconn_YOUR_CONNECTION",
"task_token": event["task_token"],
},
}
req = urllib.request.Request(
"https://api.cnvy.ai/cargo/load",
data=json.dumps(body).encode(),
headers={
"Content-Type": "application/json",
"X-API-Key": os.environ["CONVOY_API_KEY"],
},
)
with urllib.request.urlopen(req) as resp:
return json.loads(resp.read())The task_token is treated as an opaque capability token: encrypted at rest, never logged.
The result envelope
Your Task state resumes with this JSON as its output:
{
"cargo_id": "cargo_abc123",
"success": true,
"response": { "...full model response..." },
"response_truncated": false,
"result_url": "https://api.cnvy.ai/cargo/cargo_abc123/result",
"metadata": { "your": "correlation-values" }
}Step Functions caps task output at 256 KB. Convoy inlines responses up to ~240 KB; larger responses arrive with response: null, response_truncated: true, and the full payload retrievable from result_url (the hosted results mailbox).
On cargo failure, SendTaskFailure fires instead — the execution takes your Catch path with:
Error:ConvoyCargoFailedCause: the error detail (up to 32 768 chars)
Heartbeats (optional)
For long waits, set HeartbeatSeconds on the Task state and opt into Convoy heartbeats:
"callback": {
"type": "aws_sfn_task_token",
"connection_id": "awsconn_…",
"task_token": "…",
"heartbeat": true
}While the cargo is in flight, Convoy calls states:SendTaskHeartbeat on a ~5-minute cadence — a lost cargo trips your HeartbeatSeconds early instead of waiting out the full TimeoutSeconds.
Only set HeartbeatSeconds on the Task state when you ALSO set heartbeat: true in the callback spec — otherwise the state times out waiting for heartbeats Convoy isn’t sending. Use a HeartbeatSeconds comfortably above Convoy’s ~5-minute cadence (600 is a good default).
Delivery semantics
- Exactly-once effective. The TaskToken is consumed on first success; a redelivery attempt gets
TaskDoesNotExistand stops. Your execution can never be resumed twice. - Retries. Transient AWS errors (throttles, 5xx) are retried with exponential backoff (5 attempts — 1, 3, 9, then 27 minutes between attempts, roughly 40 minutes end to end). Terminal conditions (task timed out, token consumed, permission revoked) stop immediately.
- Revocation. Delete the IAM role or the Convoy connection at any time — in-flight deliveries fail with
assume_role_deniedand the connection is flagged, with an email to your org.
Starting a NEW execution per result (aws_sfn_start)
Everything above resumes a paused workflow. The aws_sfn_start
callback type instead starts a new state machine execution with the
result envelope as its input — “result kicks off a workflow” rather than
“result resumes a workflow”:
{
"callback": {
"type": "aws_sfn_start",
"connection_id": "awsconn_…",
"state_machine_arn": "arn:aws:states:us-east-1:123456789012:stateMachine:handle-cargo-result"
}
}IAM grant (resource-scoped — never *; the CloudFormation template’s
StateMachineArns parameter does the same):
{
"Effect": "Allow",
"Action": "states:StartExecution",
"Resource": "arn:aws:states:us-east-1:123456789012:stateMachine:handle-cargo-result"
}Idempotent by construction (Standard workflows): Convoy names each
execution convoy-{cargo_id} — deterministic — so a redelivery after a
lost response hits ExecutionAlreadyExists, which Convoy treats as
success. Standard workflows therefore get exactly-once semantics.
Express state machines don’t dedupe on name — they get at-least-once
(dedupe on the input’s cargo_id, same contract as Lambda invoke).
The optional sfn_start verify probe (run from the dashboard’s capability
check with your state machine’s ARN) starts a real marked execution
(name=convoy-verify-<nonce>, input={"convoy_verification": true}).
Either add a one-line Choice state at the top of your state machine that
no-ops on that input, or skip the probe entirely (probes are opt-in;
rate-capped 1/minute per connection).
Troubleshooting
| Symptom | Likely cause | Fix |
|---|---|---|
Verify returns sfn_token: failed | states:SendTask* not granted | Add the IAM statement above and re-verify |
| Execution times out at the wait state | Cargo still processing, or TimeoutSeconds too low | Batch processing takes minutes–hours; size TimeoutSeconds accordingly |
Execution fails with States.Timeout before cargo finishes | HeartbeatSeconds set without heartbeat: true | Enable heartbeats in the callback spec, or remove HeartbeatSeconds |
Delivery fails callback_expired | The task timed out (or execution was aborted) before delivery | Increase TimeoutSeconds / investigate the abort |
Delivery fails callback_not_found | Token malformed at submit, or already consumed | Pass $$.Task.Token through verbatim; don’t reuse tokens |