Exactly once, with an asterisk: except in two places it quietly isn't. Below, the pipeline from DynamoDB events table to Streams, EventBridge Pipe, SQS FIFO and the inbox.

The Outbox Without Postgres: Exactly-Once Events on DynamoDB, and the Two Places It Quietly Isn't

The problem: two writes, one truth

Every event-driven backend eventually writes this bug. A handler saves a row, then publishes an event. The save succeeds. The publish times out. Now the database says one thing and the rest of the system believes another, and nothing anywhere has failed loudly.

The textbook fix is the transactional outbox. You write the event into an outbox table in the same database transaction as the state change, and a relay ships it later. On Postgres that is a few lines: BEGIN, two inserts, COMMIT.

Our backend has no Postgres. It is a modular monolith on AWS for a pet-health platform: a mobile app, a smart food bowl, an AI chat assistant, a commerce flow for pet tags. Transactional data lives in DynamoDB. The event engine is the centre of the system. When a payment completes, four things must happen: mirror the order to Shopify, book a shipment, notify the buyer, and mark the pet tag as ordered. If any of those runs twice, or never, a real person notices.

This article walks through how we built an outbox, a relay and an idempotent inbox from DynamoDB, Streams, EventBridge Pipes and SQS FIFO. It also covers the part that most write-ups skip. While mapping our own code for a design review, we found two places where "exactly once" was quietly not true. Both are subtle, both passed every test, and both are worth knowing before you build the same thing.

The outbox is one TransactWriteItems

DynamoDB has no BEGIN, but it has TransactWriteItems: up to 100 conditional writes that commit together or not at all. That is enough for an outbox.

Every module emits through one function, default_engine().emit(type, payload, extra_writes=[...]). The engine builds an envelope (id, type, tenant, correlation id, actor, payload) from request context, then appends it to a dedicated events table in a single transaction with three parts:

  1. The chain head. One row per tenant, TENANT#<t>#HEAD, holding last_seq and last_hash. It is updated on the condition last_seq = :prev, so two concurrent emits for one tenant cannot both win. The loser retries, up to five times.
  2. The event row. pk = TENANT#<t>#EVENT, sk = the sequence number zero-padded to 20 digits, so a plain query returns the tenant's log in order.
  3. The caller's own writes. Whatever the module passed in extra_writes: the order row, the pet record, the counter bump.
await default_engine().emit(
    "com.order_created",
    {"order_id": order_id},
    extra_writes=[{"Put": {"TableName": orders, "Item": order_item}}],
)

That third part is the whole point. The entity and its event now commit atomically. There is no window in which the order exists and the event does not.

Diagram: emit() issues one TransactWriteItems containing the chain head update, the event row and the caller’s extra writes, above a per-tenant hash chain
One emit, one transaction: the chain head, the event row and the caller’s own writes commit together, and each event extends the tenant’s hash chain.

The head row buys a second property for free. Each event stores hash = sha256(prev_hash + canonical_json(envelope)), starting from a fixed genesis value. The log is a per-tenant hash chain, and a verify_chain() call recomputes it. An IAM policy denies UpdateItem and DeleteItem on event rows, so the log is append-only by policy as well as by code. An operator endpoint reports chain_intact next to the dead-letter count.

The cost is contention. Every emit for a tenant serialises on one head row. For a household-scale tenant that is fine. For a tenant that emits thousands of events a second it would not be, and high-rate sensor data deliberately does not go through this path.

The relay: Streams, Pipes, SQS FIFO

An outbox needs a relay that reads new rows and hands them to consumers. On Postgres you poll the table or tail the WAL. On DynamoDB the table already publishes its own change log: DynamoDB Streams. We wired it with no relay code at all.

Architecture diagram: emit() to events table to DynamoDB Streams to EventBridge Pipe to SQS FIFO on the top row; worker, claim, handler, mark_completed and delete on the bottom row, with a red failure branch from the handler to mark_dead and delete
The top row is the outbox and its relay, with no code of ours between the table and the queue. The bottom row is the inbox; the red branch is where the trouble described below begins.

The events table has a stream with new and old images. An EventBridge Pipe reads it, keeps only INSERT records, drops the chain-head rows, and writes each event to an SQS FIFO queue. Two settings on the pipe matter:

  • MessageGroupId = tenant_id. FIFO ordering holds within a group, so one tenant's events stay in order while tenants never block each other.
  • MessageDeduplicationId = event id. A retried pipe delivery inside the five-minute dedup window collapses to one message.

