System design
The full system for the load agent: what is built in this repository, what is designed only, and why. Answers to the assignment questions are in answers.md; each decision record is in docs/adr/.
1. Scope: built versus designed only
Section titled “1. Scope: built versus designed only”Everything marked “built” is in the code and exercised by uv run pytest (offline, no API key; the scenario evals run against the real runtime; the Postgres contract tests run when LOAD_AGENT_TEST_PG_URL is set, as in CI). Everything marked “designed” is described here and not implemented.
| Area | Status | Where |
|---|---|---|
| Domain contracts (events, state, asks, actions, decisions, SOP spec), all frozen Pydantic models | built | src/load_agent/domain/ |
| Per-load append-only event log, state as a projection, open asks | built | store/projection.py, ADR 0001 |
| Event store with optimistic concurrency: SQLite (local, evals) and Postgres (production), one contract test suite | built (Postgres tests run when LOAD_AGENT_TEST_PG_URL is set, as in CI) | store/sqlite.py, store/postgres.py |
| Transactional outbox and dispatcher with idempotency keys, claims, leases, backoff, retry budget | built | store/outbox.py, agent/dispatcher.py, ADR 0003 |
| Durable timers (store plus pump), idempotent firing | built | store/timers.py, agent/timer_pump.py |
Layered SOP repository, routing index, compiled specs, sop preview with behavior diff | built | sop/, sops/ |
SOP compiler: validation rules and an LLM-backed compiler behind SopExtractionClient | built; the LLM call itself needs a real client, so the shipped *.compiled.yaml files are committed artifacts | sop/compiler.py, sop/llm_compiler.py |
| Layered message templates | built | messages/templates.py, ADR 0004 |
Human overrides: AskOverride event closes asks by hand (Hero tooling, broker Slack integration) | built; the Hero tooling and the integration that emit it are designed | section 5 |
| Pipeline: filter, trigger diff, routing, triage, SOP executor, planner, validation, commit | built | agent/ |
| Action validation (load contacts, channels, forbidden roles, duplicate keys) | built | agent/validation.py |
| Document handling: classify, extract, cross-check PO and BOL numbers, escalate | built against the LLMClient.extract_document seam, with a DeepSeek vision client (see the vision row below) | agent/reactions.py |
FakeLLM, scripted per scenario; LLMClient protocol, typed failures, settings from the environment | built | llm/ |
| Real provider clients | built, optional: triage on the TypeSafe decision model plus deterministic value extraction; planner, SOP compile, documents and fallback drafts on DeepSeek. Offline evals still use FakeLLM. | llm/typesafe.py, llm/extract.py, llm/deepseek.py, llm/settings.py |
Scenario evals (30), replay with FakeClock through the real runtime | built | evals/ |
| CI: ruff, ruff format, mypy strict, pytest with a Postgres service; a separate workflow runs the live LLM tests daily and on demand, outside the per-push CI | built | .github/workflows/ci.yml, .github/workflows/live.yml |
| Ingestion API and normalizer (FastAPI webhooks for SMS, email, Slack, ping, stage, documents; read endpoints; health and readiness) | built; auth is one shared-secret header, provider signature verification (Twilio, Slack signing secret) is designed only | api/, section 3 |
Queue partitioned by load_id | built on Postgres (FOR UPDATE SKIP LOCKED, one in-flight message per load) behind the ServiceQueue protocol, with an in-memory queue for tests; an SQS FIFO adapter is designed only | store/queue_pg.py, agent/queue.py |
Worker process (queue consumer, dispatcher and timer pump on intervals, SIGTERM drain) and docker compose up (Postgres, API, worker) | built; the Service workflow (.github/workflows/service.yml) builds the image, runs docker compose up and scripts/demo.sh and asserts the result on every relevant change; autoscaling is designed only | service/, Dockerfile, docker-compose.yml |
| Real integration services (SMS, email, Slack, TMS, Hero tasks, load status) | designed only; typed clients and in-process fakes are built | channels/ |
| SOP editor UI, publish flow, shadow mode, canary | designed only; sop preview, stale_sources, the SopPublisher protocol are built | section 8 |
| Observability stack (metrics, traces, dashboards, alerts) | designed only; decision records carry the raw data (path, SOP version, LLM usage with cost and latency) | section 12 |
| Per-run deadline (120 s across conflict re-decides) and run latency in each decision record | built | agent/deadline.py, llm/budget.py, section 13 |
| Aggregation and export of the run latency metric | designed only | sections 12 and 13 |
| Document extraction with a vision model | built against deepseek-flash (with the seam, schemas and PO/BOL cross-validation), smoke-tested live only on a blank image; per-type field schemas and accuracy on real BOLs and PODs are designed or unverified | sections 10 and 14 |
| Decision model for triage (typed questions, calibrated probabilities) | built (TypeSafe Jev, ADR 0005); calibrated on 190 recorded phrasings, with no production replay set yet | section 9 |
2. Requirements that shape the design
Section titled “2. Requirements that shape the design”| Requirement | Number | Consequence |
|---|---|---|
| Volume | 100,000 loads per month, 50 to 100 inbound messages each: 5M to 10M messages per month, plus pings | Most events must finish without a model call |
| Cost | about $1 of inference per load, about $0.013 per inbound message if every message called a model | Funnel: deterministic first, small model second, strong model rarely |
| Latency | event to last action within 5 minutes | Serial per load, short prompts, deadlines on every external call |
| Load life | 2 days to 2 weeks | Timers and state must survive deploys; SOP edits happen mid-load |
| SOP editors | non-technical | Markdown in, compiled spec out, with a plain-English preview |
| Unknown situations | “ok” with no question must do nothing | Noop is an explicit, tested decision; unknown means a Hero, not improvisation |
Load metadata
Section titled “Load metadata”The metadata follows the assignment’s example, with three stated changes:
| Change | Why |
|---|---|
load_type (dry_van, reefer, …) is an explicit field, defaulting to dry_van when the feed omits it | SOP routing for client + load_type must not depend on parsing the free-text commodity. Alternative: infer it from commodity with a model, which adds a call and an error source to every load |
reference_numbers (po_number, bol_number, pickup_number) is an optional structured field | Answering “What’s the PO number?” and cross-checking documents must not depend on reading prose |
broker.id is the client id, and fixtures use client_a instead of the example’s northline | The broker is our client, and the SOP directory is keyed by it. A real deployment uses the broker’s stable id; the example value is only an example |
When reference_numbers is absent, the PO number is parsed deterministically from instructions at registration (built): “PO 55821 required at gate” yields po_number 55821, which is what the 07:13 answer uses. The parse is a pattern match in code, not a model call. If no PO can be found, the agent never invents one: a question about the PO becomes a Hero task (eval po_question_without_po_number_becomes_hero_task). Alternative: ask the model to extract it from instructions; it can hallucinate a plausible number, which is the worst failure for a gate requirement.
3. Architecture
Section titled “3. Architecture”flowchart LR
subgraph Ext["Integration services (other team)"]
SMSs["SMS"]
EMs["Email"]
SLs["Slack"]
TMSs["TMS"]
LOC["Location and risk service"]
HEROs["Hero tasks and internal Slack"]
STS["Load status"]
end
subgraph Ingest["Ingestion (built)"]
API["Ingestion API"]
NORM["Normalizer: InboundEvent, event_id"]
end
SMSs --> API
EMs --> API
SLs --> API
LOC --> API
TMSs -->|"stop status"| API
GEO["Geofence"] -->|"stop status"| API
HEROs -->|"Hero stop status"| API
API --> NORM
NORM --> Q["Load queue: FIFO, group = load_id"]
PUMP["Timer pump"] -->|"TimerFired"| Q
subgraph Worker["Load worker (built)"]
RUN["process_event: filter, triggers, routing, triage, SOP executor, planner, validate"]
end
Q --> RUN
RUN <-->|"read snapshot, append at expected_version"| DB[("Postgres: event log, outbox, timers")]
RUN -.->|"triage, draft, plan, vision (optional)"| LLM["TypeSafe (triage), DeepSeek (rest)"]
RUN -->|"read"| SOPS[("SOP tree: Markdown, compiled specs, templates")]
DB --> DISP["Dispatcher: claims outbox rows"]
DISP -->|"idempotency key"| SMSs
DISP -->|"idempotency key"| EMs
DISP -->|"idempotency key"| SLs
DISP -->|"idempotency key"| TMSs
DISP -->|"idempotency key"| HEROs
DISP -->|"idempotency key"| STS
DISP -->|"schedule, cancel"| DB
DB --> PUMP
DISP -->|"ActionOutcome"| Q
subgraph Edit["SOP authoring (designed)"]
ED["Editor, compile, preview, evals, shadow, publish"]
end
ED --> SOPS
Components:
| Component | Job | State |
|---|---|---|
| Ingestion API and normalizer (built) | Authenticate the caller, map a payload to a typed InboundEvent with a deterministic event_id (provider message id, else a hash of the content and the provider’s own timestamp), reject unknown loads (404) and contacts not on the load (422), answer a retried webhook with 200 “duplicate”, enqueue and return 202 | none |
| Load queue (built on Postgres; SQS FIFO designed) | Order per load_id, at-least-once delivery, one message in flight per load | load_queue rows |
Load worker (built; load-agent worker) | One run per event, optionally under a per-load lease (LoadLock, an optimization against wasted LLM calls; the version check is what guarantees correctness). Read log, project state, decide, commit decision plus outbox in one transaction | none; all state is in the store |
| Event store (built) | Per-load log with expected_version, outbox rows in the same transaction | Postgres (partitioning designed, section 13) |
| Dispatcher (built) | Claim due outbox rows, call the owning service with the idempotency key, publish ActionOutcome | outbox row status |
| Timer pump (built) | Turn due timers into TimerFired events | timer rows |
| SOP tree (built) | Markdown, compiled specs, templates, versioned | files in git, or an object store plus index |
LLM providers (built, optional; evals use FakeLLM) | TypeSafe decision model for triage; DeepSeek for the planner, drafts for uncovered intents, documents (vision) and SOP compile | none |
Alternatives to the hybrid (full table in ADR 0001):
| Alternative | Trade-off |
|---|---|
| One ReAct agent with tools reading the full history per event | Simplest. Cost grows with history on every event, attempt counting rests on model recall, hard to test step by step. Does not fit $1 per load. |
| Workflow engine only (Temporal) | Durable and cheap, but cannot read “What’s the PO number?” or draft a message. Temporal fits as the timer and retry substrate; ADR 0003 explains why we kept the log and timers in Postgres for now. |
| Multi-agent | More hand-offs for a problem whose core is one state machine per load. Cross-channel resolution needs shared state anyway. |
Why Postgres carries the log, outbox and timers: the commit of an event, its decision and its actions must be one transaction (ADR 0003). Alternative: a log in Kafka and state elsewhere, which splits that transaction. Cost: Postgres needs partitioning and retention at 10x (section 13).
Running as a service (built)
Section titled “Running as a service (built)”docker compose up --build starts Postgres 16 (healthcheck), the API and the worker from one image. API and worker run load-agent migrate first (idempotent, serialized by an advisory lock). The pieces reuse the pipeline unchanged: process_event, the outbox dispatcher and fire_due_timers are the same functions the evals call, wired to Postgres stores and a system clock in service/wiring.py.
- The handler does not append to the log.
process_eventcommits[event, decision]in one transaction and treats an event already in the log as processed, so an event pre-appended by the API would be dropped as a duplicate. The API instead checks “already in the log or already queued” and publishes; the worker stays the only writer of a load’s log (ADR 0003). Alternative: append in the handler and have the worker run on the logged event. That makes the 202 a durable log write but splits the commit of event and decision, which the outbox design avoids. - Idempotency. The event id is the provider message id verbatim, or a hash of load, kind, content and the provider’s timestamp (never the server’s clock), so a retry maps to the same id. The queue has a unique
(load_id, event_id)and keeps processed rows, so a retry after processing is also recognised. Trade-off: two byte-identical texts without a provider id or timestamp are one event. - Per-load serialization with several workers. A message is claimable only if it is the oldest open message of its load (order is by commit visibility: the sequence the database assigns, not the sender’s clock); claims use
FOR UPDATE SKIP LOCKEDand a lease (6 minutes, above the run budget) so a crashed worker’s message is redelivered. A failing message is released with backoff and dead-lettered after 5 attempts; later messages of that load then proceed. The version check in the event store remains what guarantees correctness. - Worker loop. Per tick: pump due timers, claim and process messages, dispatch the outbox. On SIGTERM the loop finishes the message in flight, takes no further message, and exits; the stop flag is checked between messages. Ack and release are fenced by the claim’s attempt count, so a worker whose lease expired cannot complete or requeue a message another worker now holds.
- Scripted fake LLM by default.
LOAD_AGENT_LLM_PROVIDERdefaults tofake;LOAD_AGENT_FAKE_SCRIPTpoints at a scenario file whosellm:block answers the triage calls, so the load #481207 demo runs offline. The image contains no keys;typesafe+deepseekreads them from the environment. - Endpoints.
POST /v1/loads,POST /v1/loads/{id}/events/{sms,email,slack,ping,stage,documents},GET /v1/loads/{id}/{state,events,decisions},GET /healthz,GET /readyz(database reachable and every migration applied). Bodies overLOAD_AGENT_MAX_BODY_BYTESget 413.
4. Per-event pipeline
Section titled “4. Per-event pipeline”flowchart TD
E["Inbound event from the queue"] --> R["Read log snapshot and version"]
R --> P["Project LoadState: stage, delay risk per stop, open asks, committed action keys"]
P --> DUP{"event_id already in log?"}
DUP -->|"yes"| N1["Noop: duplicate, not appended"]
DUP -->|"no"| K{"Event kind"}
K -->|"ping or stop status"| TR["Trigger diff: state before vs after"]
K -->|"timer"| TM{"Ask open and attempt equals step?"}
K -->|"sms, email, slack"| SND{"Sender is a contact of the load?"}
K -->|"attachment"| DOC["Extract, cross-check, confidence"]
K -->|"action outcome"| OUT{"Failed external send?"}
TR -->|"no band or stage change"| N2["Noop: filtered"]
TR -->|"trigger"| RT["Routing: client + load type, client, standard"]
RT -->|"compiled SOP"| EX["SOP executor: open ask, send, set timer"]
RT -->|"no topic"| H0["Hero task"]
TM -->|"no"| N3["Noop: stale timer"]
TM -->|"yes"| CT["Continue ask under its pinned SOP version"]
CT --> EX
SND -->|"no"| H1["Hero task"]
SND -->|"yes"| ACK{"Bare acknowledgement and no ask can be answered by it?"}
ACK -->|"yes"| N4["Noop: no LLM call"]
ACK -->|"no"| TG["Triage: one decision-model call"]
TG -->|"unclear or low confidence"| H2["Hero task"]
TG --> ST["Resolve every open ask whose resolution condition the extracted values satisfy"]
ST --> QA{"Question?"}
QA -->|"load fact exists"| ANS["Templated answer from metadata"]
QA -->|"no fact"| PL["Planner: no carrier-facing actions"]
DOC -->|"unknown, low confidence, mismatch"| H3["Hero task"]
DOC -->|"trusted"| RT
OUT -->|"yes"| H4["Hero task"]
OUT -->|"no"| N5["Noop: bookkeeping"]
EX --> V["Validate actions"]
ANS --> V
PL --> V
ST --> V
V -->|"rejected"| H5["Replace by Hero task"]
V --> C["ONE transaction: event, DecisionRecorded, outbox rows at expected_version"]
H5 --> C
C --> D["Dispatcher sends with idempotency keys"]
Rules enforced in code:
- Triggers come from a diff of projected state (
agent/triggers.py), not from event kinds. A ping into the band the stop is already in is not a trigger; a ping older than the last applied one changes nothing. - Only stop status events move the stage (geofence, TMS, a Hero). The agent’s own
update_stagemirrors the stage to other systems and never changes state, so a wrong mirror cannot fire SOPs. - The model answers typed questions about meaning; code owns values and actions. The TypeSafe Jev decision model (ADR 0005) returns probabilities for closed questions and never text or values. Code finds and normalizes value candidates (
llm/extract.py) and the model only confirms one. A DeepSeek planner proposes actions for situations no SOP covers, never carrier-facing (only a Hero task, a TMS note, or a Slack post to the broker), validated like every other action. Alternative: a generative model that reads the message and returns the value (invents values, uncalibrated confidence). - Resolution is a function of state, not of the model’s category label. The model extracts values;
satisfies(condition, extracted, sender_role, channel)decides if an ask closes. An ETA from the dispatcher by email closes the driver’s SMS ask. - Failure never crashes the run.
LLMUnavailableError,LLMInvalidOutputError, an unsafe draft, a SOP that cannot be executed and a rejected action all become an urgent Hero task carrying the event. AFakeLLMcall that no scenario scripted is deliberately not caught, so evals see a hidden model call. - A forbidden role is a rule of the client, not of one ask. The roles that the client’s own SOPs (client and client + load type layers) forbid are applied to every action of every run for the load (
SopRepository.forbidden_roles_for, enforced inrunner._process_oncenext to the roles of the asks’ pinned SOPs). Client B forbids contacting the dispatcher, so a dispatcher’s PO question gets the same treatment with or without an open ask: a Hero task carrying the answer, never a reply. The standard layer’s “do not message the driver for this event” stays scoped to its topic. Alternative: forbid only while an ask that forbids it is open (the first implementation): the dispatcher was answered at 06:30 and refused at 07:10. Trade-off: a client SOP that forbids a role in one topic forbids it everywhere; an editor who wants it narrower writes the restriction in the standard layer’s way (the topic’s own steps). - Replies are not repeated and need a safe place to go. The same answer (intent, text and contact) inside 10 minutes is not sent twice (the triage call is still spent, the SMS is not), and an email outside the main thread is escalated with the answer because the agent only replies in the main thread. One Hero task per distinct message from unknown senders per 30 minutes: the event carries no raw sender identity, so the key is the message text and the window is per load (a repeat of the same text adds nothing, new content still reaches a Hero). A repeat is not suppressed when the first copy’s delivery failed permanently (
LoadState.failed_action_keys). Both rest onLoadState.recent_actions, a bounded projection (last 20 sent messages and Hero tasks) of committed decisions. SOP requests (request_eta, repeated by design) are never suppressed. - Concurrency conflict means re-read, re-project, re-decide, with deterministic ask ids and keys so the second decision is the same decision.
Alternative to the diff-based trigger: let the model decide whether an event is “interesting”. Cheaper to write, but it puts an LLM call on every ping (the highest-volume input) and makes the “high stays high” case probabilistic.
5. Memory: log, projection, open asks, timers
Section titled “5. Memory: log, projection, open asks, timers”Full definitions are in ADR 0001 and 0003. Summary:
- Every input and every decision is appended to the per-load log.
LoadStateis a pure fold over it. Replaying the log gives the same state, which is what lets evals and incident review be exact. - An open ask holds: id, type, target contact, pinned SOP version, attempts, escalations done, next timer,
resolved_by. “Is the ETA request still open?” isstate.open_asks_of_type("eta_request"). No model reads chat history for this. - Step counter: attempts and escalations share one counter. An action carries
ask_idandattemptonly when it takes that step, which gives a deterministic idempotency key{load_id}/{ask_id}/{attempt}/{kind}. The PO reply at 07:13 takes no step and is keyed by its event. - Timers carry
{load_id, ask_id, attempt}. The agent sees aTimerFiredevent, projects the current state and checksask.is_timer_current(payload). If the ask is closed or already past that attempt, it is a noop with no model call.
Alternatives: chat history or a vector store of past messages (probabilistic, replay not exact), a mutable state row (no audit, no replay). Cost of ours: stored models must stay readable across schema changes, and a log read per run (about 1,300 rows per load with pings at the assumed rate; a snapshot row every 200 records is the designed optimization, needed beyond about 2,000 rows; storing only band-changing pings in the log and the rest in a telemetry table is the alternative, designed).
stateDiagram-v2 [*] --> Open: AskOpened, attempts 1, timer set Open --> Open: AskStepTaken, repeat request or escalation step Open --> Resolved: message satisfies the resolution condition Open --> Cancelled: arrived at the stop, load delivered, moved past the stop, risk low at the timer, or a human override Open --> Open: timer stale or duplicate, noop Resolved --> [*] Cancelled --> [*]
An ask stays open after its escalation chain ends (no further timers), so a late ETA still resolves it. Transitions are carried inside decision records only.
Time zones
Section titled “Time zones”Appointments in the metadata are naive local times at the stop. They are localized with the stop’s timezone at the edge (domain/load.py) and stored as UTC; nothing downstream sees a naive datetime. Triage receives the next stop’s timezone so “around 9:30” or “in 2 hours” becomes an aware ETA. Rendering converts back to the stop’s zone (the 481207 output shows PDT). Timers are UTC instants. Alternative: store local times and convert on use, which breaks across DST and multi-timezone loads.
Multi-stop loads
Section titled “Multi-stop loads”Delay risk is tracked per stop, and StageEntered fires on a change of stage or of the current stop, so the middle stages repeat for each stop (a move between delivery stop 2 and 3 has from_stage == stage). An ask records the stop it is about (stop_sequence): arriving at the pickup closes the pickup ETA ask and leaves a delivery ETA ask open. The NoOpenAsk condition is scoped to the stop the trigger is about, so a pickup ETA ask does not stop the delivery ETA ask when the delivery stop turns high (eval delivery_risk_high_while_pickup_eta_ask_is_open_asks_for_the_delivery_eta). Each ask carries its own stop in the triage request (AskSummary.stop_sequence, stop_name, timezone); the request’s next_stop_* fields describe the stop the open asks are about when they agree, else the current stop. When asks of one type are open for several stops, a reply resolves only the one for the lowest stop, the stop the truck reaches first (eta_reply_resolves_only_the_nearest_stop_ask): one “ETA 9:30” must not close a request nobody answered. A reply that names a stop (“delivery ETA 9:30”, a stop name) resolves the asks of that stop and is read in its timezone (named_stop_eta_resolves_only_that_stops_ask). Questions about an address or an appointment pick the stop from the words of the question (a stop name, or pickup and delivery words); no cue means the stop the load is heading to or at, and an unclear question (both stops, several matches) goes to a Hero (address_question_names_the_stop_it_is_about). Built, with tests in tests/agent/test_pipeline.py and tests/agent/test_triggers.py.
Broker Slack messages and human overrides
Section titled “Broker Slack messages and human overrides”A Slack message from a broker contact takes the same path as an SMS: sender check, acknowledgement filter, triage, resolution against open asks, a templated reply in the broker channel (built). Human overrides are built in two steps. A Hero or the TMS can send a stop_status event that moves the stage, and the next run sees the new state. To stop a chase, a broker’s free text never cancels anything: when triage marks a broker message as new info while an ask is open (and the message answers nothing), the agent creates a Hero task “broker message may affect open ask X” with the text, and the follow-ups continue. A Hero (or the broker Slack integration, for an explicit command) then emits an AskOverride event: the open asks named by ask_id, or by ask_type and optionally stop_sequence, are cancelled with their timers, deterministically and with no model call; a scope that matches nothing is a noop (eval broker_stand_down_cancels_the_eta_chain). Alternative: let the model cancel on “stand down”. It saves a human step, but a misread “I’ll call the driver” would silence a real chase, and the cost is only the delay until the Hero acts. Designed: a Hero marking a task “agent was wrong”, and a per-client “pause the agent on this load” flag that the filter checks first. Alternative to the override event: let humans edit state directly, which breaks replay.
Delay risk falling back to low closes the ask, but not on the ping. An ask opened because a stop became high is closed when its own follow-up timer fires and finds that stop low (risk_falls_back_cancels_the_eta_chain): no second text, no escalation, no event is needed besides the timer that already exists. Medium does not close it. Why not on the ping: a ping back to low is often a blip, and closing on it would let the next high open a fresh ask at attempt 1, texting the driver again inside the retry interval and restarting the escalation clock. Because the ask stays open until its timer, high, low, high (or high, medium, high) within minutes changes nothing: one text per 30 minutes and the escalation on schedule (risk_flapping_through_low_keeps_the_chain_on_schedule, risk_flapping_through_medium_keeps_the_chain_on_schedule). Alternatives: cancel on the ping and carry attempts forward when the same stop re-opens within a cooldown (needs a reopen transition and re-arming a timer id the scheduler already used), or cancel on the ping with no memory (the flapping spam above). Cost of ours: between the low ping and the timer the ask is still open, so an ETA reply in that window still resolves it, which is harmless. The low reading must be recent, recorded after the ask’s last text was sent; an older one leaves the chain running. Also: when the load moves past a stop (any stage change, including DEPARTED with no ARRIVED before it), the open asks about earlier stops are cancelled, and replies are never matched to them (departing_without_arriving_closes_the_pickup_ask).
When an ETA resolves an ask, the ETA is stored on the ask (AskResolved.eta, OpenAsk.resolved_eta, optional so older logs replay) and a templated TMS note “ETA to {stop_name}: {eta}” records it in the stop’s time zone. An ETA resolves an ask if it is not more than 15 minutes in the past (ETA_PAST_GRACE, so “there at 9:30” said at 9:31 counts). An older ETA never resolves it: a TMS note records it and a non-urgent Hero task asks a person to check, and the ask stays open. Broker news while an ask is open gets both the TMS note and the Hero task; the Hero task is non-urgent (the urgency table says so) and the same message creates one task per 30 minutes (a different message still creates one).
6. Load 481207, 07:00 to 08:02
Section titled “6. Load 481207, 07:00 to 08:02”Client A, delay risk high. Ask id ask:eta_request:ping_0700.
sequenceDiagram autonumber participant L as Location service participant Q as Load queue participant W as Load worker participant S as Store, log and outbox and timers participant D as Outbox dispatcher participant P as Timer pump participant SMS as SMS service participant EM as Email service participant H as Hero tasks participant DR as Driver participant DS as Carrier dispatcher L->>Q: 07:00 ping, risk low to high Q->>W: ping_0700 W->>S: read, project: trigger DelayRiskBecame(high, stop 1) Note over W: route client_a/delay_risk_high, no open eta_request, no LLM W->>S: 07:01 commit decision: ask opened, outbox send_sms attempt 1, set_timer 07:31 D->>S: claim outbox D->>SMS: send_sms "what is your ETA to Riverbend plant" D->>S: schedule timer tmr:...:1 at 07:31 SMS->>DR: SMS DR->>Q: 07:12 "What's the PO number?" Q->>W: sms_0712 Note over W: not an ack, one triage call: question, topic po_number, no ETA W->>S: 07:13 commit: send_sms "the PO number is 55821" (from reference_numbers), no ask change D->>SMS: send_sms (keyed by event, not by ask) Note over S: ask still open, attempts 1, timer 07:31 untouched P->>S: 07:31 due timers P->>Q: TimerFired attempt 1 Q->>W: timer Note over W: ask open and attempt 1 equals step 1, no LLM W->>S: commit: send_sms attempt 2, set_timer 08:01 D->>SMS: send_sms again P->>Q: 08:01 TimerFired attempt 2 Q->>W: timer Note over W: ask open, step 2, max attempts reached, escalation chain W->>S: 08:01 commit: send_email (main thread), create_hero_task urgent D->>EM: email dispatcher in the main thread EM->>DS: email D->>H: create task: call the driver
The 08:01 timer fires and the actions are committed in the same run, so the email and the task go out at 08:02 or within seconds of it (the assignment’s 08:02; dispatch adds a few seconds). The two possible interruptions are tested: an ETA arriving from the dispatcher at 07:20 (dispatcher_email_eta_resolves_driver_ask), after which the 07:31 timer is a noop (timer_after_eta_received_is_noop); and the driver arriving before the follow-up (driver_arrives_before_follow_up_cancels_eta_ask).
Real output of the scenario (uv run load-agent run evals/scenarios/client_a_load_481207_full_timeline.yaml; times are 12-hour with the stop’s zone, as in every text the agent produces). The “$0” is the FakeLLM price; the single model call in the whole timeline classifies the PO question. Messages are templated, so the other steps make zero calls.
Scenario client_a_load_481207_full_timeline: Delay risk high; ETA asked twice; driver asks for the PO; escalate to dispatcher and Hero.
7:01 AM ping turns delay risk high: SMS the driver for the ETA and set the 7:31 AM follow-up event ping_0700: sop [client_a/delay_risk_high@v1] why: followed client_a/delay_risk_high@v1 -> send_sms to drv_20931: "Load 481207: what is your ETA to Riverbend plant? Please reply with the time." -> set_timer tmr:481207:ask:eta_request:ping_0700:1 at Oct 13, 7:31 AM PDT (2 delivery outcome(s) recorded, nothing to do)
7:13 AM driver asks for the PO: answer it from metadata, the ETA ask stays open event sms_0712: triage, 1 LLM call(s), $0 why: answered po_number from load data -> send_sms to drv_20931: "Load 481207: the PO number is 55821." (1 delivery outcome(s) recorded, nothing to do)
7:31 AM timer: no ETA yet, ask again and set the 8:01 AM follow-up event evt:tmr:481207:ask:eta_request:ping_0700:1: sop [client_a/delay_risk_high@v1] why: ask:eta_request:ping_0700 unanswered at step 1, continued client_a/delay_risk_high@v1 -> send_sms to drv_20931: "Load 481207: what is your ETA to Riverbend plant? Please reply with the time." -> set_timer tmr:481207:ask:eta_request:ping_0700:2 at Oct 13, 8:01 AM PDT (2 delivery outcome(s) recorded, nothing to do)
8:01 AM timer: two requests unanswered, email the dispatcher in the main thread and create an urgent Hero task event evt:tmr:481207:ask:eta_request:ping_0700:2: sop [client_a/delay_risk_high@v1] why: ask:eta_request:ping_0700 unanswered at step 2, continued client_a/delay_risk_high@v1 -> send_email to dsp_4410 (main thread): "The driver on load 481207 has not answered our texts. Please send the ETA to Riverbend plant." -> create_hero_task (urgent Load 481207: Call driver for eta): "Call the driver on load 481207 to get the ETA to Riverbend plant. No ETA arrived after our requests." (2 delivery outcome(s) recorded, nothing to do)
8:30 AM nothing is left to do: the escalation chain is finished no event reached the agent (no timer was due)
Expectations: all steps match the scenario file.7. SOPs: organization, retrieval, compilation
Section titled “7. SOPs: organization, retrieval, compilation”Detailed in ADR 0002. Summary:
sops/standard/delay_risk_high.md # applies to every clientsops/clients/client_a/delay_risk_high.md # overrides standard for client_asops/clients/client_b/reefer/departed_pickup.md # client + load type- Layering: per topic,
client + load_typeoverclientoverstandard. The winning file replaces the topic whole. No section merge. - Compilation at edit time: each Markdown file becomes a
CompiledSop(trigger, conditions, ask withmax_attempts,retry_afterand resolution condition, ordered escalation steps, forbidden contact roles) plus verbatim guidance snippets that are never executed. Exact things become structure; judgment stays text. Ambiguity is rejected with a reason, not defaulted (“if the driver doesn’t answer soon” fails asambiguous_step). - Retrieval is deterministic: the compiler builds a routing index from trigger to topics;
SopRepository.topics_for(trigger, client_id, load_type)is a lookup. No vector search on the load path. - Unmapped situation: triage, then the planner (no carrier-facing actions (a Hero task, a TMS note, or a Slack post to the broker)), whose conservative default is a Hero task.
- Version pinning: each publish is an immutable version; an open ask keeps the
SopRef(with version) it was opened under. An ask started under v1 finishes under v1. - Messages: a compiled step stores an
intent, never a body. Wording comes from layered templates (ADR 0004).
Trade-offs of the method: predictable, cheap, testable, and a client override cannot be shadowed by a retrieved chunk; but a new situation needs a new topic and a compile, and whole-topic override duplicates text when a client changes one line. Raw-Markdown-to-LLM on every event was rejected (cost, nondeterminism, retry counting by the model); ADR 0002 has the full table.
8. SOP editing workflow (designed; the preview is built)
Section titled “8. SOP editing workflow (designed; the preview is built)”People without a technical background edit SOPs, so safety comes from the pipeline, not from the editor’s care.
flowchart TD
A["Editor changes Markdown (draft)"] --> B["Compile: strong model extracts the procedure"]
B --> C{"Validation: schema, ambiguity, forbidden roles, grounding in the text"}
C -->|"rejected"| A
C -->|"ok"| D["Plain-English preview, rendered from the spec by code"]
D --> E["Behavior diff: this client's and the standard scenario evals on old vs new version"]
E --> F{"Editor approves what the agent understood and the diff"}
F -->|"no"| A
F -->|"yes"| G["Eval suites: standard + this client, must pass"]
G -->|"fail"| A
G -->|"pass"| H["Shadow mode on live loads: decide, do not act, compare with current version"]
H -->|"unexpected differences"| A
H -->|"clean for N loads or N hours"| I["Publish: immutable version N+1, stale publish refused"]
I --> J["New triggers use N+1; open asks stay pinned to the version they started under"]
I --> K["Rollback: republish the old spec as N+2, one click"]
What is built: the preview renderer and behavior diff (sop/preview.py), the sop preview command, the compile validation rules, FileSopRepository.stale_sources() (flags a Markdown file edited without a recompile) and the SopPublisher protocol. What is designed only: the editor UI, the diff against scenario results, shadow mode and the publish service.
Screen mockup (the left side is the editor’s own words, the right side is what the runtime will do):
+-----------------------------------------+-----------------------------------------------------+| client_a / delay_risk_high draft v2| What the agent understood [compile ok]|+-----------------------------------------+-----------------------------------------------------+| # Delay risk high | When delay risk becomes high: || | Only if no ETA request is already open. || ## Procedure | 1. Text the driver (ask for the ETA). || 1. Text the driver and ask for the ETA. | 2. If no ETA within 30 minutes, repeat step 1 || 2. If there is no ETA after 30 minutes, | (attempt 2 of 2). || text again. | 3. If still no ETA 30 minutes later, in order: || 3. After two texts with no ETA, email | a. Email the dispatcher in the main thread. || the dispatcher in the main thread. | b. Create an urgent Hero task. || 4. Then create a Hero task. | Answered by an ETA from any contact, any channel. || +-----------------------------------------------------+| ## Tone | Behavior change vs published v1 || Keep texts short. Include the load | - Wait before the second text: 30 min (was 20) || number. | - 4 scenarios changed, 31 unchanged || | client_a_load_481207_full_timeline: 7:21 AM now || | 7:31 AM (follow-up SMS) || | Checks || | standard suite 9/9 client_a suite 5/6 (1 diff) || | open asks affected: 12 (stay on v1) || [Save draft] [Compile] | [Request shadow run] [Approve and publish] [Undo] |+-----------------------------------------+-----------------------------------------------------+The numbers in the mockup are illustrative. The wording on the right is the real output of the built renderer, which also shows the behavior diff between layers. The same function (render_behavior_diff) is what compares two versions of one topic.
Real output of uv run load-agent sop preview client_a delay_risk_high:
SOP for client client_a, topic delay_risk_highSource: client_a/delay_risk_high@v1 (client layer)
When delay risk becomes high: Only if no ETA request is already open. 1. Text the driver (ask for the ETA). Wait for the ETA. 2. If there is no ETA within 30 minutes, repeat step 1 (attempt 2 of 2). 3. If there is still no ETA 30 minutes after the last attempt, in order: a. Email the dispatcher in the main thread (say the driver does not answer and ask for the ETA). b. Create an urgent Hero task (call the driver to get the ETA). The request is answered as soon as we receive the ETA from any contact of the load on any channel. The remaining steps are then skipped.
Behavior change compared with standard/delay_risk_high@v1 (standard):- Attempts at the request: 2 (was 1).- Escalation: now also does: email the dispatcher in the main thread (say the driver does not answer and ask for the ETA).- Escalation: the Hero task is now urgent.- Guidance (affects wording only, not what happens): added 'Email to the dispatcher'; changed 'Tone'.Real output of uv run load-agent sop preview client_b delay_risk_high:
SOP for client client_b, topic delay_risk_highSource: client_b/delay_risk_high@v1 (client layer)
When delay risk becomes high: Only if no ETA request is already open. 1. Text the driver (ask for the ETA). Wait for the ETA. 2. If there is no ETA 30 minutes after the last attempt, in order: a. Post to the broker's Slack channel (alert that the ETA is missing). The request is answered as soon as we receive the ETA from any contact of the load on any channel. The remaining steps are then skipped. Never contact: dispatcher.
Behavior change compared with standard/delay_risk_high@v1 (standard):- Escalation: no longer does: create a Hero task (call the driver to get the ETA).- Escalation: now also does: post to the broker's Slack channel (alert that the ETA is missing).- Now never contacts the dispatcher.- Guidance (affects wording only, not what happens): changed 'Tone'.Client B’s preview shows the override: one text, then Slack to the broker, never the dispatcher. Client B reefer adds a topic at the client + load_type layer (sops/clients/client_b/reefer/departed_pickup.md): when the driver leaves the pickup, ask for the trailer temperature and the seal number.
Mid-load edits: an ask opened under v1 finishes under v1, because shortening a wait that already started, or changing who gets escalated in the middle of a chain, is worse than a late change. Alternative: apply the new version to open asks immediately. That is simpler to explain and surprises operators; a client who wants an urgent change can have a Hero cancel the open asks (designed: a Hero action to cancel an ask).
9. LLM strategy
Section titled “9. LLM strategy”Principle: a model is used only where language or an unlisted situation requires it, with the smallest model that passes its evals. All calls go through LLMClient (triage, draft, extract_document, plan, extract_sop), take a Pydantic request, return a validated Pydantic result plus LLMUsage (model id, tokens, cache tokens, cost, latency). Output that does not validate is LLMInvalidOutputError and becomes a Hero task; there is no lenient parsing.
| Task | Model class | When | Per-load calls (estimate) |
|---|---|---|---|
| Triage of a message | decision model (TypeSafe Jev, below) | inbound message that survives the filter | 30 to 75 |
| Message text | none: templates filled in code (ADR 0004). A generative model (DeepSeek) drafts only for an intent with no template, and every named fact must appear verbatim, else a Hero task | outbound messages | 0 to 1 |
| Planner | strong generative model | a question no load fact answers; a situation no SOP covers. May propose only create_hero_task, tms_note, or post_slack to the broker, at most 3 steps, confidence at least 0.8 | 0 to 2 |
| Document extraction | vision model | attachments | 2 to 6 |
| SOP compilation | strong generative model | SOP save, never on the load path | none per load |
Times in text. Every time shown to a person (messages, Hero tasks, TMS notes, the CLI trace) goes through messages/timefmt.py: US style, 12-hour with AM/PM and the stop’s zone abbreviation (“Oct 14, 9:30 AM PDT”, windows “Oct 13, 6:00 AM to 2:00 PM PDT”). Storage stays UTC. A test replays every scenario and fails if any text contains a 24-hour time (verbatim quotes of a sender’s own message excepted).
Triage as a decision model (built, ADR 0005). llm/typesafe.py calls the TypeSafe AI decision model (Jev) once per message with closed, typed questions built from the TriageRequest: a probability per open ask’s required field, one for “does it ask a question?”, two choices (kind of message; topic of a question), and, for each value candidate, a yes/no question plus guard questions. The model returns decisions only, never text or values. The model decides meaning, code decides syntax, and either being unsure escalates: llm/extract.py finds candidate spans (“9:30”, “in 45 min”, “1:30 out”, “-10F”, “seal 123456”) and normalizes each to a value or None (timezones, dates, ranges, ambiguous hours, DST and stale times are None); the model is asked per candidate “is this the time this truck will arrive at unclear and a Hero decides. More than three candidates is unreadable. Thresholds are named constants (PRESENT_MIN 0.60, ABSENT_MAX 0.30, GUARD_MAX 0.30 for ETA and temperature and 0.20 for seal, KIND_MIN_CONFIDENCE 0.60, TOPIC_MIN_CONFIDENCE 0.80). A kind that contradicts the other signals also yields unclear. A field judged present with no usable candidate raises UnparseableValueError and the run escalates: the agent never guesses a value. Cost: about 1,055 input and 260 output tokens per call with one ETA candidate. Why a decision model: calibrated probabilities per typed question feed deterministic code, it is cheaper and faster than a generative call, and thresholds are tunable from replays. Alternatives: a small generative model with structured output (can invent a value, uncalibrated confidence), and meaning encoded in code (a grammar plus deny lists), which three review rounds showed cannot keep up with how drivers write. Trade-off: more escalations (on the recorded calibration set, 82% to 84% of 74 labelled positives are accepted across two recordings (62/74 and 61/74), 19/23 and 17/23 on a held-out 30%, and 6 of 7 clean controls; the rest go to a Hero) in exchange for no wrong value on 116 adversarial negatives in two independent recordings of one model version (calibration, variance and what it does not show in ADR 0005; a held-out production set is required before launch); the decision model’s calibration should be checked against labeled replays before launch. The recorded model answers for 190 phrasings (two independent samples) are replayed offline in tests/llm/test_gate.py.
Generative client (built). llm/deepseek.py speaks the OpenAI-compatible chat completions API for plan, extract_sop, extract_document and the fallback draft. Each call sends the versioned system prompt plus the JSON schema of the expected output, requests a JSON object and validates it against the pydantic schema; invalid JSON, a schema mismatch, an empty or truncated completion, or a draft over max_chars is LLMInvalidOutputError. Default models: deepseek-flash for drafts and documents (it accepts images), deepseek-v4-pro for the planner and SOP compilation, all configurable. CompositeLLMClient routes triage to TypeSafe and the rest to DeepSeek; LOAD_AGENT_LLM_PROVIDER=typesafe+deepseek builds it from TYPESAFE_AI_API_KEY and DEEPSEEK_API_KEY. Both clients share llm/http.py: timeout, 2 retries on 429, 5xx and transport errors with capped backoff, 4xx (including a 403 with its detail) is LLMUnavailableError without retry, and the key never appears in an error or log line. Usage (tokens, latency, cost from configured prices, default 0) is returned in every LLMResult and lands in the decision record.
Run deadline (built). process_event starts a 120 s budget on the injected clock that spans every conflict re-decide. DeadlineLLM refuses a model call once it is spent (LLMDeadlineExceededError, an LLMUnavailableError) and, for a call that starts, hands the remaining seconds to the HTTP client through llm/budget.py: each attempt’s timeout is the smaller of the configured timeout and the time left, and a retry whose backoff would end past the deadline is not made. The run then commits an urgent Hero task carrying the event. A run that needs no model call is unaffected. The record carries run_latency_ms. The dispatch retry budget is 150 s, so decide plus deliver stays under 270 s of the 5 minute target. Residual: httpx applies the attempt timeout to each phase (connect, write, read), not to the whole exchange, so a server that trickles bytes can exceed it slightly; responses here are small non-streamed JSON bodies.
Status: evals use a scripted FakeLLM by default: each scenario file declares the model outputs it expects the pipeline to ask for, and an unscripted call fails the test. That makes “this step must not call a model” an assertion (llm_calls: 0). The real clients are covered offline with httpx.MockTransport and by pytest -m live smoke tests that skip unless both keys are set. The live tests run in .github/workflows/live.yml daily and on demand, outside the per-push CI, so a provider outage never blocks a merge.
Model ids, timeouts (30 s default) and retries (2) come from LLMSettings, not from call sites.
10. Documents: BOL, POD, rate confirmation, receipts
Section titled “10. Documents: BOL, POD, rate confirmation, receipts”Designed flow, with the built parts marked:
flowchart LR
A["Attachment event: file URI, media type, sender"] --> B["Fetch from object store"]
B --> C["Vision model: classify and extract with the schema of the type (designed)"]
C --> D["DocumentExtraction: type, fields, confidence (built)"]
D --> E{"Type known and confidence at least 0.7? (built)"}
E -->|"no"| H["Hero task with the file"]
E -->|"yes"| F["Cross-validate against load metadata (built for PO and BOL numbers)"]
F -->|"mismatch"| H
F -->|"match"| G["DocumentReceived trigger: SOP for the type, for example TMS note (built)"]
- Classify: one vision call returns
DocumentType(bol,pod,rate_confirmation,receipt,unknown). Unknown goes to a Hero. - Extract with a schema per type: BOL (BOL number, PO number, seal number, piece count, weight, signature present), POD (receiver name, delivery time, signature present, exceptions noted), rate confirmation (rate, stops, appointments, carrier MC number, load id), receipt (amount, vendor, date, type such as lumper). Today
DocumentExtraction.fieldsis a string map and the per-type schemas are designed. - Cross-validate against the load metadata: PO and BOL numbers are checked in code today; the designed set adds weight, load id, MC number, stop names and, for a rate confirmation, the rate against the TMS. The model extracts, code compares. A model asked “does this match?” can agree with a wrong document; an equality check cannot.
- Low confidence or mismatch: a Hero task naming the field and both values, and no SOP step runs. Alternative considered: auto-correct the load from the document. Rejected, because a misread digit would silently change a contract value.
- Cost: vision at about $0.01 to $0.02 per document, 2 to 6 documents per load. Photos are downscaled before the call (designed). Duplicate files (same hash) are dropped at ingestion (designed).
- Trade-off: a vision model reads poor phone photos worse than clean PDFs. Confidence gating sends those to a Hero, which costs Hero time but never silent errors.
11. Safety and failure modes
Section titled “11. Safety and failure modes”| Failure | What happens | Where |
|---|---|---|
| Message redelivered | Event id already in the log, dropped as a duplicate, nothing appended | EventStore.append raises DuplicateEntryError |
| Two runs for one load | The second commit fails on expected_version, re-reads and re-decides | ADR 0003 |
| Crash after commit, before send | The outbox row is still pending; the dispatcher sends it | outbox |
| Send succeeds, mark fails | Retry with the same idempotency key; the service dedupes | channels/base.py contract |
| Timer fires twice | Same event_id; second copy dropped | timer_fired_event_id |
| Timer fires after the ask closed | Noop, no model call | OpenAsk.is_timer_current |
| LLM timeout or invalid output | Urgent Hero task with the event; nothing sent | agent/runner.py |
| Draft misses a required fact | No message; Hero task | ADR 0004 |
| Planner proposes contact with a driver | Rejected; Hero task | _planned_action |
| Action targets a contact not on the load, an unsupported channel, a forbidden role, or a used key | Rejected before commit; replaced by a Hero task | agent/validation.py |
| SMS service returns a permanent error | Failed ActionOutcome, urgent Hero task (“call the driver”) | agent/reactions.py |
| Hero task creation fails | Ops alert, no second Hero task (would loop) | ADR 0003 |
| Unknown sender | Hero task, no reply; the same text again within 30 minutes adds nothing | agent/messages.py |
| Dispatcher asks something and the client forbids contacting the dispatcher, or the email is in a side thread, or the stop is unclear | Hero task with the answer; no reply | agent/messages.py |
| Same answer to the same contact within 10 minutes | Not sent again; the triage call is still spent | agent/messages.py |
| Broker news while an ask is open | TMS note and a non-urgent Hero task (one per distinct broker message per 30 minutes); asks keep running until an AskOverride | agent/messages.py |
| Unknown situation | Planner, then Hero | pipeline |
12. Observability (designed; data is built)
Section titled “12. Observability (designed; data is built)”Every run produces a DecisionRecord in the log: event, path (filtered, sop, triage, planner, override, escalated_unknown), SOP ref and version, ask transitions, actions, rationale, and one LLMUsage per model call. Metrics are queries over those records, so nothing needs separate instrumentation to answer “why did it do that”.
| Signal | Source | Alert or use |
|---|---|---|
| Cost per load, by purpose and model | sum of llm_usage.cost_usd | warn at $0.60; hard ceiling at $2 per load, twice the budget, so a runaway load (a chatty driver, a looping SOP) is stopped by disabling the planner and escalating to Hero tasks instead (designed) |
| Path mix (share of filtered, sop, triage, planner, escalated) | DecisionRecord.path | planner share rising means missing SOP coverage |
| Escalation rate per client and topic | Hero tasks per 100 events | spike after an SOP publish means a bad SOP |
| Hero override rate | Hero marks a task “agent was wrong” or reverses an action | primary quality metric |
| Run latency, event received to last action committed | timestamps on the event and the commit | p99 above 60 s pages |
| Outbox age (oldest pending), dead rows per hour | outbox | pages |
Timer pump lag (now minus oldest fire_at not fired) | timers | pages above 2 minutes |
| Conflict rate per load | ConcurrencyConflictError count | investigate queue config |
| Triage confidence histogram | TriageResult.confidence | drift detection after a model change |
| Noop share | is_noop | sudden drop means the filter broke |
Traces: one trace per run with spans for read, project, triage, each LLM call, commit; per-load timeline view from the log (a Hero opens a load and sees every event, decision and action in order). Logs are structured with load_id, event_id, decision_id. Stack: OpenTelemetry to the existing metrics backend. Alternative: a hosted LLM observability product; useful for prompt diffing, but the decision record already has the data, and the production log must stay inside our database for audit.
13. Scale, response time and cost
Section titled “13. Scale, response time and cost”Volume
Section titled “Volume”| Quantity | Value |
|---|---|
| Inbound messages | 5M to 10M per month (100k loads times 50 to 100) |
| Pings | assumed one per 15 minutes while active, about 5 days average: about 450 per load, 45M per month (assumption) |
| Business-hours rate | 7.5M messages over about 220 business hours is about 10 messages per second on average, 30 per second at peak |
| Log rows | per message about 4 rows (event, decision, two delivery outcomes): 20M to 40M rows per month. Per ping 2 rows (event, filtered decision): 90M rows per month. Together about 120M per month, 1.4B per year before retention, about 1,300 rows per load |
| Insert rate | about 50 per second average and 150 at the peak at 1x; the pings are spread over 24 hours |
5-minute budget
Section titled “5-minute budget”| Step | Typical | Worst case | Control |
|---|---|---|---|
| Ingest and queue | under 1 s | seconds | FIFO per load, no head-of-line across loads |
| Read and project | under 50 ms | 200 ms | about 1,300 rows per load with pings; a snapshot row every 200 records (read = snapshot plus tail) is needed once a log passes about 2,000 rows or the read p99 passes 100 ms (designed) |
| Filter, trigger diff, SOP executor, templates | milliseconds | milliseconds | pure code |
| Triage call | about 0.2 s (measured, ADR 0005) | 30 s timeout, 2 retries, capped by the run deadline | LLMSettings.timeout; failure becomes a Hero task |
| Planner call | 3 to 10 s | same | rare |
| Commit | under 20 ms | conflict retry | optimistic concurrency, bounded retries |
| Dispatch | seconds | 150 s retry budget from first send | RetryPolicy; permanent failure escalates |
| Timer lag | under 1 minute | pump interval 30 s | alert at 2 minutes |
Most runs (filtered, SOP, timer) finish in milliseconds with no model call. The timeline for load 481207 makes one model call. Retries inside the run: up to 3 re-decides after a concurrency conflict (max_conflict_retries = 3, each re-reading and re-deciding), up to 2 retries per LLM call, and dispatch retries with capped backoff inside RetryPolicy.budget (150 s from the first send attempt).
Bound on the model path: a triage call and a planner call, each with a 30 s timeout and 2 retries, could add up to 180 s, more across conflict re-decides. A 120 s run deadline spans all conflict re-decides: after it the run makes no model call and commits an urgent Hero task (built, agent/deadline.py). With the 150 s dispatch budget, event to last action is at most 270 s on every path, because an in-flight call is capped at the time left (per phase: httpx applies the timeout to connect, write and read separately, so the bound assumes small non-streamed responses). Alternative: shorter timeouts only, which bounds one call but not the sum.
LLM cost per load
Section titled “LLM cost per load”Assumed prices (replace with contract prices; the structure holds): decision model about $0.0005 to $0.003 per call with a cached system prompt, strong model about $0.02 to $0.05 per planner call, vision about $0.01 to $0.02 per document.
| Item | Per load | Cost |
|---|---|---|
| 50 to 100 inbound messages | assumed 25% to 40% never reach a model: duplicates, acknowledgements with no open ask, unknown senders. Basis: driver SMS is short and often “ok” or “thanks”, and redelivery is a few percent; to be measured from the path mix in production. Timers and delivery outcomes are system events, not inbound messages, and also cost $0 | $0 |
| Pings (about 450) | band-change diff in code | $0 |
| Outbound messages (about 40 to 70) | templates | $0 |
| Triage calls | 30 to 75 (the surviving 60% to 75% of 50 to 100) | $0.015 to $0.23 |
| Planner calls | 0 to 2 | $0 to $0.10 |
| Documents | 2 to 6 | $0.02 to $0.12 |
| Template-less drafts | 0 to 1 | up to $0.01 |
| SOP compile, amortized | 100 clients times 20 saves per month at $0.10 | about $0.002 |
| Total | about $0.04 to $0.46 |
That is a headroom of 2x to 25x under $1. Controls that keep it there: caching the system prompt and SOP guidance; short requests (open asks and the one message, not history); debounce of several SMS in a row into one triage call (designed); a per-load cost counter from the decision records with the $2 ceiling that disables the planner (designed); a first-class noop path for acknowledgements. If the real number drifts, the first levers are the acknowledgement list, the triage prompt size, and routing planner situations to Hero tasks.
At 10x volume (1M loads, 50M to 100M messages per month)
Section titled “At 10x volume (1M loads, 50M to 100M messages per month)”| Pressure | Change |
|---|---|
| Postgres write rate (about 500 inserts per second average and 1,500 at the peak, with pings in the log) and size | Partition the log and outbox by load_id hash; archive closed loads to object storage after 30 days (the log is replayable from the archive); read replicas for dashboards; shard by client group if one cluster saturates |
| Queue | SQS FIFO handles it with many message groups; if per-group throughput or cost bites, Kafka keyed by load with consumer groups (ADR 0003 compares) |
| Timer pump | Partition by fire_at bucket and by hash, several pumps with SKIP LOCKED |
| Dispatcher | More instances; per-service rate limits and circuit breakers so one failing service does not block others |
| LLM | Provider rate limits and quotas become the bottleneck: reserved throughput, a fallback provider, and a queue-side concurrency cap. Triage is the largest cost line, and already runs on the decision model; the next levers are debouncing consecutive SMS into one call (designed) and a distilled in-house classifier validated by the recorded-answer gate (designed) |
| Pings | Filter at ingestion with a cached current band per stop, so unchanged pings never reach the queue |
| Hero capacity | Escalation rate times 10 is a human cost, not a compute cost. Watch tasks per load and tighten SOPs where Heroes keep resolving the same case by hand |
| Cost | Cost per load must stay flat: it does if the path mix does. The path-mix metric is the early warning |
14. Known gaps in this repository
Section titled “14. Known gaps in this repository”- The triage thresholds are calibrated on 190 phrasings we wrote (two recordings), not on production data; a labelled production replay set is required before launch. The default model
jev-latestfloats, so production must pin the calibrated version. The extractor covers the extractor covers the phrasings intests/llm/test_extract.py, not every driver’s way of writing a time. - Webhook auth is a shared secret compared in constant time. Production would verify each provider’s signature (Twilio
X-Twilio-Signature, Slack signing secret with a timestamp window); that is designed only. The channel adapters behind the dispatcher are still fakes that print what the agent “sent” as JSON; there is no SQS adapter and no autoscaling. The Service workflow runs the compose stack and the demo in CI. Processed queue rows are kept so webhook retries are recognised, with no retention job. A message that keeps failing, or whose lease expires on every delivery (a crash loop), is dead-lettered after 5 attempts: the worker logs it and raises an ops alert through theOpsAlertsseam (a JSONops_alertline on stdout here, a pager in production). There is no replay tool for dead letters. The read endpoints (state,events,decisions) share the webhook secret; production needs separate operator authentication for them, since they expose message content. The service uses the real clock, so the 30 minute ETA follow-up fires 30 real minutes after the ask; there is no time-scale override, because any such knob would changefire_atand therefore decision logic. - Per-run latency is recorded in each decision record but not aggregated or exported.
- Eval files are in one folder, named by client; the
standard/andclients/<id>/split described in answers.md is the designed layout. - Behavior diff compares layers and versions of a spec (built); running scenario suites against two versions is designed.
- Document extraction through the vision path (
deepseek-flashwith animage_urlpart) is smoke-tested live only with a generated blank image: the API accepts the request and returns valid JSON. Reading accuracy on real BOLs and PODs is unverified.