AWS SQS & Kinesis
Deliver cargo results straight into your queue or your data stream — no webhook endpoint:
- SQS (
aws_sqs): the classic queue-consumer pattern. Results land in a queue your workers poll at their own pace, with SQS’s native visibility timeouts, redrive policies, and DLQs. FIFO queues get exactly-once delivery. - Kinesis (
aws_kinesis): the streaming/analytics pattern. Results flow into pipelines (managed Flink, Firehose, custom consumers, real-time dashboards) with per-shard ordering and replay.
Both use the same AWS connection and the same result envelope as the other direct-AWS targets.
SQS (aws_sqs)
cargo completes ─▶ AssumeRole ─▶ sqs:SendMessage(
QueueUrl, # derived from your ARN
MessageBody=<envelope>,
MessageAttributes={cargo_id, success},
[MessageGroupId, MessageDeduplicationId]) # FIFOGrant sqs:SendMessage on your role
{
"Effect": "Allow",
"Action": "sqs:SendMessage",
"Resource": "arn:aws:sqs:us-east-1:123456789012:cargo-results"
}Or use the CloudFormation role template
with the QueueArns parameter. Convoy derives the QueueUrl from the ARN —
no sqs:GetQueueUrl grant needed.
If the queue uses SSE with a customer-managed KMS key, also grant the
role kms:GenerateDataKey and kms:Decrypt on that key — standard
SQS-producer setup, not a Convoy-specific requirement. Without it,
deliveries fail terminally with target_misconfigured.
Submit cargo with an aws_sqs callback
{
"callback": {
"type": "aws_sqs",
"connection_id": "awsconn_…",
"queue_arn": "arn:aws:sqs:us-east-1:123456789012:cargo-results"
}
}For FIFO queues (.fifo suffix), message_group_id is required (and
rejected on standard queues):
{
"callback": {
"type": "aws_sqs",
"connection_id": "awsconn_…",
"queue_arn": "arn:aws:sqs:us-east-1:123456789012:cargo-results.fifo",
"message_group_id": "tenant-acme"
}
}Consume
Messages carry two attributes so consumers can filter/route without parsing the body:
def handle(message):
attrs = message["MessageAttributes"]
if attrs.get("convoy-verify"): # skip verification probes
return
cargo_id = attrs["cargo_id"]["StringValue"]
success = attrs["success"]["StringValue"] == "true"
envelope = json.loads(message["Body"])Delivery semantics
| Queue type | Semantics |
|---|---|
| FIFO | Exactly-once: Convoy sets MessageDeduplicationId to the cargo_id, so a redelivery inside SQS’s 5-minute dedup window is silently dropped — the best duplicate story of any ARN-addressed target. |
| Standard | At-least-once: dedupe on the cargo_id message attribute. |
Kinesis (aws_kinesis)
cargo completes ─▶ AssumeRole ─▶ kinesis:PutRecord(
StreamARN,
Data=<envelope>,
PartitionKey=<partition_key or cargo_id>)Grant kinesis:PutRecord on your role
{
"Effect": "Allow",
"Action": "kinesis:PutRecord",
"Resource": "arn:aws:kinesis:us-east-1:123456789012:stream/cargo-results"
}Or the CloudFormation template’s StreamArns parameter. Streams with SSE
(customer-managed KMS key) additionally need kms:GenerateDataKey on the
key.
Verify the capability from the dashboard (see Verifying the capabilities below).
Submit cargo with an aws_kinesis callback
{
"callback": {
"type": "aws_kinesis",
"connection_id": "awsconn_…",
"stream_arn": "arn:aws:kinesis:us-east-1:123456789012:stream/cargo-results"
}
}Partition key: defaults to the cargo_id (even spread across shards).
Pass a custom partition_key (1–256 chars — e.g. a tenant id) when you
need per-key ordering on a shard:
{
"callback": {
"type": "aws_kinesis",
"connection_id": "awsconn_…",
"stream_arn": "arn:aws:kinesis:us-east-1:123456789012:stream/cargo-results",
"partition_key": "tenant-acme"
}
}Consume
Records are the envelope JSON. Skip verification probes and dedupe on
cargo_id:
def handle(record):
envelope = json.loads(record["Data"])
if envelope.get("convoy_verification"):
return
if seen_before(envelope["cargo_id"]): # at-least-once
returnKinesis delivery is at-least-once — there is no server-side dedup, so
a lost response + retry can produce a duplicate record. Deduping on the
envelope’s cargo_id is standard practice for Kinesis consumers.
Verifying the capabilities (optional)
In the dashboard, open your AWS connection and run the SQS or Kinesis capability check with your queue or stream ARN.
Neither SendMessage nor PutRecord has a dry-run mode, so the sqs and
kinesis verify probes create one real marked side effect each — an
SQS message with body {"convoy_verification": true} plus a
convoy-verify message attribute, or a Kinesis record with that body and
PartitionKey=convoy-verify. Both are rate-capped at 1/minute per
connection, and both are opt-in (skip them and rely on fail-at-delivery if
you prefer).
The result envelope
Both targets deliver the standard envelope (message body / record data):
| Field | Description |
|---|---|
cargo_id | The cargo this result belongs to — dedupe key |
success | true / false |
response | Full model response (or null when truncated/failed) |
response_truncated | true when the response was too large to inline |
result_url | GET /cargo/{cargo_id}/result — fetch the full result |
error_message | Present on failure envelopes only |
metadata | Your submit-time metadata, echoed verbatim |
Troubleshooting
| Symptom | Cause / fix |
|---|---|
message_group_id_required (422) | FIFO queue without a message_group_id — add one. |
message_group_id_not_allowed (422) | Standard queue with a message_group_id — remove it. |
region_mismatch (422) | The queue/stream is in a different region than the connection. |
target_not_allowed (422) | The connection carries a target allow-list and this ARN isn’t on it. |
target_not_found in the DLQ | Queue/stream deleted (or renamed) after submit. |
target_misconfigured (KMS) | The SSE key can’t be used by Convoy’s session — grant the role kms:GenerateDataKey (+ kms:Decrypt for SQS) on the key. |
assume_role_denied + connection flagged failed | Trust policy or the send/put grant was removed. Fix the role and re-verify. |