A worker process (its own ECS service, scaled on queue depth) long-polls the queue, ten messages at a time, and processes them one by one to keep group order.

One detail cost us more than it should have. The pipe's input template forwards the raw stream image, <$.dynamodb.NewImage>, as one blob, and the worker un-types it back into an envelope. We tried naming envelope fields individually in the template first. It cannot work: the payload is a DynamoDB map whose keys change per event type, and optional fields are simply absent from real rows. A template path that does not resolve does not raise anything. The pipe just never delivers, the queue stays empty, every durable subscriber stops running, and the test suite stays green. Both pipes now log errors to CloudWatch so a broken transform is at least visible.

Worth stating plainly: stream order is per partition key, not across transactions. The log itself is strictly ordered by the per-tenant sequence number. Consumers must not assume they see events in that order, and must be idempotent.

The inbox: claim before you act

SQS is at-least-once. FIFO dedup narrows duplicates; it does not remove them. So the consumer side needs an inbox: a record that says "this handler already ran for this event".

Ours is an idempotency table with one row per (event, subscription), keyed IDEMP#{event_id}:{subscription}. The worker's delivery loop for one event looks like this:

for sub in matching_durable_subscriptions(env):
    key = f"{env.id}:{sub.name}"
    if not await idempotency.claim(key):   # put_item, attribute_not_exists(pk)
        continue                            # someone already has it
    try:
        await sub.handler(env)
        await idempotency.mark_completed(key)
    except Exception as exc:
        await idempotency.mark_dead(key, str(exc))
        await dead_letters.record(env, sub.name, str(exc))

Three choices here are deliberate.

  • The key is per subscription, not per event. One com.payment_completed fans out to four handlers. If the Shopify mirror fails, the shipment, the email and the tag update still run, each exactly once.
  • The tenant comes from the envelope. The handler runs inside the event's tenant context, so every write it makes lands in the right partition, and the dead-letter row does too.
  • Failures become data. A dead handler gets a row in a tenant-scoped dead_letter table, readable through the same reconciliation endpoint that reports chain integrity. A message the worker cannot even parse is a poison pill: it is recorded under a global partition and deleted at once, so one bad body cannot block a tenant's FIFO group forever.

Behind all of this, each SQS queue has its own dead-letter queue with a CloudWatch alarm at a threshold of zero. A DLQ with anything in it is, by definition, work nobody did.

A lane is reachability, not a preference

Not every subscriber wants a queue between it and the event. The engine has four delivery lanes, and each subscription declares one:

LaneRunsGuaranteeDispatched by
INLINEIn the caller, awaitedFail-closed: an error fails the requestemit()
LOCALSame process, fire-and-forgetBest effort, errors loggedemit()
DURABLEA worker, via SQS FIFOAt-least-once, idempotent inboxThe durable consumer
BEST_EFFORTAny worker, via a record-change queueNo order, no dead letterThe record-change consumer

The lesson we learned the hard way: each lane has exactly one dispatcher, and a subscription on the wrong lane is not slow or flaky. It is unreachable, in silence. A handler listening for row changes on LOCAL never runs, because those changes surface on another worker. A domain handler on BEST_EFFORT never runs, because that consumer only dispatches row changes.

Diagram of three dispatchers feeding four lanes; a row-change handler on LOCAL and a domain handler on BEST_EFFORT are marked as never running
Each lane has exactly one dispatcher. A subscription on a lane its dispatcher never feeds binds without error and never runs.

The matching is just as quiet. Subscriptions match event types with fnmatch, and a glob with no wildcard is an exact string compare. We had a subscriber whose glob said CONTACTADDR for an entity called CONTACTADDRESS, on the wrong lane as well. It bound without error, appeared in the generated catalog, dead-lettered nothing and never ran. Users who changed their email kept getting notifications at the old one.

The fix was not more care. It was a contract test that checks the catalog both ways: every emitted type is declared, every declared emitter has a producer, every glob matches something emittable, and every row-change subscriber sits on BEST_EFFORT. The engine also logs a warning when a durable event reaches the worker and matches nothing, because "delivered to nobody" and "delivered to five handlers" otherwise look the same.

