Amazon EventBridge Pipes: Wire DynamoDB Streams / SQS to Coding-Agent Tool Runners Without Glue Lambdas

2 views

Every coding-agent platform eventually grows a mapper Lambda: read DynamoDB Streams for session expiry, parse SQS tool jobs, reshape the payload, call Step Functions or an HTTP tool runner. That glue is where idempotency bugs hide, where cold starts add 200ms to every tool hop, and where on-call discovers “we forgot the DLQ.” Amazon EventBridge Pipes collapses source → filter → optional enrichment → target into one managed pipe. Pair with DynamoDB Streams session expiry hooks, SQS FIFO ordered agent jobs, and Step Functions agent graphs. If you already use Pipes→Lambda enrich/fail-closed, this post is the next step: skip the Lambda target when the runner can consume the event directly.

⚡ TL;DR: Create a Pipe with DynamoDB Streams or SQS as source, a JSONPath filter that drops non-tool events, optional enrichment only when you truly need a lookup, and a target of Step Functions StartExecution, API Destination (HTTPS tool), or another SQS/ECS. Enable partial batch failure + DLQ. Measure end-to-end latency vs your glue Lambda and delete the mapper when P95 wins. Related: Lambda Destinations, EventBridge Scheduler, Recursive Loop Protection.

Why glue Lambdas rot in agent fleets

Coding-agent event paths are boring until they are not:

  1. Session row TTL expires → Stream REMOVE → kick cleanup / notify tenant
  2. Tool-job SQS message → start a Step Functions graph or POST to a tool microservice
  3. Ledger MODIFY with status=needs_retry → re-invoke a specific tool runner

A dedicated Lambda for each path means five IAM roles, five concurrency knobs, five places to forget ReportBatchItemFailures. Pipes gives you one pipe per path, managed scaling, and first-class filter/enrich/target configuration. Keep Lambdas for true business logic — not for reshaping CloudEvents.

Pattern When it wins When it loses
Glue Lambda Complex branching, multi-API fan-out in code Simple filter + one target
Pipe → Step Functions Multi-step tool graphs, waiters, retries Single fire-and-forget HTTP
Pipe → API Destination Existing HTTPS tool runners Needs VPC-only private DNS (use Lattice/PrivateLink)
Pipe → SQS/ECS Buffer + worker pool You needed sync response in-band

Pipe anatomy for coding-agent runners

A Pipe has four stages. For agent fleets, treat each as a contract:

  1. Source — DynamoDB Stream ARN or SQS queue ARN
  2. Filter — drop heartbeat noise; keep tool_invoke, session_expired
  3. Enrichment (optional) — small Lambda/Step Functions/API only if you must join tenant config
  4. Target — Step Functions, API Destination, Kinesis, SQS, EventBridge bus, ECS task
bash
# ✅ create a Pipe: DynamoDB Streams → Step Functions (no glue Lambda target)
aws pipes create-pipe \
  --name coding-agent-session-expiry-to-sfn \
  --role-arn arn:aws:iam::111122223333:role/AgentPipesRole \
  --source arn:aws:dynamodb:us-east-1:111122223333:table/AgentSessions/stream/2026-10-01T12:00:00.000 \
  --source-parameters '{
    "DynamoDBStreamParameters": {
      "BatchSize": 10,
      "StartingPosition": "LATEST",
      "MaximumRetryAttempts": 3,
      "MaximumRecordAgeInSeconds": 3600,
      "OnPartialBatchItemFailure": "AUTOMATIC_BISECT"
    }
  }' \
  --filter-criteria '{
    "Filters": [{
      "Pattern": "{\"eventName\":[\"REMOVE\"],\"dynamodb\":{\"OldImage\":{\"entityType\":{\"S\":[\"session\"]}}}"
    }]
  }' \
  --target arn:aws:states:us-east-1:111122223333:stateMachine:AgentSessionCleanup \
  --target-parameters '{
    "StepFunctionStateMachineParameters": {"InvocationType": "FIRE_AND_FORGET"}
  }'
json
// ✅ SQS → API Destination pipe filter: only tool jobs ready to run
{
  "Filters": [{
    "Pattern": "{\"body\": {\"type\": [\"tool_invoke\"], \"status\": [\"queued\"]}}"
  }]
}

❌ Mapping every Stream event through a Lambda that only renames fields — put that rename in the Pipe input transformer instead.

Input transformers beat mapper code

Pipes input transformers let you reshape DynamoDB Stream images into the JSON your Step Functions or tool API expects — without deploying code.

json
// ✅ input transformer: Stream REMOVE → cleanup execution input
{
  "InputTemplate": "{\"sessionId\": <$.dynamodb.OldImage.sessionId.S>, \"tenantId\": <$.dynamodb.OldImage.tenantId.S>, \"reason\": \"ttl_expiry\", \"source\": \"eventbridge-pipes\"}"
}
typescript
// ✅ CDK: Pipe from SQS tool queue → Step Functions graph
import * as pipes from "aws-cdk-lib/aws-pipes";
import * as iam from "aws-cdk-lib/aws-iam";

