Skip to content

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 RuleActivated in the same atomic unit as the activation; the relay hands it to Billing’s subscriber.
  • In-process PubSub after 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 owns outbox_events, created by a normal migration with the outbox prefix (see DB migrations).
  • Postgres / SQLite: the use-case calls Outbox.append(event) inside Transactions.withTransaction. D1: the repository method that owns the batch adds Outbox.appendStatement(sql, event) to it; the use-case builds the event and passes it in (see SQL and transactions).
  • The row’s id is a UUIDv7 from newId(OutboxEventId), and is every consumer’s dedup key. The append is a plain INSERT with no ON 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 (BIGINT on Postgres, INTEGER on SQLite and D1).

❌ Incorrect — on D1, the event written in a second, separate call:

yield* rules.activate(orgId, ruleId) // batch 1 commits
yield* 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 subscribers list exported by main.ts. Claimed rows run concurrently, each under a 30s timeout. Claim size × timeout never exceeds the lease.
  • Settle: success sets processed_at. A failure sets available_at = now + Retry.durableDelay(attempts) (see scheduling and retry); after 10 attempts the row gets failed_at with a Warn. A payload SchemaError fails at once. An unknown event_tag is 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.contention only; 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 twice
SELECT * 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_events
SET claimed_until = ${now + 2 minutes}, attempts = attempts + 1
WHERE 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

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 is INSERT … 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 NOTHING on 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_id inside withTransaction; zero rows means already applied. D1: guard every statement with NOT 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, reset attempts) 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 history
DELETE 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

  • 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_at after exhausting all attempts on an error marked non-retryable.
  • An inbound Idempotency-Key header 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.