The two places exactly-once quietly isn't

Everything above is sound, and every test passed. Then we read the inbox code again, line by line, while drawing the system for a design review. Two properties we believed did not hold.

Three timelines: a handler error marked dead and deleted after one try; a worker crash leaving a pending row that blocks the redelivered claim; a redrive blocked by the dead row
Three traces through the inbox. In each one SQS does its job, and the claim or the error path throws the work away.

1. A failing handler is never retried

The durable queue has a redrive policy: five receives, then the SQS dead-letter queue. The dead-letter module's docstring says a handler is dead-lettered "after exhausting SQS redrive". The code does something else. Look at the loop again: on the first exception the handler is marked dead, the dead-letter row is written, and the worker deletes the message as it would after a success.

So the five redrives never apply to handler errors. They only cover faults below the handler: the worker crashing, the idempotency table throttling, the delete call failing. A two-second Bedrock throttle, a Shopify 503, a cold DNS lookup: each one kills that handler's run for that event, permanently, after a single try.

2. The claim blocks every second chance

claim() is a put_item conditioned on attribute_not_exists(pk). That condition fails for any existing row: completed, but also pending and dead. Two consequences follow.

  • A redrive does not rerun dead work. The runbook says: fix the handler, deploy, then move the DLQ back to its source queue. The redelivered message reaches the worker, the claim fails because the dead row exists, and the handler is skipped. The idempotency rows expire after seven days, so the fix works only for events older than a week.
  • A crash mid-handler loses the run. If the worker dies after the claim and before mark_completed (a deploy, an out-of-memory kill, a task replaced by the scheduler), the row stays pending. SQS correctly redelivers the message after the visibility timeout. The claim fails. The handler is skipped. That is at-most-once, wearing an exactly-once badge.

What the fix looks like

The claim has to know the difference between "done" and "in progress" and "failed":

  1. Claim with a lease. Store status plus lease_expires_at. A claim succeeds if the row is absent, or pending with an expired lease, or dead with attempts < N. One conditional expression covers all three.
  2. Let SQS do the retrying. On a retryable error, release the claim and do not delete the message. The visibility timeout and the existing maxReceiveCount become the retry policy you already configured. Only a terminal error, or the last receive, marks the row dead.
  3. Give operators a real replay. A redrive tool that clears the dead rows for the selected events before moving messages back. Without it, "redrive after the fix" is a ritual that does nothing.
State machine for the idempotency row: absent, pending with lease, completed and dead, with take-over on lease expiry, release on retryable error and replay from dead
The claim as a small state machine. Only completed is terminal; a stale pending row and a dead row with attempts left can both be claimed again.

None of this needs a new service. It is a different condition expression and a different branch in the error path.

What to take away

You can build a credible outbox on DynamoDB with no relay code: one TransactWriteItems for state plus event, Streams and Pipes for transport, SQS FIFO keyed by tenant for order, and a claim table for the inbox. If you do, check these before you trust it:

  • State and event commit in one transaction, through one emit function. Nothing writes the event table directly.
  • The pipe forwards the raw stream image, and its failures log somewhere you look.
  • FIFO group = tenant, dedup id = event id.
  • Idempotency is per (event, subscription), not per event.
  • Your claim can take over a stale pending row and retry a dead one. Write the crash test: kill the worker between claim and complete, and assert the handler runs on redelivery.
  • A retryable handler error leaves the message for SQS. Write the flaky test: fail once, succeed on the second receive.
  • Redrive clears the claims it needs to. Test the runbook, not just the code.
  • A contract test proves every subscription is reachable: right lane, a glob that matches a real event type.
  • "Matched no subscriber" is logged, and every DLQ alarms at zero.

The outbox was the easy part. The hard part was the word "once", and the only reliable way we found to check it was to read our own error paths as if someone else had written them.

Read more

Preventing Unsafe Reconfiguration While Subsystems Are Active

Preventing Unsafe Reconfiguration While Subsystems Are Active

Introduction Modern embedded products are becoming increasingly modular. A single hardware platform may support: * Multiple sensor configurations * Different communication modules * Optional accessories * Replaceable compute modules * Expandable peripheral boards This flexibility improves product scalability, but it introduces a hidden reliability challenge: What happens when someone tries to change the system configuration

By Vinayak M K