const pipeRole = new iam.Role(this, "AgentPipeRole", {
  assumedBy: new iam.ServicePrincipal("pipes.amazonaws.com"),
});
toolQueue.grantConsumeMessages(pipeRole);
stateMachine.grantStartExecution(pipeRole);

new pipes.CfnPipe(this, "ToolJobsPipe", {
  roleArn: pipeRole.roleArn,
  source: toolQueue.queueArn,
  sourceParameters: {
    sqsQueueParameters: { batchSize: 5 },
  },
  filterCriteria: {
    filters: [{ pattern: JSON.stringify({ body: { type: ["tool_invoke"] } }) }],
  },
  target: stateMachine.stateMachineArn,
  targetParameters: {
    stepFunctionStateMachineParameters: { invocationType: "FIRE_AND_FORGET" },
  },
});

Enrichment: use sparingly

Enrichment is the escape hatch when the event lacks tenant policy or model routing hints. Prefer reading those from AppConfig inside the target graph, or from DynamoDB in the state machine. If enrichment must exist:

  • Keep it idempotent and under 1s P95
  • Fail closed: on enrichment error, do not deliver a half-baked invoke
  • Put poison messages on a DLQ via source retry policy
python
# ✅ enrichment Lambda ONLY when join is unavoidable (keep tiny)
def handler(event, _ctx):
    # event is a list of records from the Pipe
    out = []
    for rec in event:
        tenant = rec["dynamodb"]["NewImage"]["tenantId"]["S"]
        # ❌ Do not call Bedrock or heavy APIs here — enrichment is not your agent loop
        rec["enrichment"] = {"tier": lookup_tier_cached(tenant)}
        out.append(rec)
    return out

Failure modes that matter for agents

Failure What you want How with Pipes
Bad filter miss Never invoke tool Tighten pattern; metric on matched count
Target 5xx Retry then DLQ Source retry + DLQ on SQS/Stream
Duplicate delivery Idempotent tool ledger DynamoDB Transactions at target
Recursive storm Cap fan-out Recursive Loop Protection on any Lambda still in path
Silent drop Alert Lambda Destinations-style alarms on Pipe failures
bash
# ✅ alarm when Pipe target invocations fail
aws cloudwatch put-metric-alarm \
  --alarm-name agent-pipe-target-failures \
  --namespace AWS/EventBridge/Pipes \
  --metric-name TargetInvocationFailed \
  --dimensions Name=PipeName,Value=coding-agent-session-expiry-to-sfn \
  --statistic Sum --period 60 --threshold 1 \
  --comparison-operator GreaterThanOrEqualToThreshold \
  --evaluation-periods 1

Production checklist

  • [ ] One Pipe per event path (session expiry, tool queue, ledger retry) — not one mega-pipe
  • [ ] Filter patterns reviewed: heartbeats and INSERT noise dropped
  • [ ] Input transformer covers field rename; glue Lambda deleted or scheduled for removal
  • [ ] Target is Step Functions / API Destination / worker queue — not “Lambda that only forwards”
  • [ ] Partial batch failure + DLQ configured; alarm on TargetInvocationFailed
  • [ ] IAM role for Pipes is least-privilege (source read + target invoke only)
  • [ ] Idempotency keys from eventID / message MessageId stamped into tool ledger
  • [ ] Latency dashboard: Stream/SQS age → target start vs old glue Lambda baseline
  • [ ] Documented ownership tag workload=coding-agent on every Pipe

FAQ

Q: How is this different from your earlier EventBridge Pipes → Lambda post?
A: That post is about safe enrichment when Lambda is the right enricher. This post is about eliminating the glue target so Streams/SQS land directly on the runner (SFN/API/ECS).

Q: Can Pipes replace Step Functions?
A: No. Pipes is the wire. Step Functions remains the multi-step brain for tool graphs, waiters, and human approval.

Q: DynamoDB Streams + Pipes vs Lambda event source mapping?
A: ESM is fine for compute-heavy handlers. Prefer Pipes when filter + transform + managed target covers 90% of the path and you want less code in production.

Related reading

Wire the stream to the runner with a Pipe, keep business logic in Step Functions or the tool service, and stop paying cold-start tax for JSON rename functions.

Last updated on October 1, 2026

Deep-dive PDF

Get the expanded guide for this post — extra diagrams-style checklists, failure modes, and a production walkthrough. Free when you subscribe to CheatCoders.

Already subscribed? or open the subscribe page.


Discover more from CheatCoders

Subscribe to get the latest posts sent to your email.

Comments

No comments yet. Why don’t you start the discussion?

Leave a comment

No account needed. Name and email are optional.