Skip to Content
IntegrationsAWS Lambda Durable Functions

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} } ──┘ └────────────────────────┘
  1. The durable workflow calls wait_for_callback. The SDK hands your submitter function a one-time callback ID.
  2. The submitter POSTs to /cargo/load with metadata: {"de_callback_id": "<id>"} and suspends.
  3. Convoy batches and processes the request (minutes to hours).
  4. Convoy POSTs the result to your callback_url, echoing your metadata back.
  5. The receiver extracts metadata.de_callback_id and 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 the 2025-12-01 Lambda API and first shipped in boto3 1.42.1. The boto3 bundled in the Lambda runtime may be older, so bundle your own via requirements.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_callback submitter — 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 from GET /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.Alias
sam build sam deploy --guided --parameter-overrides ConvoyApiKeySecretName=convoy/api-key

AuthType: 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.json

The 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

ConcernGuidance
Result sizeSendDurableExecutionCallbackSuccess 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 timeoutAlways set one. Convoy’s completion window is up to 24h — 26h is a sensible wait, with the durable ExecutionTimeout above it (48h)
Callback ID hygieneIt’s a capability token: pass it only in metadata, never log it, never put it in URLs you might log
Replay safetyConvoy may retry webhook deliveries. The Lambda callback API consumes each callback ID exactly once — duplicate deliveries get ResourceNotFoundException, which the receiver acks harmlessly
Metadata limits50 entries, 128-char keys, 2,048-char values, 4 KB total — see metadata limits. A callback ID fits comfortably
Treat the envelope as untrustedThe 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:alias or function:version), never $LATEST

The workflow resumed with oversize: true


Next Steps

Last updated on