ADR 0003: Outbox, durable timers and per-load concurrency
Status: accepted (issue #1)
Context
Section titled “Context”ADR 0001 makes the per-load event log the source of truth and actions data. Three production problems remain:
- Dual write. A decision must be logged and its actions sent. Writing the log and then calling the SMS service can crash in between: either the SMS goes out with no record (the agent will ask again, a duplicate to the driver) or the record exists with no SMS (the agent believes it asked).
- Timers. Client A waits 30 minutes between attempts; loads live 2 days to 2 weeks. Timers must survive deploys and crashes, fire at least once, and be harmless when they fire late, twice, or after the ask was answered.
- Concurrency. Two runs for the same load at once (queue redelivery, a slow run past its visibility timeout, a timer and an SMS arriving together) would project the same state and both act. The design requires one load to be processed serially; that must be enforced, not assumed.
Scale: 100,000 loads per month, 50 to 100 messages per load, mostly in US business hours, growing tenfold. One run must finish within 5 minutes.
Decision
Section titled “Decision”Optimistic concurrency on the log
Section titled “Optimistic concurrency on the log”EventStore.read(load_id) returns the records and the current version (seq of the last record). EventStore.append(load_id, entries, expected_version, outbox) succeeds only if the version is still expected_version; otherwise it raises ConcurrencyConflictError and writes nothing. The run re-reads, re-projects and re-decides (bounded retries, then the queue redelivers). A duplicate log entry id raises DuplicateEntryError, meaning “already processed”, so a redelivered event is never processed twice. A duplicate outbox key raises a different error, DuplicateActionKeyError, meaning a decision bug: the run re-decides with a Hero task keyed by the event instead, so the inbound event is never lost.
The FIFO queue (message group per load_id) makes conflicts rare; the version check makes them harmless. An optional LoadLock lease (Postgres advisory lock or similar) avoids wasted LLM calls on a conflict; it is an optimization, never the correctness mechanism.
Transactional outbox
Section titled “Transactional outbox”One transaction writes the inbound event, the DecisionRecorded entry and one outbox row per accepted action (outbox_id = idempotency_key, unique). The dispatcher:
- claims due rows (
Outbox.pending(limit, now), PostgresFOR UPDATE SKIP LOCKED), - calls the owning service’s client with
idempotency_key=outbox_id(services dedupe on it), - on success
mark_dispatched(id, receipt), - on
RetryableDeliveryErrormark_failed(id, failure, retry_at)with capped exponential backoff, - on
PermanentDeliveryErroror retries exhaustedmark_failed(id, failure, None)(dead), - publishes an
ActionOutcome(delivered or failed) to the load queue with a deterministic event id.
The dispatcher never writes the log. Outcomes reach the log through the load queue, appended by the load worker like any other input. This keeps exactly one writer per load, so a burst of delivery receipts can never cause concurrency conflicts with decision runs. The cost: every delivered outcome is one more worker run that ends in the deterministic filter and appends a filtered noop decision (no LLM, milliseconds, two log rows). At 50 to 100 inbound messages and a similar number of actions per load, that roughly doubles log rows, which is cheap next to the alternative of conflict retries that re-run LLM calls.
A failed outcome of an external action is a situation the agent handles (usually an urgent Hero task: “SMS to driver failed, call them”). A failed outcome of a Hero-channel action (create_hero_task, internal Hero Slack post) must not create another Hero task, because that loops when the Hero tooling is down. The dispatcher raises an OpsAlert (pager or metric, independent of Hero tooling) and marks the outcome hero_channel; the agent records it as a filtered noop (needs_hero_task is false). Internal actions (set_timer, cancel_timer, update_stage) follow the same rule: their permanent failure is an infrastructure incident, so it raises an OpsAlert and never a Hero task.
Claims, leases and the retry budget
Section titled “Claims, leases and the retry budget”- Claims.
Outbox.pendinghands out rows with aclaim_tokenand a 90 second lease.mark_dispatchedandmark_failedrequire the token; a dispatcher whose lease expired and whose row was taken over getsStaleClaimErrorand changes nothing. The dispatcher also re-checks the claim right before each send. - Late FAILED next to DELIVERED. A send that outlasts the lease can still complete at the service while another dispatcher re-sends (the service dedupes on the idempotency key, so the contact gets one message). Each dispatcher may then publish an outcome, so the log can hold a
deliveredand afailed_permanentoutcome for one key, for example when the second dispatcher’s attempt hit a transient error after the first had succeeded. This is acceptable: outcome event ids differ by status so both are recorded, the agent sees both outcomes and the delivered one proves the message went out, and the worst case is one extra Hero task about a message that arrived, which a human closes in seconds. Avoiding it would need a fenced, exactly-once send, which the services cannot offer. Keeping sends well inside the lease (client timeouts of a few seconds) keeps this rare. - Failures after the send (outcome publish or marking fails, for example a queue outage) count as an attempt and back off with the normal curve, so an unavailable queue is probed with growing gaps. When the attempts are spent the row is buried; the failed outcome is published best effort. Ops alerts follow a successful publish, so an outage does not page on every pass.
- Retry budget.
RetryPolicy.budget(150 s, so a 120 s run plus delivery stays under 5 minutes) is measured from the first send attempt (first_attempt_at, set when the row is first claimed), not from commit: a dispatcher backlog should not spend the budget before anything was tried. The first retry is always allowed so one transient error never kills a row; later retries that would land after the budget kill it, and the failed outcome makes the agent escalate. Alternative considered: measuring from commit, which is simpler but turns a backlog into spurious escalations.
Durable timers through the same pipeline
Section titled “Durable timers through the same pipeline”set_timer and cancel_timer are outbox actions, so the intent to follow up commits atomically with the decision; the dispatcher applies them to the TimerStore (schedule is idempotent on timer_id, cancel and mark_fired are idempotent). The timer pump reads due(now, limit), publishes TimerFired with event_id = timer_fired_event_id(timer_id) to the load queue, then mark_fired. A crash between publish and mark publishes twice; the event store drops the second copy by event id. When the event is processed, OpenAsk.is_timer_current decides: a timer for a resolved ask, or for a step already passed, is a noop (ADR 0001). Cancellation is housekeeping only.
Anchor and tolerance: fire_at = decided_at + retry_after, where decided_at is the time of the decision that took the step. The pump polls at least every 30 seconds, so a timer fires within one minute after fire_at; outbox delivery of the SMS itself adds seconds. Evals assert the follow-up in the first scenario step at or after fire_at.
In evals, FakeClock drives the same pump, dispatcher and worker code: advance the clock, pump timers, process the queue, dispatch the outbox.
sequenceDiagram participant Q as Load queue participant W as Load worker participant S as Event store + outbox (one DB) participant D as Dispatcher participant SMS as SMS service participant T as Timer store participant P as Timer pump Note over Q,W: ask id A = ask:eta_request:ping_0700 Q->>W: LocationPing ping_0700 (risk high), received 07:00 W->>S: read -> version 7 W->>S: 07:01 append([ping, decision], expected 7, outbox [send_sms A/1, set_timer A/1]) D->>S: pending() D->>SMS: 07:01 send_sms(key 481207/A/1/send_sms) D->>T: schedule(tmr:481207:A:1, fire_at 07:01 + 30m = 07:31) D->>Q: ActionOutcome delivered (filtered noop when processed) P->>T: due(07:31) P->>Q: TimerFired(evt:tmr:481207:A:1) P->>T: mark_fired Q->>W: TimerFired: ask open, step 1 == attempt 1, take step 2 (SMS again, timer 08:01)
Alternatives considered
Section titled “Alternatives considered”| Alternative | Strengths | Why not (now) |
|---|---|---|
| Temporal (one workflow per load, durable timers, activities with retries) | Durable timers and retries built in; workflow history is an event log; strong visibility tooling. | A second source of truth next to our log, or our log moves into Temporal history (harder to query for evals, analytics, and “state as projection”). Determinism constraints on workflow code, versioning of long-running workflows (loads live up to 2 weeks) across SOP and code changes, and an extra cluster to operate or pay for. A good fit if workflow count and complexity grow; our contracts (TimerStore, Outbox) could be implemented on it. |
| SQS FIFO + EventBridge Scheduler (AWS managed) | Managed, scales without operating a database poller; FIFO message group per load gives ordering; one-time schedules for timers. | Still needs the outbox (SQS send is not in the database transaction) and idempotent consumers. EventBridge one-time schedules have quotas and per-schedule cost at our timer volume; cancellation is an API call that can fail. We plan to use SQS FIFO as the load queue in production behind LoadQueue, and keep timers in the database next to the outbox, where set_timer commits atomically with its decision. EventBridge Scheduler stays a drop-in TimerStore option. |
| Kafka partitioned by load_id | Strict per-partition ordering, high throughput, replayable topic. | Ordering is per partition, not per load, so one slow load blocks others in its partition (head-of-line). Delayed delivery (timers) is not native. Operating Kafka for about 10M messages per month is more than we need; a log per load in Postgres is simpler to query. Reconsider at 10x if Postgres write throughput becomes the bottleneck. |
| Postgres advisory locks as the serialization mechanism (pessimistic) | Simple mental model: hold the lock for the run. | A run can take minutes (LLM calls); holding locks that long across connection pools is fragile, and a lost connection silently releases the lock mid-run. Correctness would depend on lock hygiene. We use advisory locks only as the optional LoadLock lease; the version check is what guarantees correctness. |
| Direct send, no outbox (log then call the service) | Fewer moving parts, lower latency. | The dual-write failure above: duplicates to drivers or phantom asks after a crash. Not acceptable for messages to external people. |
Trade-offs
Section titled “Trade-offs”- Gains: no lost or duplicated actions across crashes; timers survive deploys; at-least-once everywhere is made safe by deterministic ids and idempotency keys; one database transaction per decision; same code paths in evals and production.
- Costs: a dispatcher and a timer pump to run and monitor (lag, dead rows); delivery adds seconds of latency; Postgres carries log, outbox and timers, which needs partitioning and retention at 10x volume; conflicts re-run a decision (rare, and the optional lease reduces them).
Consequences
Section titled “Consequences”EventStore,OutboxandTimerStorehave contract test suites run against SQLite (local and CI) and Postgres (production pipeline).- Every externally visible effect has an idempotency key or a deterministic event id: action keys,
timer_id,timer_fired_event_id,action_outcome_event_id. - Monitoring: outbox age of oldest pending row, dead rows per hour, timer pump lag, conflict rate per load. Dead rows of external actions produce a Hero task through the failed
ActionOutcome; dead rows of Hero-channel and internal actions produce an ops alert. - Services owned by the integrations team must dedupe on the idempotency key. This is a stated requirement of the client contracts in
channels/base.py.