ADR 0001: Hybrid event-driven architecture with an append-only log and open asks
Status: accepted (issue #1)
Context
Section titled “Context”- Volume: 100,000 loads per month, 50 to 100 inbound messages per load, plus location pings. About 7.5M messages per month before pings.
- Budget: about $1 of LLM inference per load over its life. If every message called an LLM, that is about $0.013 per event, which rules out a large model reading full history on every event.
- Latency: one run must finish within 5 minutes from event to last action, infrastructure included.
- Most SOP behavior is procedural (wait 30 minutes, count attempts, escalate in order). The assignment scenario (load #481207) is a retry state machine with one language step in the middle (the driver asks for the PO number).
- The agent must remember its own past actions and messages, and it must know when not to act (“ok” with no open question; a timer for an ask that was already answered).
Decision
Section titled “Decision”An event-driven pipeline, serial per load, where deterministic code owns state and procedure and the LLM is called only for language and for uncovered situations.
flowchart LR
subgraph Sources[Integrations team services]
SMS[SMS] --> N
EM[Email] --> N
SL[Slack] --> N
PG[Pings] --> N
AT[Attachments] --> N
end
N[Normalizer: InboundEvent] --> Q[Load queue, FIFO per load_id]
TP[Timer pump] -- TimerFired --> Q
subgraph Worker[Load worker: one run per event]
R[Read log snapshot + version] --> P[Project LoadState]
P --> F{Deterministic filter}
F -- duplicate, stale timer, no band change, ack with no open ask --> D
F --> TR[Trigger diff + triage: decision model]
TR --> X[SOP executor: compiled spec]
TR --> PL[Planner: strong model, rare]
X --> V[Validate actions]
PL --> V
V --> D[DecisionRecord]
end
Q --> R
D --> C[(ONE transaction: log entries + outbox rows, at expected_version)]
C --> DS[Dispatcher]
DS -- idempotency key --> SVC[SMS / email / Slack / TMS / Hero tasks / load status clients]
DS -- set_timer, cancel_timer --> TS[(Timer store)]
TS --> TP
DS -- ActionOutcome --> Q
Memory is a log, state is a projection
Section titled “Memory is a log, state is a projection”- Every load has one append-only log (
EventStore). Entries (load_agent.domain.log.LogEntry):LoadRegistered: the metadata at start, so even metadata is replayable.- Inbound events:
sms,email,slack,ping,timer,attachment, andstop_status(arrived, departed or delivered at a stop, from geofencing, the TMS or a Hero). action_outcome: delivered or permanently failed, reported by the dispatcher through the load queue (so the load worker stays the only writer of the log; see ADR 0003).DecisionRecorded: oneDecisionRecordper processed event (triggering event, SOP ref and version, ask transitions, actions, rationale, decision path, LLM usage with model id, tokens and cost).
LoadStateis a pure fold over the log (store.projection.project_state). Replaying the log yields an equal state. Nothing writes state directly.- Ask state changes only through
AskTransitions (ask_opened,ask_step_taken,ask_resolved,ask_cancelled) carried inside the decision record. The decision is the unit of change, so a log never contains an action without the state change that explains it.
Open asks are first class
Section titled “Open asks are first class”OpenAsk holds: id, type, target contact, pinned SOP ref and version, opened_at, attempts, max_attempts, escalations_done, next timer (id and fire time), resolved_by and resolved_at. “Is the ETA request still open?” is state.open_asks_of_type("eta_request"), never an LLM reading chat history.
Ask ids are derived, ask_id_for(opened_by_event_id, type), so re-deciding after a concurrency conflict yields the same ids and keys. In the example below the ask opened by ping ping_0700 is ask:eta_request:ping_0700.
Step counter. Contact attempts and escalation steps share one counter, step = attempts + escalations_done. An action carries ask_id and attempt only when its decision takes that step (an ask_step_taken transition with step == attempt in the same record; DecisionRecord enforces it). Anything else related to an ask that does not advance it is keyed by its causing event. Otherwise, for example, re-explaining the question after the driver texts ”?” at 07:20 would reuse the 07:01 key, collide in the outbox and lose the 07:20 event. For load #481207 under Client A:
| Time | Action | attempt | Ask after |
|---|---|---|---|
| 07:01 | SMS driver for ETA, timer for 07:31 | 1 | attempts 1 |
| 07:13 | SMS driver with PO 55821 (no ask, keyed by triggering event) | none | unchanged |
| 07:31 | SMS driver again, timer for 08:01 | 2 | attempts 2 |
| 08:01 to 08:02 | Email dispatcher in main thread | 3 | escalations 1 |
| 08:01 to 08:02 | Urgent Hero task to call the driver | 4 | escalations 2 |
Timers carry exactly {load_id, ask_id, attempt} and a deterministic timer_id = tmr:{load_id}:{ask_id}:{attempt}. They fire at decided_at + retry_after, anchored at the decision that took the step (07:01 plus 30 minutes is 07:31), and the pump fires them within one minute after that. On fire, the timer is current only if OpenAsk.is_timer_current(payload): the ask is open and payload.attempt == ask.step. Otherwise (ask resolved by the dispatcher’s email at 07:20, step already advanced, duplicate delivery) the decision is a noop. Cancelling timers is best effort; correctness never depends on it.
Triggers are a pure diff of state before and after an event (agent.triggers.detect_triggers), not raw event kinds. A ping at “high” when risk for that stop was already high is not a trigger, and a ping recorded earlier than the last applied one is stale and changes nothing. Delay risk is tracked per stop, so “high for the pickup” does not suppress a later “high for the delivery”.
What moves the stage. Only stop_status events (geofence entry or exit from the tracking service, a TMS status update, or a Hero) change stage and current_stop_sequence in the projection. The resulting diff produces StageEntered; departing the pickup gives at_pickup to to_delivery, which fires Client B reefer’s “driver leaves the pickup”. The agent’s update_stage action only mirrors the stage to our load status service and the TMS; it never moves state, so a wrong mirror cannot trigger SOPs.
Arrival closes the ETA ask. StageEntered(at_pickup) or StageEntered(at_delivery) cancels the open asks that await only an ETA and are about the stop reached (OpenAsk.stop_sequence, set from the trigger’s stop when the ask opens; AskCancelled("arrived at the stop") plus a cancel_timer). StageEntered(delivered) cancels all open asks. Without this the 07:31 follow-up would text a driver who is already at the dock and escalate to the dispatcher and a Hero. Scoping by stop matters when the delay risk of the delivery leg goes high while the truck is still heading to the pickup: arriving at the pickup must not cancel that delivery ETA ask. Cost: one optional field on OpenAsk; asks logged without it count as about the current stop.
Actions are data, committed with their decision
Section titled “Actions are data, committed with their decision”A decision returns a tuple of Action models (send_sms, send_email with main or new thread, post_slack to the broker or the internal Hero channel, tms_note, set_timer, cancel_timer, create_hero_task with urgent, update_stage, noop). Decisions never perform side effects. Actions are validated before commit (contact belongs to the load, channel supported by the contact, role not forbidden by the SOP, key not in LoadState.committed_action_keys); a rejected action is replaced by a Hero task. The inbound event, the decision record and the accepted actions are written in one transaction at the snapshot version (log rows plus outbox rows). A dispatcher then sends each outbox row through the owning service’s client with the action’s idempotency key. Details and alternatives: ADR 0003.
Idempotency key, derived, never supplied:
- with an ask:
{load_id}/{ask_id}/{attempt}/{kind} - without an ask:
{load_id}/evt:{caused_by_event_id}/{kind} /{seq}appended when one event emits several actions of the same kind.
kind is part of the key because one step emits an SMS and a set_timer with the same ask and attempt.
A noop is an explicit Noop action in a decision of its own, so audits and evals can tell “decided to do nothing” from “crashed” or “skipped”. DecisionRecord.effective_actions is empty for a noop.
Serial per load
Section titled “Serial per load”The queue is FIFO per load_id (SQS FIFO message group, Kafka key), so normally one run per load is in flight. The guarantee does not rest on the queue: every commit carries expected_version, and a concurrent run for the same load gets ConcurrencyConflictError, re-reads and re-decides. Different loads run in parallel. See ADR 0003.
Storage and service boundaries
Section titled “Storage and service boundaries”- Event store. One contract (
store.base.EventStore), two implementations behind it: SQLite is the hermetic store for local runs and evals (no server, deterministic, one file per run); Postgres is the production store (row-level transactions for log plus outbox,SKIP LOCKEDfor the dispatcher, partitioning by load). Both run the same contract test suite, so behavior proven in evals is the behavior of the contract, not of a toy. - Services of other teams. SMS, email, Slack, TMS, Hero tasks and load status are services owned by the integrations team, reached over HTTP or a queue. The agent depends on typed clients (
channels.base) that take an idempotency key and raiseRetryableDeliveryErrororPermanentDeliveryError. Local runs and evals use in-process fakes of those services at the same contract. That is the correct production boundary: we own the agent, not the SMS provider, and our tests should pin the contract we call, exactly as the integrations team’s own contract tests pin theirs.
The LLM is behind a narrow, structured boundary
Section titled “The LLM is behind a narrow, structured boundary”LLMClient has five methods: triage (decision model, TypeSafe Jev), draft (message bodies), extract_document (vision), plan (strong model, rare) and extract_sop (strong model, edit time only, the compile-time seam; never on the load path). Each takes a Pydantic request and returns LLMResult[T]: the validated output plus LLMUsage (model id, tokens, cost, latency) that goes into the decision record, so cost per load is a query over the log. Failures are typed: LLMUnavailableError (timeout, overload) and LLMInvalidOutputError (schema validation failed). The pipeline turns either into an urgent Hero task; it never parses leniently or guesses. Free text appears only in message bodies.
Model strategy. Triage is a classification with typed outputs consumed by code, so it runs on a decision model (typed questions over the message and load state, calibrated probabilities, thresholds in code; built, ADR 0005), with a small generative model with structured output as the rejected alternative. Message text comes from templates (ADR 0004). A generative model is used only for the planner (no carrier-facing actions (a Hero task, a TMS note, or a Slack post to the broker)), template-less drafts and SOP compilation; vision for documents. Every call has a per-call timeout and bounded retries from LLMSettings. The real clients are built (ADR 0005): triage on the TypeSafe decision model with deterministic value extraction, the rest on DeepSeek, composed by build_llm_client from LOAD_AGENT_LLM_PROVIDER=typesafe+deepseek; evals use a scripted FakeLLM, so the default test run needs no key. Alternatives: one generative model for every task (simplest, but the cost per message and uncalibrated confidence do not fit the funnel); a single provider-specific SDK in the pipeline (faster to ship, locks the pipeline to one vendor and breaks offline evals). Trade-off: a decision model needs labeled data and calibration before launch. A 120 s per-run deadline spanning conflict re-decides bounds the model calls of a run; after it the run commits an urgent Hero task.
Alternatives considered
Section titled “Alternatives considered”| Alternative | Why not |
|---|---|
| Single ReAct agent with tools, reading full history each event | Simplest to build. Cost grows with history length on every event (does not fit $1 per load at 50 to 100 messages plus pings). Attempt counting and “is this still open” rely on model recall, which is where errors appear. Hard to test step by step. |
| Pure workflow engine, no LLM | Fully predictable and cheap. Cannot classify “What’s the PO number?” vs an ETA vs “ok”, cannot draft messages, cannot read documents. |
| Multi-agent (one agent per role or per channel) | More moving parts and hand-offs for a problem whose core is one state machine per load. Cross-channel resolution (ETA by email resolves an SMS ask) needs shared state anyway. |
| Memory as chat history or vector store of past messages | Answers “what did we ask?” probabilistically. Replay is not exact. Our log plus projection is exact and cheap to query. |
| Mutable state table updated in place (no log) | Simpler storage. Loses audit trail and replay, which evals and incident review depend on. |
Trade-offs
Section titled “Trade-offs”- Gains: most events end in the deterministic filter or SOP executor in milliseconds and cost nothing; behavior is reproducible from the log; every action is explainable by its decision record.
- Costs: more upfront modeling (events, asks, transitions); a situation that no SOP and no triage category covers goes to the planner or to a Hero, so coverage depends on the SOP corpus; schema evolution of logged entries needs care (stored models must stay readable); the outbox adds a dispatcher process and a few seconds of delivery latency (well inside the 5 minute budget).
Consequences
Section titled “Consequences”- Scenario evals replay a timeline with
FakeClockthrough the same worker, dispatcher and timer pump code as production, and assert the decision of each step, including explicit noops. process_event(event, deps) -> DecisionRecordis the single entry point the worker and evals drive; it never raises for business situations or LLM failures.- The
EventStore,OutboxandTimerStorecontract tests run against SQLite in CI and against Postgres in the production pipeline. - Changing the idempotency key format or the step counter semantics is a breaking change for stored logs.
- Durable timers, the outbox and per-load concurrency are specified in ADR 0003.