AWS Lambda Durable Functions
Resume an AWS Lambda durable function from a Convoy callback — your workflow submits a batch, suspends at zero compute cost for minutes or hours, and picks up exactly where it left off when the result arrives.
Testing tip — Build and verify this integration end-to-end with convoy-mock first. It’s a free, synthetic model that returns a callback in ~60 seconds and is never billed. Once your wiring works, swap "model": "convoy-mock" for a production model.
Looking for plain (non-durable) Lambda patterns — submit + receive with DynamoDB, SQS, or S3 triggers? See the AWS Lambda guide. This page is specifically about waitForCallback-style durable workflows.
Want to skip the receiver Lambda entirely? With Direct AWS Durable Callbacks, Convoy assumes a narrowly-scoped IAM role in your account and calls SendDurableExecutionCallbackSuccess itself — no webhook endpoint, no receiver, no signing. This page covers the webhook + receiver variant, which works without any IAM setup.
Why durable functions + Convoy
Convoy is a submit now, results arrive later batch API. Lambda durable functions are built for exactly this shape of problem:
- No polling loops. Your workflow calls
wait_for_callback(...)and suspends. You pay nothing while Convoy’s batch is in flight. - No lost state. Everything before the suspension is checkpointed. When the callback arrives the workflow resumes with the result in hand.
- Multi-step pipelines stay in one function. Submit → wait → post-process → submit again → wait — all as ordinary code, no state machine JSON.
The glue is Convoy’s metadata field: you pass the durable execution’s callback ID through metadata when you load the cargo, and Convoy echoes it back verbatim in the webhook payload. A tiny receiver Lambda extracts it and calls SendDurableExecutionCallbackSuccess to resume your workflow.
Architecture
┌────────────────────────┐ 1. wait_for_callback(submitter)
│ Durable workflow │─────────────────────────────────┐
│ Lambda │ │
│ (suspended, $0) │ 2. POST /cargo/load ▼
└───────────▲────────────┘ metadata: {de_callback_id} ┌─────────────┐
│ ────────────────────────────▶│ Convoy API │
│ 5. SendDurableExecution └──────┬──────┘
│ CallbackSuccess(callback_id, result) │
│ │ 3. batch
┌───────────┴────────────┐ 4. webhook POST │ processing
│ Receiver Lambda │ { cargo_id, success, response, │
│ (Function URL) │◀────── metadata: {de_callback_id} } ──┘
└────────────────────────┘- The durable workflow calls
wait_for_callback. The SDK hands your submitter function a one-time callback ID. - The submitter POSTs to
/cargo/loadwithmetadata: {"de_callback_id": "<id>"}and suspends. - Convoy batches and processes the request (minutes to hours).
- Convoy POSTs the result to your
callback_url, echoing yourmetadataback. - The receiver extracts
metadata.de_callback_idand calls the Lambda callback API — your workflow resumes with the result.
Prerequisites
- An AWS region where Lambda durable functions are available
- boto3 1.42.1 or newer — the
SendDurableExecutionCallback*APIs are on the2025-12-01Lambda API and first shipped in boto3 1.42.1. The boto3 bundled in the Lambda runtime may be older, so bundle your own viarequirements.txt - The AWS Durable Execution SDK for Python (
aws-durable-execution-sdk-python) - A Convoy API key (
convoy_sk_...) stored in AWS Secrets Manager
The durable workflow function
The key rules for durable functions apply here:
- All I/O happens inside a step or the
wait_for_callbacksubmitter — both are checkpointed and never re-run on replay. - Always set a timeout on
wait_for_callback. Convoy’s completion window is up to 24 hours; give the wait comfortable headroom (26h below) so slow batches don’t strand your execution.
# src/workflow/handler.py
import json
import os
import urllib.request
import boto3
from aws_durable_execution_sdk_python import DurableContext, durable_execution
from aws_durable_execution_sdk_python.config import Duration, WaitForCallbackConfig
from aws_durable_execution_sdk_python.exceptions import CallbackError
CONVOY_API_BASE = os.environ.get("CONVOY_API_BASE", "https://api.cnvy.ai")
API_KEY_SECRET_NAME = os.environ["CONVOY_API_KEY_SECRET_NAME"]
RECEIVER_URL = os.environ["RECEIVER_URL"]
# Convoy's completion window is up to 24h; leave headroom for retries.
CALLBACK_TIMEOUT_HOURS = 26
def _get_convoy_api_key() -> str:
sm = boto3.client("secretsmanager")
return sm.get_secret_value(SecretId=API_KEY_SECRET_NAME)["SecretString"].strip()
def _post_json(url: str, body: dict, api_key: str) -> dict:
req = urllib.request.Request(
url,
data=json.dumps(body).encode(),
headers={"Content-Type": "application/json", "X-API-Key": api_key},
method="POST",
)
with urllib.request.urlopen(req, timeout=30) as resp:
return json.loads(resp.read().decode())
@durable_execution
def handler(event: dict, context: DurableContext) -> dict:
prompt = event["prompt"]
def submit_cargo(callback_id: str, ctx) -> dict:
"""Runs exactly once (checkpointed). Safe place for I/O."""
api_key = _get_convoy_api_key()
return _post_json(
f"{CONVOY_API_BASE}/cargo/load",
{
"params": {
"model": event.get("model", "claude-3-haiku"),
"max_tokens": int(event.get("max_tokens", 1024)),
"messages": [{"role": "user", "content": prompt}],
},
"callback_url": RECEIVER_URL,
# The durable callback ID rides through Convoy untouched
# and comes back in the webhook payload.
"metadata": {"de_callback_id": callback_id},
},
api_key,
)
try:
# Suspends here (zero compute) until the receiver calls
# SendDurableExecutionCallbackSuccess/Failure with our callback_id.
result = context.wait_for_callback(
submitter=submit_cargo,
name="wait-convoy-result",
config=WaitForCallbackConfig(
timeout=Duration.from_hours(CALLBACK_TIMEOUT_HOURS),
),
)
except CallbackError as err:
if getattr(err, "error_type", "") == "Timeout":
return {"status": "convoy_timeout"}
return {"status": "convoy_failed", "error": str(err)}
envelope = json.loads(result) if isinstance(result, (str, bytes)) else result
def post_process(_ctx) -> dict:
# Your real business logic here (checkpointed step).
response = envelope.get("response") or {}
text = "".join(
block.get("text", "")
for block in response.get("content") or []
if isinstance(block, dict) and block.get("type") == "text"
)
return {"cargo_id": envelope.get("cargo_id"), "preview": text[:200]}
processed = context.step(post_process, name="post-process")
return {"status": "done", "output": processed}The callback ID is a capability token — anyone holding it can resume (or fail) your execution. Pass it only through metadata, and never log it. Convoy stores metadata verbatim, never logs values, and echoes them only to your callback_url.
The receiver function
A plain Lambda behind a Function URL. It receives Convoy’s webhook, pulls the callback ID out of metadata, and resumes the workflow:
# src/receiver/handler.py
import base64
import json
import logging
import boto3
logger = logging.getLogger()
logger.setLevel(logging.INFO)
lambda_client = boto3.client("lambda")
# AWS caps the durable callback Result at 256 KB; leave envelope headroom.
MAX_INLINE_RESULT_BYTES = 240_000
def _extract_body(event: dict) -> dict:
raw = event.get("body") or ""
if event.get("isBase64Encoded"):
raw = base64.b64decode(raw).decode()
return json.loads(raw) if raw else {}
def _build_result_envelope(payload: dict) -> bytes:
envelope = {
"cargo_id": payload.get("cargo_id"),
"response": payload.get("response"),
"oversize": False,
}
encoded = json.dumps(envelope).encode()
if len(encoded) > MAX_INLINE_RESULT_BYTES:
# Too big for a durable callback Result — send a pointer instead
# and have the resumed workflow fetch the full result via
# GET /cargo/{cargo_id}/result.
envelope = {"cargo_id": payload.get("cargo_id"), "response": None, "oversize": True}
encoded = json.dumps(envelope).encode()
return encoded
def _ok(body: dict) -> dict:
return {"statusCode": 200, "body": json.dumps(body)}
def handler(event: dict, _context) -> dict:
try:
payload = _extract_body(event)
except (ValueError, UnicodeDecodeError):
return {"statusCode": 400, "body": json.dumps({"error": "bad_body"})}
cargo_id = payload.get("cargo_id")
callback_id = (payload.get("metadata") or {}).get("de_callback_id")
if not callback_id:
# Not a durable-workflow cargo. Ack so Convoy does not retry.
return _ok({"ok": True, "skipped": True})
try:
if payload.get("success"):
lambda_client.send_durable_execution_callback_success(
CallbackId=callback_id,
Result=_build_result_envelope(payload),
)
else:
lambda_client.send_durable_execution_callback_failure(
CallbackId=callback_id,
Error={
"ErrorType": "ConvoyCargoFailed",
"ErrorMessage": (payload.get("error_message") or "Cargo processing failed")[:1024],
},
)
except lambda_client.exceptions.CallbackTimeoutException:
# The durable wait already timed out; nothing left to resume.
return _ok({"ok": True, "callback_expired": True})
except lambda_client.exceptions.ResourceNotFoundException:
# Unknown or already-consumed callback id. Retrying never helps.
return _ok({"ok": True, "callback_not_found": True})
except lambda_client.exceptions.TooManyRequestsException:
# Throttled — return 500 so Convoy's retry machinery tries again.
return {"statusCode": 500, "body": json.dumps({"error": "throttled"})}
return _ok({"ok": True})Three behaviors worth copying:
- Ack (200) terminal failures. An expired or already-consumed callback can never succeed on retry — returning 200 stops Convoy’s retry loop from hammering a dead callback.
- Return 500 on throttling so Convoy does retry.
- Respect the 256 KB Result cap. For oversize results, resume with a pointer (
oversize: true) and let the workflow fetch the full result fromGET /cargo/{cargo_id}/result.
Deploy with SAM
# template.yaml
AWSTemplateFormatVersion: '2010-09-09'
Transform: AWS::Serverless-2016-10-31
Parameters:
ConvoyApiKeySecretName:
Type: String
Default: convoy/api-key
Globals:
Function:
Runtime: python3.14
MemorySize: 256
Timeout: 60
Resources:
# Durable execution MUST be enabled at creation (DurableConfig) and the
# function MUST be invoked with a qualified ARN (AutoPublishAlias).
WorkflowFunction:
Type: AWS::Serverless::Function
Properties:
FunctionName: convoy-batch-workflow
CodeUri: src/workflow/
Handler: handler.handler
AutoPublishAlias: live
DurableConfig:
ExecutionTimeout: 172800 # 48h: covers 26h callback wait + slack
RetentionPeriodInDays: 7
Policies:
- arn:aws:iam::aws:policy/service-role/AWSLambdaBasicDurableExecutionRolePolicy
- Statement:
- Sid: ReadConvoyApiKey
Effect: Allow
Action: secretsmanager:GetSecretValue
Resource: !Sub arn:aws:secretsmanager:${AWS::Region}:${AWS::AccountId}:secret:${ConvoyApiKeySecretName}*
Environment:
Variables:
CONVOY_API_KEY_SECRET_NAME: !Ref ConvoyApiKeySecretName
RECEIVER_URL: !GetAtt ReceiverFunctionUrl.FunctionUrl
ReceiverFunction:
Type: AWS::Serverless::Function
Properties:
FunctionName: convoy-webhook-receiver
CodeUri: src/receiver/
Handler: handler.handler
Timeout: 30
Policies:
- Statement:
# Least privilege: only the callback actions, only on the
# workflow function.
- Sid: SendConvoyDurableCallbacks
Effect: Allow
Action:
- lambda:SendDurableExecutionCallbackSuccess
- lambda:SendDurableExecutionCallbackFailure
- lambda:SendDurableExecutionCallbackHeartbeat
Resource: !Sub arn:aws:lambda:${AWS::Region}:${AWS::AccountId}:function:convoy-batch-workflow:*
ReceiverFunctionUrl:
Type: AWS::Lambda::Url
Properties:
TargetFunctionArn: !GetAtt ReceiverFunction.Arn
AuthType: NONE
ReceiverUrlPermission:
Type: AWS::Lambda::Permission
Properties:
FunctionName: !Ref ReceiverFunction
Action: lambda:InvokeFunctionUrl
Principal: '*'
FunctionUrlAuthType: NONE
Outputs:
ReceiverUrl:
Value: !GetAtt ReceiverFunctionUrl.FunctionUrl
WorkflowAliasArn:
Value: !Ref WorkflowFunction.Aliassam build
sam deploy --guided --parameter-overrides ConvoyApiKeySecretName=convoy/api-keyAuthType: NONE keeps the sample simple, but Convoy’s webhook POST is unauthenticated. For production, embed a random token path segment in the callback_url, or put API Gateway with a shared-secret check in front. See webhook security best practices.
Run a workflow
Durable functions must be invoked with a qualified function name, and an idempotent --durable-execution-name prevents accidental duplicate submissions:
aws lambda invoke \
--function-name 'convoy-batch-workflow:live' \
--invocation-type Event \
--durable-execution-name "demo-$(date +%s)" \
--payload '{"prompt": "Summarize the history of containerized shipping in 3 bullets.", "model": "convoy-mock"}' \
--cli-binary-format raw-in-base64-out \
response.jsonThe function submits the cargo and suspends. When Convoy’s batch completes, the webhook hits the receiver, the receiver resumes the execution, and post-process runs. Watch progress in the Lambda console’s Durable Executions tab or CloudWatch Logs.
With "model": "convoy-mock" the entire round trip completes in about a minute — ideal for verifying the wiring before switching to a production model.
Gotchas & limits
| Concern | Guidance |
|---|---|
| Result size | SendDurableExecutionCallbackSuccess caps Result at 256 KB. For large model outputs, resume with a pointer and fetch the full result from GET /cargo/{cargo_id}/result (stored 30 days) |
| Wait timeout | Always set one. Convoy’s completion window is up to 24h — 26h is a sensible wait, with the durable ExecutionTimeout above it (48h) |
| Callback ID hygiene | It’s a capability token: pass it only in metadata, never log it, never put it in URLs you might log |
| Replay safety | Convoy may retry webhook deliveries. The Lambda callback API consumes each callback ID exactly once — duplicate deliveries get ResourceNotFoundException, which the receiver acks harmlessly |
| Metadata limits | 50 entries, 128-char keys, 2,048-char values, 4 KB total — see metadata limits. A callback ID fits comfortably |
| Treat the envelope as untrusted | The webhook body crosses the public internet to your receiver. Validate cargo_id against your own records before acting on the result |
Troubleshooting
The workflow never resumes
- Check the receiver’s CloudWatch Logs — did the webhook arrive? If not, follow the webhook troubleshooting steps
- If the webhook arrived but logged
callback_not_found, the wait may have already timed out or the execution was re-run with a fresh callback ID - Verify the receiver’s IAM policy targets the workflow function’s ARN (including the
:*qualifier suffix)
ValidationException when invoking the workflow
- Durable functions must be invoked via a qualified name (
function:aliasorfunction:version), never$LATEST
The workflow resumed with oversize: true
- The result exceeded the 256 KB callback cap. Fetch the full result in a checkpointed step:
GET /cargo/{cargo_id}/result
Next Steps
- Load Cargo API — full request reference, including
metadata - Webhooks guide — payload format, retries, and security
- AWS Lambda guide — non-durable Lambda patterns (DynamoDB, SQS, S3 triggers)