Idempotency & the outbox
Delivery is at-least-once everywhere, and exactly-once is faked by the consumer. This page answers which events go through an outbox, how the row is written and relayed, what the dedup key is, and how a consumer survives running twice.
Deliver every domain-pair event through the outbox
Section titled “Deliver every domain-pair event through the outbox”Impact: HIGH a lost fact leaves the subscriber wrong forever, with no error anywhere
- The reverse direction of a domain pair (see
coupling and cohesion) always goes through the
outbox. Labeling appends
RuleActivatedin the same atomic unit as the activation; the relay hands it to Billing’s subscriber. - In-process
PubSubafter commit stays for signals that are not domain-pair events: realtime push, cache invalidation, a metric. Losing one leaves no state wrong. - Other work that must survive a restart or deploy (an outbound webhook, an email, a provider
call) is an outbox row too.
makeWorker’s in-memory queue is for work that may be lost. - Work triggered by an upstream that redelivers needs no outbox: answer 2xx only after the durable write, and 5xx otherwise so the upstream retries.
❌ Incorrect — publish after commit; a crash in between loses the fact:
yield* transactions.withTransaction(rules.activate(orgId, ruleId))yield* PubSub.publish(events, RuleActivated.make({ id, ruleId, orgId }))✅ Correct — the event is written with the change:
yield* transactions.withTransaction(Effect.gen(function* () { yield* rules.activate(orgId, ruleId) yield* outbox.append(RuleActivated.make({ id: yield* newId(OutboxEventId), ruleId, orgId }))}))Source: notes/09-production-concerns/idempotency-and-outbox.md · Decision 1
Write the outbox row in the same atomic unit as the change
Section titled “Write the outbox row in the same atomic unit as the change”Impact: HIGH one shape on Postgres, SQLite and D1
infra/outbox/is a foundation feature. It ownsoutbox_events, created by a normal migration with theoutboxprefix (see DB migrations).- Postgres / SQLite: the use-case calls
Outbox.append(event)insideTransactions.withTransaction. D1: the repository method that owns thebatchaddsOutbox.appendStatement(sql, event)to it; the use-case builds the event and passes it in (see SQL and transactions). - The row’s
idis a UUIDv7 fromnewId(OutboxEventId), and is every consumer’s dedup key. The append is a plainINSERTwith noON CONFLICT. - The payload is the event schema encoded to JSON and evolves additively; a breaking change is a
new
event_tag. Every time column is epoch ms (BIGINTon Postgres,INTEGERon SQLite and D1).
❌ Incorrect — on D1, the event written in a second, separate call:
yield* rules.activate(orgId, ruleId) // batch 1 commitsyield* outbox.append(RuleActivated.make({ id, ruleId, orgId })) // batch 2 may never run✅ Correct — one batch in one repository method:
activate: (orgId: OrgId, ruleId: RuleId, events: ReadonlyArray<OutboxEvent>) => d1.batch([ sql`UPDATE labeling_rules SET active = 1 WHERE org_id = ${orgId} AND id = ${ruleId}`, ...events.map((event) => Outbox.appendStatement(sql, event)), ]).pipe(Effect.asVoid, Effect.catchTag("SqlError", Persistence.fail("LabelingRulesRepo.activate")))Source: notes/09-production-concerns/idempotency-and-outbox.md · Decision 2, amended
Relay with a claim lease, dispatch in-process, promise no order
Section titled “Relay with a claim lease, dispatch in-process, promise no order”Impact: HIGH concurrent relays are safe without a leader
- The relay claims a batch in one statement with a lease and an attempt count. Postgres adds
FOR UPDATE SKIP LOCKED; SQLite and D1 have one writer. An expired lease, or one more than twice its length in the future, is reclaimable while the process runs. - Dispatch is in-process, through the
subscriberslist exported bymain.ts. Claimed rows run concurrently, each under a 30s timeout. Claim size × timeout never exceeds the lease. - Settle: success sets
processed_at. A failure setsavailable_at = now + Retry.durableDelay(attempts)(see scheduling and retry); after 10 attempts the row getsfailed_atwith a Warn. A payloadSchemaErrorfails at once. An unknownevent_tagis retried normally. A defect counts as a failed attempt. - Waking: the container polls every 1s on every replica (backing off on failed ticks) and is
nudged after commit; the Worker runs
waitUntil(relay.drainOnce)after the response, plus a* * * * *cron backstop. No lease job is needed. - No ordering guarantee. A subscriber that cares guards by state or version. A subscriber retries
in process on
Retry.contentiononly; the relay owns every other retry.
❌ Incorrect — a large claim run one by one, so a slow batch outlives its lease:
-- claim 100, lease 2 minutes, process sequentially → rows processed twiceSELECT * FROM outbox_events WHERE processed_at IS NULL ORDER BY created_at LIMIT 100✅ Correct — one claim statement, ten rows, lease and attempt count:
UPDATE outbox_eventsSET claimed_until = ${now + 2 minutes}, attempts = attempts + 1WHERE id IN ( SELECT id FROM outbox_events WHERE processed_at IS NULL AND failed_at IS NULL AND available_at <= ${now} AND (claimed_until IS NULL OR claimed_until < ${now} OR claimed_until > ${now + 4 minutes}) ORDER BY created_at LIMIT 10)RETURNING *Source: notes/09-production-concerns/idempotency-and-outbox.md · Decision 3, amended
Dedup on the producer’s own identity
Section titled “Dedup on the producer’s own identity”Impact: HIGH the best key is one that already exists
| what arrives | its dedup key |
|---|---|
| a webhook | the upstream’s event or delivery id |
| a create from our web client | the entity id, generated by the client |
| an outbox event, at its subscriber | the outbox row’s id |
| a webhook with no upstream id | a hash of source, canonical JSON and the upstream’s timestamp, never our receive time |
- A retryable create takes a client-generated branded id, decoded through
makeEntityId(a non-UUID is a 400). The insert isINSERT … ON CONFLICT (id) DO NOTHING RETURNING *. On zero rows, read by id: the same body is a replay and returns the stored row; otherwise fail with a conflict error. See IDs and identity. - A mutation that is not a create is idempotent by meaning: it sets a state. One that can only “add” takes a client-generated operation id and is handled like a create.
- An outbound call to a provider that accepts a key gets one derived from our event id.
❌ Incorrect — the server mints the id, so a client retry creates a duplicate:
create: (orgId: OrgId, input: NewRule) => newId(RuleId).pipe(Effect.flatMap((id) => insertRow(toRow({ ...input, id }))))✅ Correct — the client’s id, replay-safe:
insert: (orgId: OrgId, rule: LabelingRule) => insertRow(toRow(rule)).pipe( Effect.flatMap(Option.match({ onSome: (row) => Effect.succeed(fromRow(row)), onNone: () => findById(orgId, rule.id).pipe(Effect.flatMap((existing) => Option.isSome(existing) && sameBody(existing.value, rule) ? Effect.succeed(existing.value) // replay : Effect.fail(new RuleIdConflict({ id: rule.id })))), })), )Source: notes/09-production-concerns/idempotency-and-outbox.md · Decision 4, amended
Make a consumer idempotent by its write, or write a receipt with it
Section titled “Make a consumer idempotent by its write, or write a receipt with it”Impact: HIGH a wrong call is a silent double increment
- A handler that sets state needs nothing more: an upsert, an insert with
ON CONFLICT DO NOTHINGon a natural key, or a version guard. - A handler that adds writes a receipt, keyed by the dedup key, in the same atomic unit as its
effect. Postgres/SQLite: insert the receipt with
ON CONFLICT DO NOTHING RETURNING event_idinsidewithTransaction; zero rows means already applied. D1: guard every statement withNOT EXISTS, receipt insert last. - Never claim a receipt in its own statement before the work; that turns at-least-once into
at-most-once. Each slice owns its receipts table (
billing_receipts). - An outbound call with no key goes last in the handler. In a process step that commits a transition after calling out, the call goes before the commit, under a key (see workflows and jobs).
- Review asks every subscriber and webhook handler: “does this set or add? If it adds, where is the receipt?”
❌ Incorrect — the receipt claimed first; a crash after it loses the effect:
yield* sql`INSERT INTO billing_receipts (event_id, created_at) VALUES (${e.id}, ${now})`yield* sql`UPDATE billing_quota SET active_rules = active_rules + 1 WHERE org_id = ${e.orgId}`✅ Correct — guarded effect and receipt in one batch (D1 form):
d1.batch([ sql`UPDATE billing_quota SET active_rules = active_rules + 1 WHERE org_id = ${e.orgId} AND NOT EXISTS (SELECT 1 FROM billing_receipts WHERE event_id = ${e.id})`, sql`INSERT INTO billing_receipts (event_id, created_at) VALUES (${e.id}, ${now}) ON CONFLICT DO NOTHING`,])Source: notes/09-production-concerns/idempotency-and-outbox.md · Decision 5, amended
Keep processed rows 7 days, receipts 30 days, failed rows until handled
Section titled “Keep processed rows 7 days, receipts 30 days, failed rows until handled”Impact: MEDIUM both tables otherwise grow without bound
- Processed outbox rows are deleted after 7 days; they are debugging history. Consumers dedup on their receipts, not on the outbox.
- Receipts are deleted after 30 days, longer than any window in which a duplicate can arrive.
- Failed rows are never deleted automatically. They are the dead-letter record; an operator
requeues them (clear
failed_at, resetattempts) or deletes them. - Cleanup is idempotent
DELETEs with no lease: hourly on every container replica, a daily cron on the Worker. Each slice’s receipts table registers its own retention with the cleanup.
❌ Incorrect — delete on success and keep receipts briefly:
DELETE FROM outbox_events WHERE id = ${id}; -- no delivery historyDELETE FROM billing_receipts WHERE created_at < ${now - 48h}; -- a late redelivery applies twice✅ Correct — retention by age, failed rows kept:
DELETE FROM outbox_events WHERE processed_at < ${now - 7 days};DELETE FROM billing_receipts WHERE created_at < ${now - 30 days};Source: notes/09-production-concerns/idempotency-and-outbox.md · Decision 6
Deferred
Section titled “Deferred”infra/outbox/(the table,append, the relay and the cleanup) — trigger: the first domain-pair event, or the first write whose side effect must not be lost.- Replacing the relay with Effect’s
PersistedQueue— trigger: an Effect release where its SQL store separates migrations from storage and runs on D1. - Ordered, per-aggregate dispatch — trigger: a bug traced to two events of one aggregate applied out of order.
- Failing a row at once on a non-retryable error — trigger: more than 10 rows in a week reach
failed_atafter exhausting all attempts on an error marked non-retryable. - An inbound
Idempotency-Keyheader with stored responses — trigger: a second audience, an API consumer that is not our web client. - Shorter retention — trigger: the outbox and receipts tables together exceed 1 GB, or 10% of the database.