Skip to main content

agentplane/store/
postgres_cases.rs

1//! The case layer on `PostgreSQL`.
2//!
3//! Every store here exists to settle one race, and the races are the reason this
4//! file is not a mechanical translation of the `SQLite` one:
5//!
6//! | Store | Race | How Postgres settles it |
7//! |---|---|---|
8//! | cases | two messages, one new matter | partial unique index on open keys |
9//! | events | one message, two waiters | `UPDATE … RETURNING` claims in one statement |
10//! | timers | one wake-up, two sweeps | same |
11//! | tasks | one decision, two reviewers | same |
12//! | batches | one item, two reservations | `ON CONFLICT DO NOTHING`, then read back |
13//!
14//! `UPDATE … RETURNING` is the reason several of these are *simpler* here than
15//! in `SQLite` rather than merely different. The read-then-write that `SQLite` has
16//! to wrap in a transaction becomes a single statement whose result tells you
17//! whether you won — there is no window to reason about because there is no
18//! second statement.
19//!
20//! Whether that reasoning is right is not left to the reader:
21//! `testkit::conformance_case` runs the same battery against this and against
22//! `SQLite`.
23
24use async_trait::async_trait;
25use deadpool_postgres::Pool;
26use tokio_postgres::error::SqlState;
27
28use crate::batch::{BatchCensus, BatchStore, ItemOutcome, ItemRecord};
29use crate::case::{BufferedEvent, ClaimError, TargetedDelivery};
30use crate::case::{CaseCensus, CaseStore, Correlation, EventStore, TaskStore, TimerStore};
31use crate::core::{
32    BatchId, BreachNote, Case, CaseId, CaseStatus, CaseVersion, CorrelationKey, DeadLetter,
33    Deadline, DeadlineState, Digest, EffectKey, InboundEvent, LegalHold, OnExpiry, Priority, RunId,
34    Spend, StoreError, Subscription, Task, TaskId, TaskState, Timer, Timestamp,
35};
36
37use super::postgres::{PostgresStore, amount_of, be, pool_err, sql_amount};
38
39/// A hold's stored attribution, or corruption.
40///
41/// Fails closed for the reason every decoder in this crate does: a hold this
42/// build cannot read is a preservation order a retention pass would walk past,
43/// and from the outside that is indistinguishable from one that was released.
44fn hold_row(reason: String, actor: &str, basis: &str) -> Result<super::HoldRow, StoreError> {
45    Ok(super::HoldRow {
46        reason,
47        by: super::decode_operator(actor, basis, "case_legal_holds")?,
48    })
49}
50
51pub(super) const CASE_SCHEMA: &str = "
52-- Every table here leads with the tenant, for the reason the journal schema
53-- gives: a key component turns a forgotten predicate into an empty result
54-- instead of another tenant's row. The foreign keys carry it too, so a child
55-- row cannot reference a parent in a different tenant.
56CREATE TABLE IF NOT EXISTS cases (
57    tenant    TEXT   NOT NULL,
58    case_id   TEXT   NOT NULL,
59    kind      TEXT   NOT NULL,
60    status    TEXT   NOT NULL,
61    state     TEXT   NOT NULL,
62    -- Bumped by every state write, which must name the version it read. This
63    -- backend is the one that exists for several plane instances sharing a
64    -- store, so the overlapping read-modify-write is not a corner case here —
65    -- it is the normal operating condition.
66    version   BIGINT NOT NULL DEFAULT 0 CHECK (version >= 0),
67    opened_at BIGINT NOT NULL,
68    PRIMARY KEY (tenant, case_id)
69);
70
71-- The worklist read: cases by status, newest first.
72CREATE INDEX IF NOT EXISTS cases_by_status
73    ON cases (tenant, status, opened_at DESC);
74
75CREATE TABLE IF NOT EXISTS case_correlation (
76    tenant    TEXT    NOT NULL,
77    case_id   TEXT    NOT NULL,
78    namespace TEXT    NOT NULL,
79    value     TEXT    NOT NULL,
80    open      BOOLEAN NOT NULL DEFAULT TRUE,
81    PRIMARY KEY (tenant, case_id, namespace, value),
82    FOREIGN KEY (tenant, case_id) REFERENCES cases (tenant, case_id) ON DELETE CASCADE
83);
84
85-- One open case per business key. This is the arbiter, not a hint: two inbound
86-- messages racing to open the same matter both attempt the insert, and exactly
87-- one succeeds. The loser re-reads and attaches.
88-- Scoped to the tenant, because a correlation key is a *business* value:
89-- `document`/`DOC-1` means something different to every tenant, and two of them
90-- using it is ordinary rather than a collision. Globally unique, one tenant's
91-- run would attach to another's case and the two would share a history, a
92-- deadline set and an erasure unit.
93CREATE UNIQUE INDEX IF NOT EXISTS case_correlation_open
94    ON case_correlation (tenant, namespace, value) WHERE open;
95
96CREATE TABLE IF NOT EXISTS case_runs (
97    tenant  TEXT   NOT NULL,
98    case_id TEXT   NOT NULL,
99    run_id  TEXT   NOT NULL,
100    seq     BIGINT NOT NULL CHECK (seq >= 0),
101    PRIMARY KEY (tenant, case_id, run_id),
102    FOREIGN KEY (tenant, case_id) REFERENCES cases (tenant, case_id) ON DELETE CASCADE
103);
104
105-- One run per position. `attach_run` allocates seq as MAX+1, and two instances
106-- attaching different runs concurrently would both compute the same next
107-- position; the constraint makes the collision an error the store retries
108-- rather than two runs silently sharing a place in the case's order.
109CREATE UNIQUE INDEX IF NOT EXISTS case_runs_order
110    ON case_runs (tenant, case_id, seq);
111
112-- The blobs a case produced. The case is what an erasure request names, and a
113-- digest cannot be reversed to find its case, so the link has to be recorded
114-- when the bytes are written or it cannot be recovered at all.
115CREATE TABLE IF NOT EXISTS case_blobs (
116    tenant     TEXT   NOT NULL,
117    case_id    TEXT   NOT NULL,
118    digest     BYTEA  NOT NULL,
119    written_at BIGINT NOT NULL,
120    PRIMARY KEY (tenant, case_id, digest),
121    FOREIGN KEY (tenant, case_id) REFERENCES cases (tenant, case_id) ON DELETE CASCADE
122);
123
124CREATE INDEX IF NOT EXISTS case_blobs_time
125    ON case_blobs (tenant, case_id, written_at);
126
127-- `by_actor` and `by_basis` are on the same footing as the reason: a
128-- preservation order the runtime cannot check is worth exactly the name beside
129-- it, and the basis says whether an authenticator produced that name or
130-- somebody holding the connection string typed it.
131CREATE TABLE IF NOT EXISTS case_legal_holds (
132    tenant    TEXT   NOT NULL,
133    case_id   TEXT   NOT NULL,
134    placed_at BIGINT NOT NULL,
135    reason    TEXT   NOT NULL,
136    by_actor  TEXT   NOT NULL,
137    by_basis  TEXT   NOT NULL,
138    PRIMARY KEY (tenant, case_id),
139    FOREIGN KEY (tenant, case_id) REFERENCES cases (tenant, case_id) ON DELETE CASCADE
140);
141
142-- The listing, oldest hold first. Ascending unlike every other backlog index
143-- here, because `release_hold` empties this one and the longest-standing
144-- unlifted hold is the entry somebody has to ask about.
145CREATE INDEX IF NOT EXISTS case_legal_holds_by_time
146    ON case_legal_holds (tenant, placed_at, case_id);
147
148-- The last recovery rehearsal, one row per tenant that the latest write
149-- replaces. A history of drills is a different artifact and a larger promise
150-- than the question an audit asks: *when did you last rehearse, and did it
151-- pass*. The columns are counts, an instant and the checkpoint the plane was
152-- serving — no case is referenced, so nothing cascades and a drill's verdict
153-- outlives the matters it walked, which is the point.
154CREATE TABLE IF NOT EXISTS case_last_drill (
155    tenant      TEXT    NOT NULL,
156    ran_at      BIGINT  NOT NULL,
157    sound       BOOLEAN NOT NULL,
158    cases       BIGINT  NOT NULL,
159    findings    BIGINT  NOT NULL,
160    -- Kept apart from `sound`: a pass over nothing is not a pass.
161    not_checked BIGINT  NOT NULL,
162    origin      TEXT    NOT NULL,
163    log_size    BIGINT  NOT NULL,
164    PRIMARY KEY (tenant)
165);
166
167CREATE TABLE IF NOT EXISTS case_deadlines (
168    tenant          TEXT   NOT NULL,
169    case_id         TEXT   NOT NULL,
170    name            TEXT   NOT NULL,
171    resolved_at     BIGINT NOT NULL,
172    calendar_digest BYTEA  NOT NULL,
173    warn_at         BIGINT,
174    state           TEXT   NOT NULL,
175    -- Who accounted for the breach, once somebody has. Absent on every other
176    -- state; `acknowledged_at` is what says the other two mean anything.
177    acknowledged_at    BIGINT,
178    acknowledged_by    TEXT,
179    acknowledged_basis TEXT,
180    acknowledged_note  TEXT,
181    PRIMARY KEY (tenant, case_id, name),
182    FOREIGN KEY (tenant, case_id) REFERENCES cases (tenant, case_id) ON DELETE CASCADE
183);
184-- A breach applied whose account the sweep has not yet written: set by the
185-- breach in its own statement, cleared once the sweep's notes land.
186ALTER TABLE case_deadlines ADD COLUMN IF NOT EXISTS breach_unnoted BOOLEAN NOT NULL DEFAULT FALSE;
187CREATE INDEX IF NOT EXISTS case_deadlines_unnoted
188    ON case_deadlines (tenant, resolved_at)
189    WHERE breach_unnoted;
190
191-- The sweep's read: outstanding obligations by due instant.
192CREATE INDEX IF NOT EXISTS case_deadlines_due
193    ON case_deadlines (tenant, state, resolved_at);
194
195-- The obligation listing and its gauge: breaches nobody has accounted for,
196-- longest-overdue first. Partial, because an acknowledged breach is never
197-- selected and a deployment's whole history of them should not be paged
198-- through to answer what is still open.
199CREATE INDEX IF NOT EXISTS case_deadlines_unaccounted
200    ON case_deadlines (tenant, resolved_at)
201    WHERE state = 'breached' AND acknowledged_at IS NULL;
202
203CREATE TABLE IF NOT EXISTS inbound_events (
204    -- The dedup identity is (source, id), CloudEvents' uniqueness pair. Keying
205    -- on id alone deduplicates two producers into each other: id is unique
206    -- within a producer, and the collision is silent because the second message
207    -- looks exactly like a retry of the first.
208    tenant      TEXT    NOT NULL,
209    event_id    TEXT    NOT NULL,
210    source      TEXT    NOT NULL,
211    -- The producer's own id, stored rather than split back out of the key: a
212    -- reconstructed event must be the event that arrived, and parsing the key
213    -- apart would be a second place that has to agree about the separator.
214    bare_id     TEXT    NOT NULL,
215    kind        TEXT    NOT NULL,
216    payload     TEXT    NOT NULL,
217    -- The operator this plane minted the message for, where it minted one.
218    -- Two columns rather than a rendered string, for the reason every other
219    -- operator column here is two: a name and what established it are two
220    -- facts, and `alice (asserted)` is a sentence a reader has to parse back.
221    -- Null for everything that arrived over a wire.
222    by_actor    TEXT,
223    by_basis    TEXT,
224    received_at BIGINT  NOT NULL,
225    claimed_by  TEXT,
226    claimed_at  BIGINT,
227    -- The wait a standing claim was made for, until that wait's unsubscribe
228    -- consumes it. A wait recovers only its own standing claim, and retiring a
229    -- wait sheds only what it consumed.
230    claimed_for TEXT,
231    -- Delivered to one run by name: that run's alone, dead-lettered rather
232    -- than offered to another run when its addressee concludes unconsumed.
233    targeted    BOOLEAN NOT NULL DEFAULT FALSE,
234    dead        BOOLEAN NOT NULL DEFAULT FALSE,
235    dead_reason TEXT,
236    PRIMARY KEY (tenant, event_id)
237);
238-- Set by `erase_payload` and never cleared: an erased row is never handed to
239-- a run as a value, whatever its claim says.
240ALTER TABLE inbound_events ADD COLUMN IF NOT EXISTS erased BOOLEAN NOT NULL DEFAULT FALSE;
241
242-- The sweep's read: live unclaimed events by age. Partial, so the index is
243-- exactly the sweep's candidate set rather than the whole buffer.
244CREATE INDEX IF NOT EXISTS inbound_events_live
245    ON inbound_events (tenant, received_at)
246    WHERE claimed_by IS NULL AND NOT dead;
247
248-- The dead-letter view: retired events, newest first.
249CREATE INDEX IF NOT EXISTS inbound_events_dead
250    ON inbound_events (tenant, received_at DESC)
251    WHERE dead;
252
253-- The unsubscribe strip's read: the rows a run holds claimed, whose delivered
254-- payloads are shed once the run's wait is satisfied and journaled.
255CREATE INDEX IF NOT EXISTS inbound_events_claimed
256    ON inbound_events (tenant, claimed_by)
257    WHERE claimed_by IS NOT NULL;
258
259CREATE TABLE IF NOT EXISTS inbound_correlation (
260    tenant    TEXT NOT NULL,
261    event_id  TEXT NOT NULL,
262    namespace TEXT NOT NULL,
263    value     TEXT NOT NULL,
264    PRIMARY KEY (tenant, event_id, namespace, value),
265    FOREIGN KEY (tenant, event_id)
266        REFERENCES inbound_events (tenant, event_id) ON DELETE CASCADE
267);
268
269-- The claim path's join: which buffered events carry this business key.
270CREATE INDEX IF NOT EXISTS inbound_correlation_by_key
271    ON inbound_correlation (tenant, namespace, value);
272
273-- The match path: an arriving event finds the runs waiting for it here. The
274-- tenant leads for the same reason it leads `case_correlation_open`, and this
275-- is the worse of the two failures — one tenant's message resuming another
276-- tenant's run hands it a payload it was never sent.
277CREATE TABLE IF NOT EXISTS subscriptions (
278    tenant     TEXT   NOT NULL,
279    run_id     TEXT   NOT NULL,
280    effect_key TEXT   NOT NULL,
281    case_id    TEXT,
282    step       BIGINT NOT NULL,
283    phase      TEXT   NOT NULL,
284    event_kind TEXT   NOT NULL,
285    namespace  TEXT   NOT NULL,
286    value      TEXT   NOT NULL,
287    created_at BIGINT NOT NULL,
288    -- Set while the wait holds a claimed event nothing has delivered yet: the
289    -- redelivery pass reads these rows rather than every registered wait.
290    parked_at  BIGINT,
291    PRIMARY KEY (tenant, run_id, effect_key, namespace, value)
292);
293-- The one event source this wait accepts; NULL accepts any.
294ALTER TABLE subscriptions ADD COLUMN IF NOT EXISTS from_source TEXT;
295
296-- The match path's read: who is waiting on this kind and key. The primary key
297-- leads with the run for targeted delivery; this serves the arriving event,
298-- which knows only what it carries.
299CREATE INDEX IF NOT EXISTS subscriptions_by_key
300    ON subscriptions (tenant, event_kind, namespace, value);
301
302-- The redelivery sweep's read: the oldest waits, bounded. Registration order is
303-- the index's own order, so a tick walks `limit` rows instead of sorting every
304-- subscription the tenant has — which is the shape the embedded backend already
305-- had in `SUBS_BY_TIME` and this one did not, so the two agreed on the answer
306-- and disagreed on the cost, on a path that runs every tick against a
307-- population a plane is *expected* to accumulate.
308CREATE INDEX IF NOT EXISTS subscriptions_waiting
309    ON subscriptions (tenant, created_at, run_id, effect_key);
310
311-- The redelivery pass's read: waits parked with a claimed event. Partial, so a
312-- plane's long legitimate waits never stand in front of the pairs it is for.
313CREATE INDEX IF NOT EXISTS subscriptions_parked
314    ON subscriptions (tenant, parked_at, run_id, effect_key)
315    WHERE parked_at IS NOT NULL;
316
317CREATE TABLE IF NOT EXISTS timers (
318    tenant     TEXT   NOT NULL,
319    run_id     TEXT   NOT NULL,
320    effect_key TEXT   NOT NULL,
321    case_id    TEXT,
322    step       BIGINT NOT NULL,
323    phase      TEXT   NOT NULL,
324    fire_at    BIGINT NOT NULL CHECK (fire_at >= 0),
325    claimed_at BIGINT CHECK (claimed_at >= 0),
326    PRIMARY KEY (tenant, run_id, effect_key)
327);
328-- The due sweep's read, soonest first: a tick walks the due rows rather than
329-- every timer the tenant holds.
330CREATE INDEX IF NOT EXISTS timers_due ON timers (tenant, fire_at);
331
332CREATE TABLE IF NOT EXISTS tasks (
333    tenant          TEXT   NOT NULL,
334    task_id         TEXT   NOT NULL,
335    run_id          TEXT   NOT NULL,
336    case_id         TEXT,
337    kind            TEXT   NOT NULL,
338    justification   TEXT   NOT NULL,
339    -- Arrays, never a delimited string. Role and actor names are the four-eyes
340    -- control's operands and this store does not get to constrain their
341    -- alphabet: an exclusion list stored as 'a,b'-joined text reads an actor
342    -- named 'a,b' back as two actors named neither, and the person barred from
343    -- deciding is barred no longer.
344    candidate_roles TEXT[] NOT NULL,
345    escalate_to     TEXT[] NOT NULL,
346    excluded_actors TEXT[] NOT NULL,
347    assignee        TEXT,
348    priority        TEXT   NOT NULL,
349    state           TEXT   NOT NULL,
350    on_expiry       TEXT   NOT NULL,
351    -- `Priority::rank`, stored rather than expressed in SQL. Queue order is one
352    -- rule; a `CASE priority WHEN …` here would be a second copy of the
353    -- priority vocabulary, in a dialect the compiler cannot check, whose
354    -- fallback arm sorts an unranked priority last — in the one index whose
355    -- whole job is order.
356    priority_rank   SMALLINT NOT NULL,
357    created_at      BIGINT NOT NULL,
358    due_at          BIGINT,
359    -- `Withheld::as_str` when the proposal is sealed at rest; the decorator
360    -- that sealed it writes it, so no argument value can spell it.
361    withheld        TEXT,
362    PRIMARY KEY (tenant, task_id)
363);
364
365-- The queue's read: pending work, most urgent first and oldest within a rank.
366CREATE INDEX IF NOT EXISTS tasks_queue_rank
367    ON tasks (tenant, state, priority_rank, created_at);
368
369-- The overdue sweep's read. Partial: most tasks have no window to close.
370CREATE INDEX IF NOT EXISTS tasks_due
371    ON tasks (tenant, due_at)
372    WHERE due_at IS NOT NULL;
373
374CREATE TABLE IF NOT EXISTS batches (
375    tenant      TEXT    NOT NULL,
376    batch_id    TEXT    NOT NULL,
377    plan_digest TEXT    NOT NULL,
378    exhausted   BOOLEAN NOT NULL DEFAULT FALSE,
379    PRIMARY KEY (tenant, batch_id)
380);
381
382CREATE TABLE IF NOT EXISTS batch_items (
383    tenant   TEXT   NOT NULL,
384    batch_id TEXT   NOT NULL,
385    -- Byte order, as the embedded backend compares keys: the resume cursor is
386    -- a `<` over these, and a locale collation would order a source's pages
387    -- differently from the source.
388    item_key TEXT   COLLATE \"C\" NOT NULL,
389    run_id   TEXT   NOT NULL,
390    outcome  TEXT,
391    detail   TEXT,
392    -- Spend cannot be negative; the docs say so, and a constraint is the only
393    -- form of cannot a shared store honours.
394    tokens   BIGINT NOT NULL DEFAULT 0 CHECK (tokens >= 0),
395    minor    BIGINT NOT NULL DEFAULT 0 CHECK (minor >= 0),
396    PRIMARY KEY (tenant, batch_id, item_key),
397    FOREIGN KEY (tenant, batch_id) REFERENCES batches (tenant, batch_id) ON DELETE CASCADE
398);
399
400-- The runs a tenant currently has executing.
401--
402-- A set rather than a counter, and the difference is recovery: a counter that is
403-- incremented on admission leaks a slot every time a process dies before the
404-- decrement, and nothing can say which increments were real. A set names its
405-- members, so a stranded slot is attributable to a run an operator can look up,
406-- and releasing it is idempotent by construction.
407CREATE TABLE IF NOT EXISTS quota_running (
408    tenant      TEXT   NOT NULL,
409    run_id      TEXT   NOT NULL,
410    admitted_at BIGINT NOT NULL,
411    PRIMARY KEY (tenant, run_id)
412);
413
414-- The emergency stop: one row per standing halt, none for the rest.
415--
416-- In the database rather than in a process, because a switch that stops only
417-- the instance it was thrown on is not a switch — it is the in-process-counter
418-- failure arriving during an incident.
419--
420-- `scope` is part of the key, not a column beside it: a halt on one agent and a
421-- halt on the whole tenant are two rows rather than one flag the last writer
422-- wins. An incident that widens and then partly resolves is the ordinary shape,
423-- and a single overwritable flag gets it wrong in the direction that lets work
424-- through. The scope spellings are the ones `HaltScope::parse` accepts, and
425-- they are not restated here: a grammar written twice is a grammar that loses a
426-- form the day one is added.
427--
428-- `by_actor` and `by_basis` are the whole evidentiary weight of the row. The
429-- runtime cannot check an emergency stop — there is no verdict to re-derive —
430-- so who asked, and what established that name, is the record. The basis is
431-- kept beside the name rather than inferred from the surface, because the same
432-- act arrives from an authenticated API caller and from somebody holding the
433-- database URL, and a reader years later cannot tell those apart from a name.
434CREATE TABLE IF NOT EXISTS quota_halted (
435    tenant      TEXT        NOT NULL,
436    scope       TEXT        NOT NULL,
437    reason      TEXT        NOT NULL,
438    by_actor    TEXT        NOT NULL,
439    by_basis    TEXT        NOT NULL,
440    thrown_at   BIGINT      NOT NULL,
441    PRIMARY KEY (tenant, scope)
442);
443
444-- The manifest registry: one row per published name and version.
445--
446-- The primary key is the immutability rule as a database constraint rather
447-- than as application logic, for the reason exactly-once is a unique index
448-- here: application logic can be bypassed by the next caller and a constraint
449-- cannot. What the transaction around it adds is the *decision* — a version
450-- holding different content is a refusal, not an overwrite.
451--
452-- `key_id` and `signature` are empty for a publish that happened before
453-- anybody signed. That is an operational fact rather than a defect, and it is
454-- not a sentinel overlapping a real value: a signature cannot exist without a
455-- key id.
456CREATE TABLE IF NOT EXISTS registry_manifests (
457    tenant      TEXT   NOT NULL,
458    name        TEXT   NOT NULL,
459    version     TEXT   NOT NULL,
460    digest      TEXT   NOT NULL,
461    yaml        TEXT   NOT NULL,
462    key_id      TEXT   NOT NULL DEFAULT '',
463    signature   TEXT   NOT NULL DEFAULT '',
464    PRIMARY KEY (tenant, name, version)
465);
466
467CREATE TABLE IF NOT EXISTS quota_spent (
468    tenant      TEXT   NOT NULL,
469    period      TEXT   NOT NULL,
470    tokens      BIGINT NOT NULL DEFAULT 0 CHECK (tokens >= 0),
471    minor_units BIGINT NOT NULL DEFAULT 0 CHECK (minor_units >= 0),
472    PRIMARY KEY (tenant, period)
473);
474
475-- One exact receipt per live execution pass. The payload is retained so a
476-- retry with changed accounting is damage, not a second interpretation of the
477-- same key.
478CREATE TABLE IF NOT EXISTS quota_settled (
479    tenant       TEXT    NOT NULL,
480    run_id       TEXT    NOT NULL,
481    epoch        BIGINT  NOT NULL CHECK (epoch >= 0),
482    period       TEXT,
483    tokens       BIGINT  NOT NULL CHECK (tokens >= 0),
484    minor_units  BIGINT  NOT NULL CHECK (minor_units >= 0),
485    release_slot BOOLEAN NOT NULL,
486    concludes    BOOLEAN NOT NULL,
487    PRIMARY KEY (tenant, run_id, epoch)
488);
489
490-- What an open run still holds against a period: written beside the slot at
491-- admission, reduced by every pass settlement and removed by the concluding
492-- one, each in the transaction that writes the receipt. One row per run, so a
493-- resume in a later period moves it rather than adding a second.
494CREATE TABLE IF NOT EXISTS quota_reserved (
495    tenant      TEXT   NOT NULL,
496    run_id      TEXT   NOT NULL,
497    period      TEXT   NOT NULL,
498    tokens      BIGINT NOT NULL CHECK (tokens >= 0),
499    minor_units BIGINT NOT NULL CHECK (minor_units >= 0),
500    PRIMARY KEY (tenant, run_id)
501);
502
503-- One row per dispatch of a rate-ceilinged tool. The key is the idempotency:
504-- a retry and a recovered re-dispatch carry the same run and first-attempt
505-- key and find their own row. The count is the rows for (tenant, grant_ref)
506-- inside the window, decided under a per-grant advisory lock. Rows age out
507-- with the widest window and are never refunded early.
508CREATE TABLE IF NOT EXISTS quota_rate (
509    tenant      TEXT   NOT NULL,
510    grant_ref   TEXT   NOT NULL,
511    run_id      TEXT   NOT NULL,
512    dispatch    TEXT   NOT NULL,
513    reserved_at BIGINT NOT NULL,
514    PRIMARY KEY (tenant, grant_ref, run_id, dispatch)
515);
516";
517
518fn corrupt(what: &str, e: impl std::fmt::Display) -> StoreError {
519    StoreError::Corrupt {
520        seq: 0,
521        detail: format!("{what}: {e}"),
522    }
523}
524
525/// Read a stored column through the type's own vocabulary.
526///
527/// Every decoder below goes through the core type's `parse` rather than a match
528/// of its own — see [`CaseStatus::parse`] for why nothing else may spell a
529/// stored vocabulary, and why a reader that defaults is worse than one that
530/// refuses. The whole point of `phase` is telling a step's forward pass from
531/// its compensating one, so a silent `Forward` hands the unwind logic a
532/// compensating record wearing the wrong half of the saga; `Deny` and `Normal`
533/// are safe values, and that is precisely what makes them the wrong answer — a
534/// fail-closed default is still a fact this store invented about a row it could
535/// not read. Every call site already returns `StoreError`, so refusing costs
536/// nothing but the `?`.
537fn decoded<T>(what: &str, raw: &str, parsed: Option<T>) -> Result<T, StoreError> {
538    parsed.ok_or_else(|| corrupt(&format!("unknown {what}"), raw))
539}
540
541fn status_from(s: &str) -> Result<CaseStatus, StoreError> {
542    decoded("case status", s, CaseStatus::parse(s))
543}
544
545fn deadline_state_from(s: &str) -> Result<DeadlineState, StoreError> {
546    decoded("deadline state", s, DeadlineState::parse(s))
547}
548
549fn task_state_from(s: &str) -> Result<TaskState, StoreError> {
550    decoded("task state", s, TaskState::parse(s))
551}
552
553fn priority_from(s: &str) -> Result<Priority, StoreError> {
554    decoded("task priority", s, Priority::parse(s))
555}
556
557fn expiry_from(s: &str) -> Result<OnExpiry, StoreError> {
558    decoded("expiry policy", s, OnExpiry::parse(s))
559}
560
561fn phase_from(s: &str) -> Result<crate::core::Phase, StoreError> {
562    decoded("step phase", s, crate::core::Phase::parse(s))
563}
564
565/// The stored spellings of every state a typed predicate admits.
566///
567/// A membership test in SQL takes its list from the rule rather than repeating
568/// it. A literal `IN ('open','claimed','escalated')` is the vocabulary spelled
569/// again in a dialect no compiler checks, and it keeps working when a variant
570/// is added — by leaving the new state out of the worklist, the backlog count
571/// and the expiry sweep at once, which is a task that exists and is in no index.
572fn task_states(admits: fn(TaskState) -> bool) -> Vec<&'static str> {
573    TaskState::ALL
574        .iter()
575        .copied()
576        .filter(|s| admits(*s))
577        .map(TaskState::as_str)
578        .collect()
579}
580
581/// The obligation half of [`task_states`].
582fn deadline_states(admits: fn(DeadlineState) -> bool) -> Vec<&'static str> {
583    DeadlineState::ALL
584        .iter()
585        .copied()
586        .filter(|s| admits(*s))
587        .map(DeadlineState::as_str)
588        .collect()
589}
590
591impl PostgresStore {
592    fn pool(&self) -> &Pool {
593        self.pool_ref()
594    }
595}
596
597#[async_trait]
598impl CaseStore for PostgresStore {
599    fn tenant(&self) -> &str {
600        self.tenant_str()
601    }
602
603    async fn correlate(&self, keys: &[CorrelationKey]) -> Result<Option<CaseId>, StoreError> {
604        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
605        for k in keys {
606            let row = client
607                .query_opt(
608                    "SELECT case_id FROM case_correlation
609                      WHERE tenant = $3 AND namespace = $1 AND value = $2 AND open",
610                    &[&k.namespace, &k.value, &self.tenant_name()],
611                )
612                .await
613                .map_err(|e| be(&e))?;
614            if let Some(row) = row {
615                let id: String = row.get(0);
616                return Ok(Some(
617                    CaseId::parse(&id).map_err(|e| corrupt("bad case id", e))?,
618                ));
619            }
620        }
621        Ok(None)
622    }
623
624    async fn correlate_or_open(
625        &self,
626        kind: &str,
627        keys: &[CorrelationKey],
628        at: Timestamp,
629    ) -> Result<Correlation, StoreError> {
630        // Two attempts, and the second is not a retry loop for flakiness — it is
631        // how the constraint arbitrates. Both messages read no open case, both
632        // try to insert, and the unique index picks a winner. The loser comes
633        // back here and *finds* the winner's row, which is exactly the answer it
634        // should have had.
635        for attempt in 0..2 {
636            if let Some(existing) = self.correlate(keys).await? {
637                return Ok(Correlation::Attached(existing));
638            }
639
640            let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
641            let tx = client.transaction().await.map_err(|e| be(&e))?;
642            let id = CaseId::generate();
643            tx.execute(
644                "INSERT INTO cases (case_id, kind, status, state, opened_at, tenant)
645                 VALUES ($1, $2, $3, $4, $5, $6)",
646                &[
647                    &id.to_string(),
648                    &kind.to_owned(),
649                    &CaseStatus::Open.as_str(),
650                    &"null".to_owned(),
651                    &at.unix_timestamp(),
652                    &self.tenant_name(),
653                ],
654            )
655            .await
656            .map_err(|e| be(&e))?;
657
658            let mut lost = false;
659            for k in keys {
660                let r = tx
661                    .execute(
662                        "INSERT INTO case_correlation (case_id, namespace, value, open, tenant)
663                         VALUES ($1, $2, $3, TRUE, $4)",
664                        &[&id.to_string(), &k.namespace, &k.value, &self.tenant_name()],
665                    )
666                    .await;
667                if let Err(e) = r {
668                    if e.code() == Some(&SqlState::UNIQUE_VIOLATION) {
669                        lost = true;
670                        break;
671                    }
672                    return Err(be(&e));
673                }
674            }
675            if lost {
676                // Someone else opened this matter first. Discard ours entirely —
677                // a half-opened case with no keys is a case nothing can reach.
678                drop(tx);
679                continue;
680            }
681            tx.commit().await.map_err(|e| be(&e))?;
682            let _ = attempt;
683            return Ok(Correlation::Opened(id));
684        }
685
686        // Two losses means a third party is opening and closing this key in a
687        // tight loop; say so rather than looping forever.
688        self.correlate(keys)
689            .await?
690            .map(Correlation::Attached)
691            .ok_or_else(|| {
692                StoreError::Backend(
693                    "could not open or attach a case: the correlation key is being opened \
694                     and closed concurrently"
695                        .into(),
696                )
697            })
698    }
699
700    async fn case(&self, id: CaseId) -> Result<Option<Case>, StoreError> {
701        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
702        let Some(row) = client
703            .query_opt(
704                "SELECT kind, status, state, opened_at, version FROM cases
705                  WHERE case_id = $1 AND tenant = $2",
706                &[&id.to_string(), &self.tenant_name()],
707            )
708            .await
709            .map_err(|e| be(&e))?
710        else {
711            return Ok(None);
712        };
713        let kind: String = row.get(0);
714        let status: String = row.get(1);
715        let state: String = row.get(2);
716        let opened: i64 = row.get(3);
717        let version: i64 = row.get(4);
718
719        let corr = client
720            .query(
721                "SELECT namespace, value FROM case_correlation
722                  WHERE case_id = $1 AND tenant = $2",
723                &[&id.to_string(), &self.tenant_name()],
724            )
725            .await
726            .map_err(|e| be(&e))?;
727        let runs = client
728            .query(
729                "SELECT run_id FROM case_runs
730                  WHERE case_id = $1 AND tenant = $2 ORDER BY seq ASC",
731                &[&id.to_string(), &self.tenant_name()],
732            )
733            .await
734            .map_err(|e| be(&e))?;
735
736        Ok(Some(Case {
737            id,
738            kind,
739            status: status_from(&status)?,
740            correlation: corr
741                .iter()
742                .map(|r| CorrelationKey::new(r.get::<_, String>(0), r.get::<_, String>(1)))
743                .collect(),
744            state: serde_json::from_str(&state)?,
745            version: CaseVersion(u64::try_from(version).unwrap_or(0)),
746            opened_at: Timestamp::from_unix_timestamp(opened)
747                .map_err(|e| corrupt("unrepresentable opened_at", e))?,
748            runs: runs
749                .iter()
750                .map(|r| RunId::parse(&r.get::<_, String>(0)))
751                .collect::<Result<Vec<_>, _>>()
752                .map_err(|e| corrupt("bad run id", e))?,
753        }))
754    }
755
756    async fn cases(&self, after: Option<CaseId>, limit: usize) -> Result<Vec<Case>, StoreError> {
757        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
758        // Ids first, rows second — through `case()`, so the assembly of a case
759        // from its tables exists exactly once and the export cannot drift from
760        // the ordinary reader.
761        let cursor = after.map(|c| c.to_string()).unwrap_or_default();
762        let rows = client
763            .query(
764                "SELECT case_id FROM cases
765                  WHERE tenant = $1 AND case_id > $2
766                  ORDER BY case_id ASC
767                  LIMIT $3",
768                &[
769                    &self.tenant_name(),
770                    &cursor,
771                    &i64::try_from(limit).unwrap_or(i64::MAX),
772                ],
773            )
774            .await
775            .map_err(|e| be(&e))?;
776        drop(client);
777        let mut out = Vec::with_capacity(rows.len());
778        for row in rows {
779            let id: String = row.get(0);
780            let id = CaseId::parse(&id).map_err(|e| corrupt("bad case id", e))?;
781            if let Some(case) = self.case(id).await? {
782                out.push(case);
783            }
784        }
785        Ok(out)
786    }
787
788    #[allow(clippy::too_many_lines)]
789    async fn import_case(
790        &self,
791        case: &Case,
792        deadlines: &[Deadline],
793        blobs: &[Digest],
794    ) -> Result<(), StoreError> {
795        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
796        let tx = client.transaction().await.map_err(|e| be(&e))?;
797        let key = case.id.to_string();
798        let tenant = self.tenant_name();
799        let state = serde_json::to_string(&case.state)?;
800
801        // Refused rather than upserted: a restore rebuilds a case layer, it
802        // does not merge one, and an `ON CONFLICT DO UPDATE` here would let a
803        // second restore silently rewrite a matter.
804        let inserted = tx
805            .execute(
806                "INSERT INTO cases (tenant, case_id, kind, status, state, version, opened_at)
807                 VALUES ($1, $2, $3, $4, $5, $6, $7)
808                 ON CONFLICT DO NOTHING",
809                &[
810                    &tenant,
811                    &key,
812                    &case.kind,
813                    &case.status.as_str(),
814                    &state,
815                    &i64::try_from(case.version.0).unwrap_or(i64::MAX),
816                    &case.opened_at.unix_timestamp(),
817                ],
818            )
819            .await
820            .map_err(|e| be(&e))?;
821        if inserted == 0 {
822            return Err(StoreError::Backend(format!(
823                "case {} already exists — a restore rebuilds a case layer, it does not \
824                 merge one",
825                case.id
826            )));
827        }
828
829        for k in &case.correlation {
830            // The partial unique index over open correlation rows is the
831            // arbiter, exactly as it is for `correlate_or_open`: importing an
832            // open case whose key another open case claims must fail rather
833            // than merge two matters.
834            tx.execute(
835                "INSERT INTO case_correlation (tenant, case_id, namespace, value, open)
836                 VALUES ($1, $2, $3, $4, $5)",
837                &[
838                    &tenant,
839                    &key,
840                    &k.namespace,
841                    &k.value,
842                    &(case.status != CaseStatus::Closed),
843                ],
844            )
845            .await
846            .map_err(|e| be(&e))?;
847        }
848
849        for (at, run) in case.runs.iter().enumerate() {
850            tx.execute(
851                "INSERT INTO case_runs (tenant, case_id, run_id, seq)
852                 VALUES ($1, $2, $3, $4)",
853                &[
854                    &tenant,
855                    &key,
856                    &run.to_string(),
857                    &i64::try_from(at).unwrap_or(i64::MAX),
858                ],
859            )
860            .await
861            .map_err(|e| be(&e))?;
862        }
863
864        for deadline in deadlines {
865            tx.execute(
866                "INSERT INTO case_deadlines
867                    (tenant, case_id, name, resolved_at, calendar_digest, warn_at, state,
868                     acknowledged_at, acknowledged_by, acknowledged_note, acknowledged_basis)
869                 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)",
870                &[
871                    &tenant,
872                    &key,
873                    &deadline.name,
874                    &deadline.resolved_at.unix_timestamp(),
875                    &deadline.calendar_digest.as_bytes().to_vec(),
876                    &deadline.warn_at.map(time::OffsetDateTime::unix_timestamp),
877                    &deadline.state.as_str(),
878                    &deadline
879                        .acknowledged
880                        .as_ref()
881                        .map(|a| a.at.unix_timestamp()),
882                    &deadline
883                        .acknowledged
884                        .as_ref()
885                        .map(|a| a.by.actor().to_owned()),
886                    &deadline.acknowledged.as_ref().map(|a| a.note.clone()),
887                    &deadline
888                        .acknowledged
889                        .as_ref()
890                        .map(|a| a.by.basis().as_str()),
891                ],
892            )
893            .await
894            .map_err(|e| be(&e))?;
895        }
896
897        // Blob links, stamped from `opened_at` plus an ordinal: the export
898        // carries digests without their link timestamps, the key needs
899        // uniqueness and erasure needs reachability, and neither needs the
900        // original instant.
901        for (at, digest) in blobs.iter().enumerate() {
902            tx.execute(
903                "INSERT INTO case_blobs (tenant, case_id, digest, written_at)
904                 VALUES ($1, $2, $3, $4)",
905                &[
906                    &tenant,
907                    &key,
908                    &digest.as_bytes().to_vec(),
909                    &(case.opened_at.unix_timestamp() + i64::try_from(at).unwrap_or(i64::MAX)),
910                ],
911            )
912            .await
913            .map_err(|e| be(&e))?;
914        }
915
916        tx.commit().await.map_err(|e| be(&e))?;
917        Ok(())
918    }
919
920    async fn attach_run(&self, case: CaseId, run: RunId) -> Result<(), StoreError> {
921        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
922        // Zero-based, matching `import_case`'s enumeration — one rule for what
923        // position the first run holds, whichever door it came through.
924        //
925        // The MAX+1 subquery is not atomic with the insert: two instances
926        // attaching different runs both compute the same next position, and
927        // the `case_runs_order` unique index turns that from two runs silently
928        // sharing a place in the case's history into a violation this loop
929        // retries. The `ON CONFLICT` arm still absorbs only the *idempotence*
930        // conflict — the same run attached twice — because it names the
931        // primary key's columns; a seq collision surfaces as the error the
932        // retry exists for. Bounded, because unbounded politeness under a
933        // pathological writer is a hang: past the bound the collision is
934        // reported rather than spun on.
935        for _ in 0..8 {
936            let result = client
937                .execute(
938                    "INSERT INTO case_runs (case_id, run_id, seq, tenant)
939                     VALUES ($1, $2,
940                             (SELECT COALESCE(MAX(seq) + 1, 0) FROM case_runs
941                               WHERE case_id = $1 AND tenant = $3),
942                             $3)
943                     ON CONFLICT (tenant, case_id, run_id) DO NOTHING",
944                    &[&case.to_string(), &run.to_string(), &self.tenant_name()],
945                )
946                .await;
947            match result {
948                Ok(_) => return Ok(()),
949                Err(e) if e.code() == Some(&SqlState::UNIQUE_VIOLATION) => {}
950                Err(e) => return Err(be(&e)),
951            }
952        }
953        Err(StoreError::Backend(format!(
954            "could not attach run {run} to case {case}: the run-order position was \
955             claimed concurrently on every attempt"
956        )))
957    }
958
959    async fn link_blob(
960        &self,
961        case: CaseId,
962        digest: Digest,
963        at: Timestamp,
964    ) -> Result<(), StoreError> {
965        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
966        // The primary key makes re-linking the same bytes the same record:
967        // two runs on one case writing identical content land on one digest by
968        // construction, and that is one artifact.
969        client
970            .execute(
971                "INSERT INTO case_blobs (case_id, digest, written_at, tenant)
972                 VALUES ($1, $2, $3, $4)
973                 ON CONFLICT (tenant, case_id, digest) DO NOTHING",
974                &[
975                    &case.to_string(),
976                    &digest.as_bytes().as_slice(),
977                    &at.unix_timestamp(),
978                    &self.tenant_name(),
979                ],
980            )
981            .await
982            .map_err(|e| be(&e))?;
983        Ok(())
984    }
985
986    async fn blobs_of(&self, case: CaseId) -> Result<Vec<Digest>, StoreError> {
987        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
988        let rows = client
989            .query(
990                "SELECT digest FROM case_blobs
991                  WHERE case_id = $1 AND tenant = $2 ORDER BY written_at, digest",
992                &[&case.to_string(), &self.tenant_name()],
993            )
994            .await
995            .map_err(|e| be(&e))?;
996        let mut out = Vec::with_capacity(rows.len());
997        for row in rows {
998            let raw: Vec<u8> = row.get(0);
999            let bytes: [u8; 32] = raw.try_into().map_err(|_| StoreError::Corrupt {
1000                seq: 0,
1001                detail: "a linked blob digest is not 32 bytes".into(),
1002            })?;
1003            out.push(Digest::from_bytes(bytes));
1004        }
1005        Ok(out)
1006    }
1007
1008    async fn put_state(
1009        &self,
1010        case: CaseId,
1011        expected: CaseVersion,
1012        state: serde_json::Value,
1013    ) -> Result<CaseVersion, StoreError> {
1014        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1015        let next = expected.next();
1016        // The row count is read, and that is the point. The previous version of
1017        // this method discarded it and returned `Ok(())` for a case that does
1018        // not exist — the same defect already found once in `release` on this
1019        // backend. A guard whose result nobody reads is not a guard.
1020        let n = client
1021            .execute(
1022                "UPDATE cases SET state = $2, version = $3
1023                  WHERE case_id = $1 AND version = $4 AND tenant = $5",
1024                &[
1025                    &case.to_string(),
1026                    &state.to_string(),
1027                    &i64::try_from(next.0).unwrap_or(i64::MAX),
1028                    &i64::try_from(expected.0).unwrap_or(i64::MAX),
1029                    &self.tenant_name(),
1030                ],
1031            )
1032            .await
1033            .map_err(|e| be(&e))?;
1034        if n == 1 {
1035            return Ok(next);
1036        }
1037        // Tell "gone" apart from "moved on": a missing case reported as a
1038        // conflict sends the caller into a re-read loop against nothing.
1039        let current = client
1040            .query_opt(
1041                "SELECT version FROM cases WHERE case_id = $1 AND tenant = $2",
1042                &[&case.to_string(), &self.tenant_name()],
1043            )
1044            .await
1045            .map_err(|e| be(&e))?;
1046        match current {
1047            Some(row) => Err(StoreError::CaseConflict {
1048                case: case.to_string(),
1049                expected: expected.0,
1050                current: u64::try_from(row.get::<_, i64>(0)).unwrap_or(0),
1051            }),
1052            None => Err(StoreError::NotFound(case.to_string())),
1053        }
1054    }
1055
1056    async fn set_status(&self, case: CaseId, status: CaseStatus) -> Result<(), StoreError> {
1057        // Closing releases the correlation keys and refuses an open obligation;
1058        // an ordinary status write does neither. Keeping `set_status(Closed)` a
1059        // bare column write let the `status` column and correlation-open
1060        // membership disagree — a closed case stayed correlatable and a new
1061        // matter attached to it. Route it through `close`.
1062        if status == CaseStatus::Closed {
1063            return self.close(case).await;
1064        }
1065        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1066        let tx = client.transaction().await.map_err(|e| be(&e))?;
1067        // Read the prior status under the row lock, so *which transition this is*
1068        // is decided against the row this transaction is about to write.
1069        let was: Option<String> = tx
1070            .query_opt(
1071                "SELECT status FROM cases WHERE case_id = $1 AND tenant = $2 FOR UPDATE",
1072                &[&case.to_string(), &self.tenant_name()],
1073            )
1074            .await
1075            .map_err(|e| be(&e))?
1076            .map(|r| r.get(0));
1077        let Some(was) = was else {
1078            return Err(StoreError::NotFound(case.to_string()));
1079        };
1080        tx.execute(
1081            "UPDATE cases SET status = $2 WHERE case_id = $1 AND tenant = $3",
1082            &[&case.to_string(), &status.as_str(), &self.tenant_name()],
1083        )
1084        .await
1085        .map_err(|e| be(&e))?;
1086
1087        // Leaving `Closed` takes back the correlation keys closure released, or
1088        // the matter comes back live-looking and unreachable. `open = FALSE`
1089        // rows another case has since claimed are left alone: the identifier
1090        // belongs to whichever matter is open for it now.
1091        if was == CaseStatus::Closed.as_str() {
1092            tx.execute(
1093                "UPDATE case_correlation SET open = TRUE
1094                  WHERE case_id = $1 AND tenant = $2
1095                    AND NOT EXISTS (
1096                          SELECT 1 FROM case_correlation other
1097                           WHERE other.tenant = case_correlation.tenant
1098                             AND other.namespace = case_correlation.namespace
1099                             AND other.value = case_correlation.value
1100                             AND other.open)",
1101                &[&case.to_string(), &self.tenant_name()],
1102            )
1103            .await
1104            .map_err(|e| be(&e))?;
1105        }
1106        tx.commit().await.map_err(|e| be(&e))?;
1107        Ok(())
1108    }
1109
1110    async fn close(&self, case: CaseId) -> Result<(), StoreError> {
1111        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1112        let tx = client.transaction().await.map_err(|e| be(&e))?;
1113
1114        // The case row's lock, taken **before** the count and held to commit.
1115        // Without it this is a check and `register_deadline` is a write that
1116        // walks past it: at READ COMMITTED each statement reads its own
1117        // snapshot, so a registration in flight is invisible here and this
1118        // closure is invisible there, and both commit — leaving a case closed
1119        // and owing. `register_deadline` takes the same lock, which is what
1120        // makes the two decide one at a time.
1121        let locked = tx
1122            .query_opt(
1123                "SELECT status FROM cases WHERE case_id = $1 AND tenant = $2 FOR UPDATE",
1124                &[&case.to_string(), &self.tenant_name()],
1125            )
1126            .await
1127            .map_err(|e| be(&e))?;
1128        if locked.is_none() {
1129            return Err(StoreError::NotFound(case.to_string()));
1130        }
1131
1132        // An unmet obligation survives closure, because closure is the moment
1133        // people stop looking.
1134        let open: i64 = tx
1135            .query_one(
1136                "SELECT COUNT(*) FROM case_deadlines
1137                  WHERE case_id = $1 AND tenant = $2 AND state = ANY($3::text[])",
1138                &[
1139                    &case.to_string(),
1140                    &self.tenant_name(),
1141                    &deadline_states(DeadlineState::is_open),
1142                ],
1143            )
1144            .await
1145            .map_err(|e| be(&e))?
1146            .get(0);
1147        if open > 0 {
1148            return Err(StoreError::ObligationsOutstanding {
1149                case: case.to_string(),
1150                outstanding: usize::try_from(open).unwrap_or(usize::MAX),
1151            });
1152        }
1153
1154        // The row count is read: closing a case that does not exist reported
1155        // success, while the redb backend answers `NotFound` — and the caller
1156        // of a silent close believes an audited matter was settled.
1157        let n = tx
1158            .execute(
1159                "UPDATE cases SET status = $2 WHERE case_id = $1 AND tenant = $3",
1160                &[
1161                    &case.to_string(),
1162                    &CaseStatus::Closed.as_str(),
1163                    &self.tenant_name(),
1164                ],
1165            )
1166            .await
1167            .map_err(|e| be(&e))?;
1168        if n == 0 {
1169            return Err(StoreError::NotFound(case.to_string()));
1170        }
1171        // Releasing the keys is what lets a later message open a *new* matter
1172        // rather than reanimating an audited one.
1173        tx.execute(
1174            "UPDATE case_correlation SET open = FALSE WHERE case_id = $1 AND tenant = $2",
1175            &[&case.to_string(), &self.tenant_name()],
1176        )
1177        .await
1178        .map_err(|e| be(&e))?;
1179        tx.commit().await.map_err(|e| be(&e))?;
1180        Ok(())
1181    }
1182
1183    async fn register_deadline(&self, d: &Deadline) -> Result<(), StoreError> {
1184        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1185        let tx = client.transaction().await.map_err(|e| be(&e))?;
1186
1187        // The same row lock `close` takes, for the same reason: these two are
1188        // the pair that has to be serialized, and the anomaly they would
1189        // otherwise produce — both committing on each other's pre-state — is a
1190        // case closed with an obligation it will go on to miss.
1191        let status: Option<String> = tx
1192            .query_opt(
1193                "SELECT status FROM cases WHERE case_id = $1 AND tenant = $2 FOR UPDATE",
1194                &[&d.case.to_string(), &self.tenant_name()],
1195            )
1196            .await
1197            .map_err(|e| be(&e))?
1198            .map(|r| r.get(0));
1199        let Some(status) = status else {
1200            return Err(StoreError::NotFound(d.case.to_string()));
1201        };
1202        if status == CaseStatus::Closed.as_str() {
1203            return Err(StoreError::CaseClosed {
1204                case: d.case.to_string(),
1205            });
1206        }
1207
1208        tx.execute(
1209            "INSERT INTO case_deadlines
1210                   (case_id, name, resolved_at, calendar_digest, warn_at, state, tenant,
1211                    acknowledged_at, acknowledged_by, acknowledged_note, acknowledged_basis)
1212                 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
1213                 ON CONFLICT (tenant, case_id, name) DO NOTHING",
1214            &[
1215                &d.case.to_string(),
1216                &d.name,
1217                &d.resolved_at.unix_timestamp(),
1218                &d.calendar_digest.as_bytes().to_vec(),
1219                &d.warn_at.map(Timestamp::unix_timestamp),
1220                &d.state.as_str(),
1221                &self.tenant_name(),
1222                &d.acknowledged.as_ref().map(|a| a.at.unix_timestamp()),
1223                &d.acknowledged.as_ref().map(|a| a.by.actor().to_owned()),
1224                &d.acknowledged.as_ref().map(|a| a.note.clone()),
1225                &d.acknowledged.as_ref().map(|a| a.by.basis().as_str()),
1226            ],
1227        )
1228        .await
1229        .map_err(|e| be(&e))?;
1230        tx.commit().await.map_err(|e| be(&e))?;
1231        Ok(())
1232    }
1233
1234    async fn deadlines(&self, case: CaseId) -> Result<Vec<Deadline>, StoreError> {
1235        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1236        let rows = client
1237            .query(
1238                "SELECT case_id, name, resolved_at, calendar_digest, warn_at, state,
1239                        acknowledged_at, acknowledged_by, acknowledged_note, acknowledged_basis
1240                   FROM case_deadlines
1241                  WHERE case_id = $1 AND tenant = $2 ORDER BY resolved_at ASC",
1242                &[&case.to_string(), &self.tenant_name()],
1243            )
1244            .await
1245            .map_err(|e| be(&e))?;
1246        rows.iter().map(deadline_from).collect()
1247    }
1248
1249    async fn set_deadline_state(
1250        &self,
1251        case: CaseId,
1252        name: &str,
1253        state: DeadlineState,
1254    ) -> Result<(), StoreError> {
1255        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1256        let tx = client.transaction().await.map_err(|e| be(&e))?;
1257
1258        // How an obligation ended is not editable, and the row's lock is what
1259        // decides that against the row this transaction writes rather than
1260        // against one a concurrent sweep has already moved.
1261        let current: Option<String> = tx
1262            .query_opt(
1263                "SELECT state FROM case_deadlines
1264                  WHERE case_id = $1 AND name = $2 AND tenant = $3 FOR UPDATE",
1265                &[&case.to_string(), &name.to_owned(), &self.tenant_name()],
1266            )
1267            .await
1268            .map_err(|e| be(&e))?
1269            .map(|r| r.get(0));
1270        if let Some(from) = current {
1271            let parsed = deadline_state_from(&from)?;
1272            if !parsed.may_become(state) {
1273                return Err(StoreError::DeadlineFinal {
1274                    case: case.to_string(),
1275                    obligation: name.to_owned(),
1276                    from,
1277                    to: state.as_str().to_owned(),
1278                });
1279            }
1280        }
1281
1282        // The row count is read: a transition for an obligation that does not
1283        // exist reported success, while the redb backend answers `NotFound` —
1284        // and the sweep that believes it breached a deadline nobody registered
1285        // has written its decision into nothing.
1286        let n = tx
1287            .execute(
1288                "UPDATE case_deadlines SET state = $3
1289                  WHERE case_id = $1 AND name = $2 AND tenant = $4",
1290                &[
1291                    &case.to_string(),
1292                    &name.to_owned(),
1293                    &state.as_str(),
1294                    &self.tenant_name(),
1295                ],
1296            )
1297            .await
1298            .map_err(|e| be(&e))?;
1299        if n == 0 {
1300            return Err(StoreError::NotFound(format!("{case}/{name}")));
1301        }
1302        tx.commit().await.map_err(|e| be(&e))?;
1303        Ok(())
1304    }
1305
1306    async fn breach_deadline(
1307        &self,
1308        case: CaseId,
1309        name: &str,
1310        now: Timestamp,
1311    ) -> Result<bool, StoreError> {
1312        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1313        let tx = client.transaction().await.map_err(|e| be(&e))?;
1314        let (key, tenant) = (case.to_string(), self.tenant_name());
1315        // The case row's lock first, as `close` takes it: closure and this
1316        // breach then decide one at a time, and a closure that won has already
1317        // met or cancelled the obligation this checks.
1318        let status: Option<String> = tx
1319            .query_opt(
1320                "SELECT status FROM cases WHERE case_id = $1 AND tenant = $2 FOR UPDATE",
1321                &[&key, &tenant],
1322            )
1323            .await
1324            .map_err(|e| be(&e))?
1325            .map(|r| r.get(0));
1326        let Some(status) = status else {
1327            return Err(StoreError::NotFound(key));
1328        };
1329        let breached = tx
1330            .execute(
1331                "UPDATE case_deadlines SET state = $4, breach_unnoted = TRUE
1332                  WHERE case_id = $1 AND name = $2 AND tenant = $3
1333                    AND state = ANY($5::text[]) AND resolved_at <= $6",
1334                &[
1335                    &key,
1336                    &name.to_owned(),
1337                    &tenant,
1338                    &DeadlineState::Breached.as_str(),
1339                    &deadline_states(DeadlineState::is_open),
1340                    &now.unix_timestamp(),
1341                ],
1342            )
1343            .await
1344            .map_err(|e| be(&e))?;
1345        if breached == 0 {
1346            let exists = tx
1347                .query_opt(
1348                    "SELECT 1 FROM case_deadlines
1349                      WHERE case_id = $1 AND name = $2 AND tenant = $3",
1350                    &[&key, &name.to_owned(), &tenant],
1351                )
1352                .await
1353                .map_err(|e| be(&e))?
1354                .is_some();
1355            tx.commit().await.map_err(|e| be(&e))?;
1356            return if exists {
1357                Ok(false)
1358            } else {
1359                Err(StoreError::NotFound(format!("{case}/{name}")))
1360            };
1361        }
1362        if status != CaseStatus::Closed.as_str() {
1363            tx.execute(
1364                "UPDATE cases SET status = $3 WHERE case_id = $1 AND tenant = $2",
1365                &[&key, &tenant, &CaseStatus::Escalated.as_str()],
1366            )
1367            .await
1368            .map_err(|e| be(&e))?;
1369        }
1370        tx.commit().await.map_err(|e| be(&e))?;
1371        Ok(true)
1372    }
1373
1374    async fn due(&self, now: Timestamp, limit: usize) -> Result<Vec<Deadline>, StoreError> {
1375        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1376        let rows = client
1377            .query(
1378                "SELECT case_id, name, resolved_at, calendar_digest, warn_at, state,
1379                        acknowledged_at, acknowledged_by, acknowledged_note, acknowledged_basis
1380                   FROM case_deadlines
1381                  WHERE tenant = $3 AND state = ANY($4::text[])
1382                    AND (resolved_at <= $1 OR (warn_at IS NOT NULL AND warn_at <= $1))
1383                  ORDER BY resolved_at ASC LIMIT $2",
1384                &[
1385                    &now.unix_timestamp(),
1386                    &i64::try_from(limit).unwrap_or(i64::MAX),
1387                    &self.tenant_name(),
1388                    &deadline_states(DeadlineState::is_open),
1389                ],
1390            )
1391            .await
1392            .map_err(|e| be(&e))?;
1393        rows.iter().map(deadline_from).collect()
1394    }
1395
1396    async fn breaches_to_note(&self, limit: usize) -> Result<Vec<Deadline>, StoreError> {
1397        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1398        let rows = client
1399            .query(
1400                "SELECT case_id, name, resolved_at, calendar_digest, warn_at, state,
1401                        acknowledged_at, acknowledged_by, acknowledged_note, acknowledged_basis
1402                   FROM case_deadlines
1403                  WHERE tenant = $2 AND breach_unnoted
1404                  ORDER BY resolved_at ASC LIMIT $1",
1405                &[
1406                    &i64::try_from(limit).unwrap_or(i64::MAX),
1407                    &self.tenant_name(),
1408                ],
1409            )
1410            .await
1411            .map_err(|e| be(&e))?;
1412        rows.iter().map(deadline_from).collect()
1413    }
1414
1415    async fn mark_breach_noted(&self, case: CaseId, name: &str) -> Result<(), StoreError> {
1416        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1417        let n = client
1418            .execute(
1419                "UPDATE case_deadlines SET breach_unnoted = FALSE
1420                  WHERE tenant = $1 AND case_id = $2 AND name = $3",
1421                &[&self.tenant_name(), &case.to_string(), &name.to_owned()],
1422            )
1423            .await
1424            .map_err(|e| be(&e))?;
1425        if n == 0 {
1426            return Err(StoreError::NotFound(format!("{case}/{name}")));
1427        }
1428        Ok(())
1429    }
1430
1431    async fn breached(&self, limit: usize) -> Result<Vec<Deadline>, StoreError> {
1432        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1433        // Served by `case_deadlines_due`, whose leading columns are the two this
1434        // filters on. No clock: a breach is a recorded state, not a comparison
1435        // against now.
1436        let rows = client
1437            .query(
1438                "SELECT case_id, name, resolved_at, calendar_digest, warn_at, state,
1439                        acknowledged_at, acknowledged_by, acknowledged_note, acknowledged_basis
1440                   FROM case_deadlines
1441                  WHERE tenant = $2 AND state = 'breached' AND acknowledged_at IS NULL
1442                  ORDER BY resolved_at ASC LIMIT $1",
1443                &[
1444                    &i64::try_from(limit).unwrap_or(i64::MAX),
1445                    &self.tenant_name(),
1446                ],
1447            )
1448            .await
1449            .map_err(|e| be(&e))?;
1450        rows.iter().map(deadline_from).collect()
1451    }
1452
1453    async fn acknowledge_breach(
1454        &self,
1455        case: CaseId,
1456        name: &str,
1457        note: &BreachNote,
1458    ) -> Result<bool, StoreError> {
1459        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1460        let tx = client.transaction().await.map_err(|e| be(&e))?;
1461        let row = tx
1462            .query_opt(
1463                "SELECT state, acknowledged_at FROM case_deadlines
1464                  WHERE case_id = $1 AND name = $2 AND tenant = $3 FOR UPDATE",
1465                &[&case.to_string(), &name.to_owned(), &self.tenant_name()],
1466            )
1467            .await
1468            .map_err(|e| be(&e))?;
1469        let Some(row) = row else {
1470            return Err(StoreError::NotFound(format!("{case}/{name}")));
1471        };
1472        let state: String = row.get(0);
1473        if state != DeadlineState::Breached.as_str() {
1474            return Err(StoreError::NotBreached {
1475                case: case.to_string(),
1476                obligation: name.to_owned(),
1477                state,
1478            });
1479        }
1480        // The first account stands: a retry must not rewrite who looked or when.
1481        if row.get::<_, Option<i64>>(1).is_some() {
1482            return Ok(false);
1483        }
1484        tx.execute(
1485            "UPDATE case_deadlines
1486                SET acknowledged_at = $4, acknowledged_by = $5, acknowledged_note = $6,
1487                    acknowledged_basis = $7
1488              WHERE case_id = $1 AND name = $2 AND tenant = $3",
1489            &[
1490                &case.to_string(),
1491                &name.to_owned(),
1492                &self.tenant_name(),
1493                &note.at.unix_timestamp(),
1494                &note.by.actor(),
1495                &note.note,
1496                &note.by.basis().as_str(),
1497            ],
1498        )
1499        .await
1500        .map_err(|e| be(&e))?;
1501        tx.commit().await.map_err(|e| be(&e))?;
1502        Ok(true)
1503    }
1504
1505    async fn place_hold(&self, case: CaseId, hold: &LegalHold) -> Result<bool, StoreError> {
1506        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1507        let tx = client.transaction().await.map_err(|e| be(&e))?;
1508        // The case row's lock, for the reason `register_deadline` takes it: two
1509        // snapshots each reading the other's pre-state is how a hold lands on a
1510        // case that an erasure in the other transaction is already destroying.
1511        let exists = tx
1512            .query_opt(
1513                "SELECT 1 FROM cases WHERE tenant = $1 AND case_id = $2 FOR UPDATE",
1514                &[&self.tenant_name(), &case.to_string()],
1515            )
1516            .await
1517            .map_err(|e| be(&e))?;
1518        if exists.is_none() {
1519            return Err(StoreError::NotFound(case.to_string()));
1520        }
1521        // First placement wins: `DO NOTHING` rather than `DO UPDATE`, so a retry
1522        // cannot move the instant or rewrite the reason.
1523        let placed = tx
1524            .execute(
1525                "INSERT INTO case_legal_holds
1526                      (tenant, case_id, placed_at, reason, by_actor, by_basis)
1527                      VALUES ($1, $2, $3, $4, $5, $6)
1528                 ON CONFLICT (tenant, case_id) DO NOTHING",
1529                &[
1530                    &self.tenant_name(),
1531                    &case.to_string(),
1532                    &hold.placed_at.unix_timestamp(),
1533                    &hold.reason,
1534                    &hold.by.actor(),
1535                    &hold.by.basis().as_str(),
1536                ],
1537            )
1538            .await
1539            .map_err(|e| be(&e))?;
1540        tx.commit().await.map_err(|e| be(&e))?;
1541        Ok(placed == 1)
1542    }
1543
1544    async fn release_hold(&self, case: CaseId) -> Result<bool, StoreError> {
1545        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1546        let exists = client
1547            .query_opt(
1548                "SELECT 1 FROM cases WHERE tenant = $1 AND case_id = $2",
1549                &[&self.tenant_name(), &case.to_string()],
1550            )
1551            .await
1552            .map_err(|e| be(&e))?;
1553        if exists.is_none() {
1554            return Err(StoreError::NotFound(case.to_string()));
1555        }
1556        let removed = client
1557            .execute(
1558                "DELETE FROM case_legal_holds WHERE tenant = $1 AND case_id = $2",
1559                &[&self.tenant_name(), &case.to_string()],
1560            )
1561            .await
1562            .map_err(|e| be(&e))?;
1563        Ok(removed == 1)
1564    }
1565
1566    async fn release_hold_if(
1567        &self,
1568        case: CaseId,
1569        standing: &LegalHold,
1570    ) -> Result<bool, StoreError> {
1571        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1572        let exists = client
1573            .query_opt(
1574                "SELECT 1 FROM cases WHERE tenant = $1 AND case_id = $2",
1575                &[&self.tenant_name(), &case.to_string()],
1576            )
1577            .await
1578            .map_err(|e| be(&e))?;
1579        if exists.is_none() {
1580            return Err(StoreError::NotFound(case.to_string()));
1581        }
1582        // The compare is the statement's own predicate, so a hold placed
1583        // after the caller read the old one fails it and stays.
1584        let removed = client
1585            .execute(
1586                "DELETE FROM case_legal_holds
1587                  WHERE tenant = $1 AND case_id = $2 AND placed_at = $3
1588                    AND reason = $4 AND by_actor = $5 AND by_basis = $6",
1589                &[
1590                    &self.tenant_name(),
1591                    &case.to_string(),
1592                    &standing.placed_at.unix_timestamp(),
1593                    &standing.reason,
1594                    &standing.by.actor(),
1595                    &standing.by.basis().as_str(),
1596                ],
1597            )
1598            .await
1599            .map_err(|e| be(&e))?;
1600        Ok(removed == 1)
1601    }
1602
1603    async fn hold(&self, case: CaseId) -> Result<Option<LegalHold>, StoreError> {
1604        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1605        let row = client
1606            .query_opt(
1607                "SELECT placed_at, reason, by_actor, by_basis FROM case_legal_holds
1608                  WHERE tenant = $1 AND case_id = $2",
1609                &[&self.tenant_name(), &case.to_string()],
1610            )
1611            .await
1612            .map_err(|e| be(&e))?;
1613        row.map(|r| {
1614            Ok(super::hold_from_row(
1615                Timestamp::from_unix_timestamp(r.get::<_, i64>(0))
1616                    .map_err(|e| corrupt("unrepresentable placed_at", e))?,
1617                hold_row(r.get(1), &r.get::<_, String>(2), &r.get::<_, String>(3))?,
1618            ))
1619        })
1620        .transpose()
1621    }
1622
1623    async fn holds(
1624        &self,
1625        after: Option<CaseId>,
1626        limit: usize,
1627    ) -> Result<Vec<(CaseId, LegalHold)>, StoreError> {
1628        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1629        // `(placed_at, case_id)` as one tuple is what makes the cursor total:
1630        // several holds can share an instant, and comparing the stamp alone
1631        // would either repeat them or skip them depending on which side of the
1632        // boundary the page fell.
1633        let rows = match after {
1634            Some(cursor) => client
1635                .query(
1636                    "SELECT case_id, placed_at, reason, by_actor, by_basis FROM case_legal_holds
1637                          WHERE tenant = $1
1638                            AND (placed_at, case_id) > (
1639                                SELECT placed_at, case_id FROM case_legal_holds
1640                                 WHERE tenant = $1 AND case_id = $2)
1641                          ORDER BY placed_at, case_id
1642                          LIMIT $3",
1643                    &[
1644                        &self.tenant_name(),
1645                        &cursor.to_string(),
1646                        &i64::try_from(limit).unwrap_or(i64::MAX),
1647                    ],
1648                )
1649                .await,
1650            None => client
1651                .query(
1652                    "SELECT case_id, placed_at, reason, by_actor, by_basis FROM case_legal_holds
1653                          WHERE tenant = $1
1654                          ORDER BY placed_at, case_id
1655                          LIMIT $2",
1656                    &[
1657                        &self.tenant_name(),
1658                        &i64::try_from(limit).unwrap_or(i64::MAX),
1659                    ],
1660                )
1661                .await,
1662        }
1663        .map_err(|e| be(&e))?;
1664
1665        rows.into_iter()
1666            .map(|r| {
1667                Ok((
1668                    CaseId::parse(&r.get::<_, String>(0)).map_err(|e| corrupt("bad case id", e))?,
1669                    super::hold_from_row(
1670                        Timestamp::from_unix_timestamp(r.get::<_, i64>(1))
1671                            .map_err(|e| corrupt("unrepresentable placed_at", e))?,
1672                        hold_row(r.get(2), &r.get::<_, String>(3), &r.get::<_, String>(4))?,
1673                    ),
1674                ))
1675            })
1676            .collect()
1677    }
1678
1679    async fn by_status(&self, status: CaseStatus, limit: usize) -> Result<Vec<Case>, StoreError> {
1680        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1681        let rows = client
1682            .query(
1683                "SELECT case_id FROM cases
1684                  WHERE status = $1 AND tenant = $3 ORDER BY opened_at DESC LIMIT $2",
1685                &[
1686                    &status.as_str(),
1687                    &i64::try_from(limit).unwrap_or(i64::MAX),
1688                    &self.tenant_name(),
1689                ],
1690            )
1691            .await
1692            .map_err(|e| be(&e))?;
1693        let mut out = Vec::with_capacity(rows.len());
1694        for row in rows {
1695            let id: String = row.get(0);
1696            let id = CaseId::parse(&id).map_err(|e| corrupt("bad case id", e))?;
1697            if let Some(c) = self.case(id).await? {
1698                out.push(c);
1699            }
1700        }
1701        Ok(out)
1702    }
1703
1704    async fn record_drill(&self, record: &crate::case::DrillRecord) -> Result<(), StoreError> {
1705        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1706        client
1707            .execute(
1708                "INSERT INTO case_last_drill
1709                   (tenant, ran_at, sound, cases, findings, not_checked, origin, log_size)
1710                 VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
1711                 ON CONFLICT (tenant) DO UPDATE SET
1712                   ran_at = EXCLUDED.ran_at, sound = EXCLUDED.sound,
1713                   cases = EXCLUDED.cases, findings = EXCLUDED.findings,
1714                   not_checked = EXCLUDED.not_checked, origin = EXCLUDED.origin,
1715                   log_size = EXCLUDED.log_size",
1716                &[
1717                    &self.tenant_name(),
1718                    &record.at.unix_timestamp(),
1719                    &record.sound,
1720                    &i64::try_from(record.cases).unwrap_or(i64::MAX),
1721                    &i64::try_from(record.findings).unwrap_or(i64::MAX),
1722                    &i64::try_from(record.not_checked).unwrap_or(i64::MAX),
1723                    &record.origin,
1724                    &i64::try_from(record.size).unwrap_or(i64::MAX),
1725                ],
1726            )
1727            .await
1728            .map_err(|e| be(&e))?;
1729        Ok(())
1730    }
1731
1732    async fn last_drill(&self) -> Result<Option<crate::case::DrillRecord>, StoreError> {
1733        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1734        let row = client
1735            .query_opt(
1736                "SELECT ran_at, sound, cases, findings, not_checked, origin, log_size
1737                   FROM case_last_drill WHERE tenant = $1",
1738                &[&self.tenant_name()],
1739            )
1740            .await
1741            .map_err(|e| be(&e))?;
1742        row.map(|r| {
1743            Ok(crate::case::DrillRecord {
1744                at: Timestamp::from_unix_timestamp(r.get::<_, i64>(0))
1745                    .map_err(|e| corrupt("unrepresentable drill instant", e))?,
1746                sound: r.get(1),
1747                cases: u64::try_from(r.get::<_, i64>(2)).unwrap_or(0),
1748                findings: u64::try_from(r.get::<_, i64>(3)).unwrap_or(0),
1749                not_checked: u64::try_from(r.get::<_, i64>(4)).unwrap_or(0),
1750                origin: r.get(5),
1751                size: u64::try_from(r.get::<_, i64>(6)).unwrap_or(0),
1752            })
1753        })
1754        .transpose()
1755    }
1756
1757    async fn census(&self, now: Timestamp) -> Result<CaseCensus, StoreError> {
1758        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1759        let row = client
1760            .query_one(
1761                "SELECT COUNT(*), MIN(opened_at) FROM cases
1762                  WHERE tenant = $1 AND status <> $2::text",
1763                &[&self.tenant_name(), &CaseStatus::Closed.as_str()],
1764            )
1765            .await
1766            .map_err(|e| be(&e))?;
1767        let open: i64 = row.get(0);
1768        let oldest: Option<i64> = row.get(1);
1769        let due: i64 = client
1770            .query_one(
1771                "SELECT COUNT(*) FROM case_deadlines
1772                  WHERE tenant = $2 AND state = ANY($3::text[]) AND resolved_at <= $1",
1773                &[
1774                    &now.unix_timestamp(),
1775                    &self.tenant_name(),
1776                    &deadline_states(DeadlineState::is_open),
1777                ],
1778            )
1779            .await
1780            .map_err(|e| be(&e))?
1781            .get(0);
1782
1783        let breached: i64 = client
1784            .query_one(
1785                "SELECT COUNT(*) FROM case_deadlines
1786                  WHERE tenant = $1 AND state = 'breached' AND acknowledged_at IS NULL",
1787                &[&self.tenant_name()],
1788            )
1789            .await
1790            .map_err(|e| be(&e))?
1791            .get(0);
1792
1793        Ok(CaseCensus {
1794            open: u64::try_from(open).unwrap_or(0),
1795            oldest_age_secs: oldest.map(|o| {
1796                crate::runtime::metrics::age_secs(
1797                    Timestamp::from_unix_timestamp(o).unwrap_or(now),
1798                    now,
1799                )
1800            }),
1801            due: u64::try_from(due).unwrap_or(0),
1802            breached: u64::try_from(breached).unwrap_or(0),
1803        })
1804    }
1805}
1806
1807fn deadline_from(row: &tokio_postgres::Row) -> Result<Deadline, StoreError> {
1808    let case: String = row.get(0);
1809    let digest: Vec<u8> = row.get(3);
1810    let warn: Option<i64> = row.get(4);
1811    let state: String = row.get(5);
1812    let arr: [u8; 32] = digest
1813        .try_into()
1814        .map_err(|_| corrupt("calendar digest", "not 32 bytes"))?;
1815    Ok(Deadline {
1816        case: CaseId::parse(&case).map_err(|e| corrupt("bad case id", e))?,
1817        name: row.get(1),
1818        resolved_at: Timestamp::from_unix_timestamp(row.get::<_, i64>(2))
1819            .map_err(|e| corrupt("unrepresentable deadline", e))?,
1820        calendar_digest: Digest::from_bytes(arr),
1821        warn_at: warn
1822            .map(Timestamp::from_unix_timestamp)
1823            .transpose()
1824            .map_err(|e| corrupt("unrepresentable warn_at", e))?,
1825        state: deadline_state_from(&state)?,
1826        acknowledged: acknowledged_from(row)?,
1827    })
1828}
1829
1830/// The account's three columns, read as one optional fact.
1831///
1832/// Fallible rather than defaulting: a row whose `acknowledged_at` is set and
1833/// whose `acknowledged_by` is not is a corrupt row, and reading it as *nobody
1834/// has looked* would put a breach back on a listing somebody had already
1835/// answered — or, the other way, hide one nobody had.
1836fn acknowledged_from(row: &tokio_postgres::Row) -> Result<Option<BreachNote>, StoreError> {
1837    let at: Option<i64> = row.get(6);
1838    let by: Option<String> = row.get(7);
1839    let note: Option<String> = row.get(8);
1840    let basis: Option<String> = row.get(9);
1841    match (at, by, note, basis) {
1842        (None, None, None, None) => Ok(None),
1843        (Some(at), Some(by), Some(note), Some(basis)) => Ok(Some(BreachNote {
1844            by: super::decode_operator(&by, &basis, "case_deadlines")?,
1845            note,
1846            at: Timestamp::from_unix_timestamp(at)
1847                .map_err(|e| corrupt("unrepresentable acknowledged_at", e))?,
1848        })),
1849        _ => Err(corrupt(
1850            "breach acknowledgement",
1851            "only some of its columns are set",
1852        )),
1853    }
1854}
1855
1856#[async_trait]
1857#[allow(clippy::too_many_lines)]
1858impl EventStore for PostgresStore {
1859    fn tenant(&self) -> &str {
1860        self.tenant_str()
1861    }
1862
1863    async fn buffer(&self, event: &InboundEvent, at: Timestamp) -> Result<bool, StoreError> {
1864        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1865        let tx = client.transaction().await.map_err(|e| be(&e))?;
1866        let inserted = tx
1867            .execute(
1868                "INSERT INTO inbound_events
1869                   (event_id, source, bare_id, kind, payload, received_at, tenant,
1870                    by_actor, by_basis)
1871                 VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
1872                 ON CONFLICT (tenant, event_id) DO NOTHING",
1873                &[
1874                    &event.dedup_key(),
1875                    &event.source,
1876                    &event.id,
1877                    &event.kind,
1878                    &event.payload.to_string(),
1879                    &at.unix_timestamp(),
1880                    &self.tenant_name(),
1881                    &event.by.as_ref().map(|o| o.actor().to_owned()),
1882                    &event.by.as_ref().map(|o| o.basis().as_str().to_owned()),
1883                ],
1884            )
1885            .await
1886            .map_err(|e| be(&e))?;
1887        if inserted == 0 {
1888            tx.commit().await.map_err(|e| be(&e))?;
1889            return Ok(false);
1890        }
1891        for k in &event.correlation {
1892            tx.execute(
1893                "INSERT INTO inbound_correlation (event_id, namespace, value, tenant)
1894                 VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING",
1895                // The dedup key throughout: correlation rows reference the row
1896                // `buffer` wrote, and that is keyed by `(source, id)`.
1897                &[
1898                    &event.dedup_key(),
1899                    &k.namespace,
1900                    &k.value,
1901                    &self.tenant_name(),
1902                ],
1903            )
1904            .await
1905            .map_err(|e| be(&e))?;
1906        }
1907        tx.commit().await.map_err(|e| be(&e))?;
1908        Ok(true)
1909    }
1910
1911    async fn subscribe(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError> {
1912        // One transaction for the whole wait. A subscription with several
1913        // correlation keys is one wait, and the per-statement autocommit this
1914        // replaces could fail halfway — a wait reachable by some of its keys
1915        // and not others, which matches or misses depending on which key the
1916        // event happens to carry. The redb backend writes them in one
1917        // transaction; so does this.
1918        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1919        let tx = client.transaction().await.map_err(|e| be(&e))?;
1920        for k in &sub.correlation {
1921            tx.execute(
1922                "INSERT INTO subscriptions
1923                       (run_id, effect_key, case_id, step, phase, event_kind,
1924                        namespace, value, created_at, tenant, from_source)
1925                     VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
1926                     ON CONFLICT DO NOTHING",
1927                &[
1928                    &sub.run.to_string(),
1929                    &sub.effect.to_hex(),
1930                    &sub.case.map(|c| c.to_string()),
1931                    &i64::from(sub.step.0),
1932                    &crate::core::Phase::as_str(sub.phase),
1933                    &sub.kind,
1934                    &k.namespace,
1935                    &k.value,
1936                    &at.unix_timestamp(),
1937                    &self.tenant_name(),
1938                    &sub.from,
1939                ],
1940            )
1941            .await
1942            .map_err(|e| be(&e))?;
1943        }
1944        tx.commit().await.map_err(|e| be(&e))?;
1945        Ok(())
1946    }
1947
1948    async fn claim_for(
1949        &self,
1950        sub: &Subscription,
1951        at: Timestamp,
1952    ) -> Result<Option<BufferedEvent>, StoreError> {
1953        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
1954        for k in &sub.correlation {
1955            // One statement. The claim predicate, the write and the wait's
1956            // parking are evaluated together, so two waiters cannot both come
1957            // away with the row and the wait takes no second message.
1958            //
1959            // Unclaimed — or claimed **for this very wait** and not yet
1960            // consumed. The second arm is crash recovery: `match_waiter`
1961            // claims durably and the run resumes in a separate step, so a
1962            // crash between the two leaves an event claimed for a wait that
1963            // never saw it. Scoped to the wait and to a standing claim, so a
1964            // message the run already journaled is never handed to its next
1965            // wait on the same key.
1966            let row = client
1967                .query_opt(
1968                    "WITH claimed AS (
1969                       UPDATE inbound_events SET claimed_by = $1, claimed_at = $2,
1970                                                 claimed_for = $8
1971                        WHERE tenant = $6 AND event_id = (
1972                            SELECT e.event_id FROM inbound_events e
1973                              JOIN inbound_correlation c
1974                                ON c.tenant = e.tenant AND c.event_id = e.event_id
1975                             WHERE e.tenant = $6 AND e.kind = $3
1976                               AND c.namespace = $4 AND c.value = $5
1977                               AND ($7::TEXT IS NULL OR e.source = $7)
1978                               AND ((e.claimed_by IS NULL
1979                                     AND NOT EXISTS (
1980                                         SELECT 1 FROM subscriptions p
1981                                          WHERE p.tenant = $6 AND p.run_id = $1
1982                                            AND p.effect_key = $8
1983                                            AND p.parked_at IS NOT NULL))
1984                                    OR (e.claimed_by = $1 AND e.claimed_for = $8))
1985                               AND NOT e.dead AND NOT e.erased
1986                             ORDER BY e.received_at ASC
1987                             FOR UPDATE SKIP LOCKED
1988                             LIMIT 1)
1989                    RETURNING bare_id, kind, payload, received_at, source, event_id,
1990                              by_actor, by_basis),
1991                     parked AS (
1992                       UPDATE subscriptions SET parked_at = $2
1993                        WHERE tenant = $6 AND run_id = $1 AND effect_key = $8
1994                          AND EXISTS (SELECT 1 FROM claimed))
1995                   SELECT * FROM claimed",
1996                    &[
1997                        &sub.run.to_string(),
1998                        &at.unix_timestamp(),
1999                        &sub.kind,
2000                        &k.namespace,
2001                        &k.value,
2002                        &self.tenant_name(),
2003                        &sub.from,
2004                        &sub.effect.to_hex(),
2005                    ],
2006                )
2007                .await
2008                .map_err(|e| be(&e))?;
2009            if let Some(row) = row {
2010                return Ok(Some(
2011                    buffered_from(&row, &client, &self.tenant_name()).await?,
2012                ));
2013            }
2014        }
2015        Ok(None)
2016    }
2017
2018    async fn match_waiter(
2019        &self,
2020        event: &InboundEvent,
2021        at: Timestamp,
2022    ) -> Result<Option<Subscription>, StoreError> {
2023        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2024        let tx = client.transaction().await.map_err(|e| be(&e))?;
2025
2026        for k in &event.correlation {
2027            // `FOR UPDATE` on the subscription row, for the race `deliver_to`
2028            // already locks against: two *different* events for one wait. The
2029            // claim below guards the event row, so one event never resumed two
2030            // runs — but nothing serialised two events electing one waiter.
2031            // Both selected it, both claimed their own (unclaimed) event, and
2032            // the loser's event was left claimed for a run whose wait the
2033            // winner had already satisfied: parked forever, because the sweep
2034            // only retires rows with no claim. Locking the row makes the loser
2035            // wait; the winner's retirement then deletes it, and the loser's
2036            // re-evaluation finds no waiter and leaves its event live.
2037            //
2038            // A parked wait is passed over: it already holds its claimed
2039            // event and waits only for redelivery, so a second event elected
2040            // for it would be claimed for a satisfied wait and never consumed.
2041            let Some(row) = tx
2042                .query_opt(
2043                    "SELECT run_id, effect_key, case_id, step, phase, from_source
2044                       FROM subscriptions
2045                      WHERE tenant = $4 AND event_kind = $1
2046                        AND namespace = $2 AND value = $3
2047                        AND (from_source IS NULL OR from_source = $5)
2048                        AND parked_at IS NULL
2049                        AND NOT EXISTS (
2050                              SELECT 1 FROM run_seal s
2051                               WHERE s.tenant = subscriptions.tenant
2052                                 AND s.run_id = subscriptions.run_id)
2053                      ORDER BY created_at ASC LIMIT 1
2054                      FOR UPDATE OF subscriptions",
2055                    &[
2056                        &event.kind,
2057                        &k.namespace,
2058                        &k.value,
2059                        &self.tenant_name(),
2060                        &event.source,
2061                    ],
2062                )
2063                .await
2064                .map_err(|e| be(&e))?
2065            else {
2066                continue;
2067            };
2068            let run: String = row.get(0);
2069            let effect: String = row.get(1);
2070
2071            // Claiming the event in the same transaction is what stops one
2072            // message resuming two runs.
2073            let claimed = tx
2074                .execute(
2075                    "UPDATE inbound_events SET claimed_by = $2, claimed_at = $3,
2076                                              claimed_for = $5
2077                      WHERE tenant = $4 AND event_id = $1
2078                        AND claimed_by IS NULL AND NOT dead",
2079                    &[
2080                        &event.dedup_key(),
2081                        &run,
2082                        &at.unix_timestamp(),
2083                        &self.tenant_name(),
2084                        &effect,
2085                    ],
2086                )
2087                .await
2088                .map_err(|e| be(&e))?;
2089            if claimed == 0 {
2090                tx.commit().await.map_err(|e| be(&e))?;
2091                return Ok(None);
2092            }
2093
2094            let case: Option<String> = row.get(2);
2095            let step: i64 = row.get(3);
2096            let phase: String = row.get(4);
2097            // The claim parks the subscription, in the same transaction: a
2098            // parked wait is matched no second event, and a crash before the
2099            // resume leaves a pair the redelivery pass finds and finishes.
2100            tx.execute(
2101                "UPDATE subscriptions SET parked_at = $4
2102                  WHERE tenant = $1 AND run_id = $2 AND effect_key = $3",
2103                &[&self.tenant_name(), &run, &effect, &at.unix_timestamp()],
2104            )
2105            .await
2106            .map_err(|e| be(&e))?;
2107            tx.commit().await.map_err(|e| be(&e))?;
2108
2109            return Ok(Some(Subscription {
2110                run: RunId::parse(&run).map_err(|e| corrupt("bad run id", e))?,
2111                case: case
2112                    .map(|c| CaseId::parse(&c))
2113                    .transpose()
2114                    .map_err(|e| corrupt("bad case id", e))?,
2115                effect: EffectKey::from_hex(&effect).map_err(|e| corrupt("bad effect key", e))?,
2116                step: crate::core::StepId(u32::try_from(step).unwrap_or(0)),
2117                phase: phase_from(&phase)?,
2118                kind: event.kind.clone(),
2119                correlation: event.correlation.clone(),
2120                from: row.get(5),
2121            }));
2122        }
2123        tx.commit().await.map_err(|e| be(&e))?;
2124        Ok(None)
2125    }
2126
2127    async fn deliver_to(
2128        &self,
2129        target: RunId,
2130        event: &InboundEvent,
2131        at: Timestamp,
2132    ) -> Result<TargetedDelivery, StoreError> {
2133        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2134        let tx = client.transaction().await.map_err(|e| be(&e))?;
2135        let tenant = self.tenant_name();
2136        let event_id = event.dedup_key();
2137
2138        // A message this run already holds claimed resumes the wait it was
2139        // claimed for, while that claim stands — a retry after a crash between
2140        // claim and resume. Consumed, or claimed by another run, it is a
2141        // duplicate. A new message goes to a wait holding no undelivered one.
2142        let existing_claim: Option<(Option<String>, Option<String>)> = tx
2143            .query_opt(
2144                "SELECT claimed_by, claimed_for FROM inbound_events
2145                  WHERE tenant = $1 AND event_id = $2 FOR UPDATE",
2146                &[&tenant, &event_id],
2147            )
2148            .await
2149            .map_err(|e| be(&e))?
2150            .map(|row| (row.get(0), row.get(1)));
2151        let target_text = target.to_string();
2152        let claimed_for: Option<String> = match &existing_claim {
2153            Some((Some(by), Some(effect))) if *by == target_text => Some(effect.clone()),
2154            _ => None,
2155        };
2156
2157        // Lock this run's candidate subscription rows. Two continuations for
2158        // one task then serialize before either can insert its event.
2159        //
2160        // Ordered by effect key, because that is the order the embedded
2161        // backend walks its subscription table in — its key is
2162        // `(tenant, run, effect, …)` — and when one run has several waits that
2163        // all match one event, *which* wait the delivery satisfies is
2164        // observable: the resumed step, the effect key on the wake record. Two
2165        // backends electing different waiters is a run that behaves
2166        // differently depending on which store it journals to. Lowest effect
2167        // key wins on both.
2168        let rows = tx
2169            .query(
2170                "SELECT effect_key, case_id, step, phase, namespace, value, from_source,
2171                        parked_at
2172                   FROM subscriptions
2173                  WHERE tenant = $1 AND run_id = $2 AND event_kind = $3
2174                    AND (from_source IS NULL OR from_source = $4)
2175                  ORDER BY effect_key ASC
2176                  FOR UPDATE",
2177                &[&tenant, &target.to_string(), &event.kind, &event.source],
2178            )
2179            .await
2180            .map_err(|e| be(&e))?;
2181        let Some(selected) = rows.iter().find(|row| {
2182            let effect: String = row.get(0);
2183            let namespace: String = row.get(4);
2184            let value: String = row.get(5);
2185            let parked: Option<i64> = row.get(7);
2186            let eligible = match (&existing_claim, &claimed_for) {
2187                (Some(_), Some(claimed)) => *claimed == effect,
2188                (Some(_), None) => false,
2189                (None, _) => parked.is_none(),
2190            };
2191            eligible
2192                && event
2193                    .correlation
2194                    .iter()
2195                    .any(|key| key.namespace == namespace && key.value == value)
2196        }) else {
2197            tx.commit().await.map_err(|e| be(&e))?;
2198            return Ok(if existing_claim.is_some() {
2199                TargetedDelivery::Duplicate
2200            } else {
2201                TargetedDelivery::NotWaiting
2202            });
2203        };
2204
2205        let effect: String = selected.get(0);
2206        let case: Option<String> = selected.get(1);
2207        let step: i64 = selected.get(2);
2208        let phase: String = selected.get(3);
2209        let subscription = Subscription {
2210            run: target,
2211            case: case
2212                .map(|value| CaseId::parse(&value))
2213                .transpose()
2214                .map_err(|e| corrupt("bad case id", e))?,
2215            effect: EffectKey::from_hex(&effect).map_err(|e| corrupt("bad effect key", e))?,
2216            step: crate::core::StepId(u32::try_from(step).unwrap_or(0)),
2217            phase: phase_from(&phase)?,
2218            kind: event.kind.clone(),
2219            correlation: rows
2220                .iter()
2221                .filter(|row| row.get::<_, String>(0) == effect)
2222                .map(|row| CorrelationKey::new(row.get::<_, String>(4), row.get::<_, String>(5)))
2223                .collect(),
2224            from: selected.get(6),
2225        };
2226
2227        if existing_claim.is_some() {
2228            // Selected only as the wait its standing claim was made for.
2229            tx.commit().await.map_err(|e| be(&e))?;
2230            return Ok(TargetedDelivery::Matched(subscription));
2231        }
2232
2233        let inserted = tx
2234            .execute(
2235                "INSERT INTO inbound_events
2236                   (event_id, source, bare_id, kind, payload, received_at,
2237                    claimed_by, claimed_at, tenant, by_actor, by_basis, claimed_for,
2238                    targeted)
2239                 VALUES ($1, $2, $3, $4, $5, $6, $7, $6, $8, $9, $10, $11, TRUE)
2240                 ON CONFLICT (tenant, event_id) DO NOTHING",
2241                &[
2242                    &event_id,
2243                    &event.source,
2244                    &event.id,
2245                    &event.kind,
2246                    &event.payload.to_string(),
2247                    &at.unix_timestamp(),
2248                    &target.to_string(),
2249                    &tenant,
2250                    &event.by.as_ref().map(|o| o.actor().to_owned()),
2251                    &event.by.as_ref().map(|o| o.basis().as_str().to_owned()),
2252                    &effect,
2253                ],
2254            )
2255            .await
2256            .map_err(|e| be(&e))?;
2257        if inserted == 0 {
2258            // Inserted by a concurrent delivery between the read above and
2259            // this write: the same rule, read again.
2260            let row = tx
2261                .query_one(
2262                    "SELECT claimed_by, claimed_for FROM inbound_events
2263                      WHERE tenant = $1 AND event_id = $2",
2264                    &[&tenant, &event_id],
2265                )
2266                .await
2267                .map_err(|e| be(&e))?;
2268            let (by, wait): (Option<String>, Option<String>) = (row.get(0), row.get(1));
2269            tx.commit().await.map_err(|e| be(&e))?;
2270            return Ok(
2271                if by.as_deref() == Some(target_text.as_str()) && wait.as_deref() == Some(&effect) {
2272                    TargetedDelivery::Matched(subscription)
2273                } else {
2274                    TargetedDelivery::Duplicate
2275                },
2276            );
2277        }
2278        for key in &event.correlation {
2279            tx.execute(
2280                "INSERT INTO inbound_correlation (event_id, namespace, value, tenant)
2281                 VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING",
2282                &[&event_id, &key.namespace, &key.value, &tenant],
2283            )
2284            .await
2285            .map_err(|e| be(&e))?;
2286        }
2287        // Parked until the run's unsubscribe: a crash between this claim and
2288        // the resume leaves a pair the redelivery pass finds.
2289        tx.execute(
2290            "UPDATE subscriptions SET parked_at = $4
2291              WHERE tenant = $1 AND run_id = $2 AND effect_key = $3",
2292            &[&tenant, &target.to_string(), &effect, &at.unix_timestamp()],
2293        )
2294        .await
2295        .map_err(|e| be(&e))?;
2296
2297        tx.commit().await.map_err(|e| be(&e))?;
2298        Ok(TargetedDelivery::Matched(subscription))
2299    }
2300
2301    async fn unsubscribe(&self, run: RunId, effect: EffectKey) -> Result<(), StoreError> {
2302        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2303        client
2304            .execute(
2305                "DELETE FROM subscriptions
2306                  WHERE tenant = $3 AND run_id = $1 AND effect_key = $2",
2307                &[&run.to_string(), &effect.to_hex(), &self.tenant_name()],
2308            )
2309            .await
2310            .map_err(|e| be(&e))?;
2311        // A wait's unsubscribe is the store's signal that its delivery was
2312        // journaled, so the buffer's copy of what this wait consumed is shed
2313        // here — only this wait's: another wait of the run may hold a message
2314        // not yet journaled. The row keeps its `(source, id)` identity, claim
2315        // and dead-letter fields: dedup and accounting need those, and only
2316        // the content was ever the erasure concern. Stripping at the *claim*
2317        // instead would lose the payload for a run that crashed between claim
2318        // and resume, whose recovery re-reads it from the buffer.
2319        client
2320            .execute(
2321                "UPDATE inbound_events SET payload = 'null', claimed_for = NULL
2322                  WHERE tenant = $1 AND claimed_by = $2 AND claimed_for = $3",
2323                &[&self.tenant_name(), &run.to_string(), &effect.to_hex()],
2324            )
2325            .await
2326            .map_err(|e| be(&e))?;
2327        Ok(())
2328    }
2329
2330    async fn unsubscribe_run(
2331        &self,
2332        run: RunId,
2333        unanswered: &[EffectKey],
2334    ) -> Result<crate::case::Retired, StoreError> {
2335        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2336        let tx = client.transaction().await.map_err(|e| be(&e))?;
2337        let retired = tx
2338            .query(
2339                "DELETE FROM subscriptions WHERE tenant = $1 AND run_id = $2
2340                 RETURNING effect_key",
2341                &[&self.tenant_name(), &run.to_string()],
2342            )
2343            .await
2344            .map_err(|e| be(&e))?;
2345        // A message claimed for a wait the run never answered reached nobody:
2346        // back to the buffer, unclaimed, with its payload.
2347        let unanswered: Vec<String> = unanswered.iter().map(|effect| effect.to_hex()).collect();
2348        let released = tx
2349            .query(
2350                "UPDATE inbound_events
2351                    SET claimed_by = NULL, claimed_at = NULL, claimed_for = NULL
2352                  WHERE tenant = $1 AND claimed_by = $2 AND claimed_for = ANY($3)
2353                    AND NOT dead AND NOT targeted
2354              RETURNING bare_id, kind, payload, received_at, source, event_id,
2355                        by_actor, by_basis",
2356                &[&self.tenant_name(), &run.to_string(), &unanswered],
2357            )
2358            .await
2359            .map_err(|e| be(&e))?;
2360        // One addressed to this run is its alone: dead-lettered, never offered
2361        // to another run waiting on the same key.
2362        tx.execute(
2363            "UPDATE inbound_events
2364                SET claimed_by = NULL, claimed_at = NULL, claimed_for = NULL,
2365                    dead = TRUE, dead_reason = $4
2366              WHERE tenant = $1 AND claimed_by = $2 AND claimed_for = ANY($3)
2367                AND NOT dead AND targeted",
2368            &[
2369                &self.tenant_name(),
2370                &run.to_string(),
2371                &unanswered,
2372                &crate::case::ADDRESSEE_CONCLUDED_REASON,
2373            ],
2374        )
2375        .await
2376        .map_err(|e| be(&e))?;
2377        // The same shedding `unsubscribe` does, for every other wait of the run.
2378        tx.execute(
2379            "UPDATE inbound_events SET payload = 'null', claimed_for = NULL
2380              WHERE tenant = $1 AND claimed_by = $2",
2381            &[&self.tenant_name(), &run.to_string()],
2382        )
2383        .await
2384        .map_err(|e| be(&e))?;
2385        tx.commit().await.map_err(|e| be(&e))?;
2386        let waits: std::collections::BTreeSet<String> =
2387            retired.iter().map(|row| row.get(0)).collect();
2388        let mut events = Vec::with_capacity(released.len());
2389        for row in &released {
2390            events.push(
2391                buffered_from(row, &client, &self.tenant_name())
2392                    .await?
2393                    .event,
2394            );
2395        }
2396        Ok(crate::case::Retired {
2397            waits: waits.len(),
2398            released: events,
2399        })
2400    }
2401
2402    async fn park_wait(&self, sub: &Subscription, at: Timestamp) -> Result<(), StoreError> {
2403        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2404        let tx = client.transaction().await.map_err(|e| be(&e))?;
2405        for k in &sub.correlation {
2406            tx.execute(
2407                "INSERT INTO subscriptions
2408                       (run_id, effect_key, case_id, step, phase, event_kind,
2409                        namespace, value, created_at, tenant, parked_at, from_source)
2410                     VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $9, $11)
2411                     ON CONFLICT (tenant, run_id, effect_key, namespace, value)
2412                     DO UPDATE SET parked_at = EXCLUDED.parked_at",
2413                &[
2414                    &sub.run.to_string(),
2415                    &sub.effect.to_hex(),
2416                    &sub.case.map(|c| c.to_string()),
2417                    &i64::from(sub.step.0),
2418                    &crate::core::Phase::as_str(sub.phase),
2419                    &sub.kind,
2420                    &k.namespace,
2421                    &k.value,
2422                    &at.unix_timestamp(),
2423                    &self.tenant_name(),
2424                    &sub.from,
2425                ],
2426            )
2427            .await
2428            .map_err(|e| be(&e))?;
2429        }
2430        tx.commit().await.map_err(|e| be(&e))?;
2431        Ok(())
2432    }
2433
2434    async fn parked_waits(&self, limit: usize) -> Result<Vec<Subscription>, StoreError> {
2435        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2436        // A sealed run can record no delivery. Its parked pairs are retired
2437        // here rather than listed, or the redelivery pass tries each, fails,
2438        // and finds it again every tick — and the payloads it held claimed
2439        // are shed as a conclusion's retirement sheds them.
2440        let tx = client.transaction().await.map_err(|e| be(&e))?;
2441        let retired = tx
2442            .query(
2443                "DELETE FROM subscriptions s
2444                  WHERE s.tenant = $1 AND s.parked_at IS NOT NULL
2445                    AND EXISTS (SELECT 1 FROM run_seal r
2446                                 WHERE r.tenant = s.tenant AND r.run_id = s.run_id)
2447              RETURNING s.run_id",
2448                &[&self.tenant_name()],
2449            )
2450            .await
2451            .map_err(|e| be(&e))?;
2452        let sealed: std::collections::BTreeSet<String> =
2453            retired.iter().map(|row| row.get(0)).collect();
2454        for run in &sealed {
2455            tx.execute(
2456                "UPDATE inbound_events SET payload = 'null', claimed_for = NULL
2457                  WHERE tenant = $1 AND claimed_by = $2",
2458                &[&self.tenant_name(), run],
2459            )
2460            .await
2461            .map_err(|e| be(&e))?;
2462        }
2463        tx.commit().await.map_err(|e| be(&e))?;
2464        // One row per wait, its keys gathered: a parked wait with two keys is
2465        // one pair to redeliver, not two.
2466        let rows = client
2467            .query(
2468                "SELECT run_id, effect_key, MIN(case_id), MIN(step), MIN(phase),
2469                        MIN(event_kind), array_agg(namespace), array_agg(value),
2470                        MIN(from_source)
2471                   FROM subscriptions
2472                  WHERE tenant = $2 AND parked_at IS NOT NULL
2473                  GROUP BY run_id, effect_key
2474                  ORDER BY MIN(parked_at) ASC, run_id, effect_key
2475                  LIMIT $1",
2476                &[
2477                    &i64::try_from(limit).unwrap_or(i64::MAX),
2478                    &self.tenant_name(),
2479                ],
2480            )
2481            .await
2482            .map_err(|e| be(&e))?;
2483        let mut out = Vec::with_capacity(rows.len());
2484        for row in rows {
2485            let run: String = row.get(0);
2486            let effect: String = row.get(1);
2487            let case: Option<String> = row.get(2);
2488            let step: i64 = row.get(3);
2489            let phase: String = row.get(4);
2490            let namespaces: Vec<String> = row.get(6);
2491            let values: Vec<String> = row.get(7);
2492            out.push(Subscription {
2493                run: RunId::parse(&run).map_err(|e| corrupt("bad run id", e))?,
2494                case: case
2495                    .map(|c| CaseId::parse(&c))
2496                    .transpose()
2497                    .map_err(|e| corrupt("bad case id", e))?,
2498                effect: EffectKey::from_hex(&effect).map_err(|e| corrupt("bad effect key", e))?,
2499                step: crate::core::StepId(u32::try_from(step).unwrap_or(0)),
2500                phase: phase_from(&phase)?,
2501                kind: row.get(5),
2502                correlation: namespaces
2503                    .into_iter()
2504                    .zip(values)
2505                    .map(|(ns, value)| CorrelationKey::new(ns, value))
2506                    .collect(),
2507                from: row.get(8),
2508            });
2509        }
2510        Ok(out)
2511    }
2512
2513    async fn erase_payload(&self, source: &str, id: &str) -> Result<bool, StoreError> {
2514        let mut client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2515        let key = crate::core::origin_key(source, id);
2516        let tenant = self.tenant_name();
2517        let tx = client.transaction().await.map_err(|e| be(&e))?;
2518        let Some(row) = tx
2519            .query_opt(
2520                "SELECT claimed_by, claimed_for IS NOT NULL FROM inbound_events
2521                  WHERE tenant = $1 AND event_id = $2 FOR UPDATE",
2522                &[&tenant, &key],
2523            )
2524            .await
2525            .map_err(|e| be(&e))?
2526        else {
2527            tx.commit().await.map_err(|e| be(&e))?;
2528            return Ok(false);
2529        };
2530        let claimed_by: Option<String> = row.get(0);
2531        let standing: bool = row.get(1);
2532        // Undelivered: nobody claimed it, or a wait claimed it and has not
2533        // yet journaled it — the claim stands until that wait's unsubscribe
2534        // consumes it. Either way it leaves the claimable set as a dead
2535        // letter, in this write, so no recovery hands the run the emptied row;
2536        // the claim is released and the claimant's wait is unparked, so the
2537        // wait stays open for its deadline to bound. A delivered row keeps its
2538        // claim and its accounting. Every right-hand side reads the row as it
2539        // was before the update.
2540        let undelivered = claimed_by.is_none() || standing;
2541        tx.execute(
2542            "UPDATE inbound_events SET payload = 'null', erased = TRUE,
2543                    dead = dead OR $3,
2544                    dead_reason = CASE WHEN $3 AND NOT dead THEN $4 ELSE dead_reason END,
2545                    claimed_by = CASE WHEN $3 THEN NULL ELSE claimed_by END,
2546                    claimed_at = CASE WHEN $3 THEN NULL ELSE claimed_at END,
2547                    claimed_for = NULL
2548              WHERE tenant = $1 AND event_id = $2",
2549            &[&tenant, &key, &undelivered, &crate::case::ERASED_REASON],
2550        )
2551        .await
2552        .map_err(|e| be(&e))?;
2553        if let (true, Some(run)) = (undelivered, claimed_by) {
2554            tx.execute(
2555                "UPDATE subscriptions s SET parked_at = NULL
2556                  WHERE s.tenant = $1 AND s.run_id = $2 AND s.parked_at IS NOT NULL
2557                    AND EXISTS (SELECT 1 FROM inbound_correlation c
2558                                 WHERE c.tenant = s.tenant AND c.event_id = $3
2559                                   AND c.namespace = s.namespace AND c.value = s.value)",
2560                &[&tenant, &run, &key],
2561            )
2562            .await
2563            .map_err(|e| be(&e))?;
2564        }
2565        tx.commit().await.map_err(|e| be(&e))?;
2566        Ok(true)
2567    }
2568
2569    async fn minter(
2570        &self,
2571        source: &str,
2572        id: &str,
2573    ) -> Result<Option<crate::case::Minter>, StoreError> {
2574        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2575        let row = client
2576            .query_opt(
2577                "SELECT by_actor, by_basis FROM inbound_events
2578                  WHERE tenant = $1 AND event_id = $2",
2579                &[&self.tenant_name(), &crate::core::origin_key(source, id)],
2580            )
2581            .await
2582            .map_err(|e| be(&e))?;
2583        let Some(row) = row else {
2584            return Ok(None);
2585        };
2586        match (
2587            row.get::<_, Option<String>>(0),
2588            row.get::<_, Option<String>>(1),
2589        ) {
2590            (Some(actor), Some(basis)) => Ok(Some(crate::case::Minter::Operator(
2591                super::decode_operator(&actor, &basis, "inbound_events")?,
2592            ))),
2593            (None, None) => Ok(Some(crate::case::Minter::Nobody)),
2594            _ => Err(StoreError::Corrupt {
2595                seq: 0,
2596                detail: "inbound_events holds half an operator: a minted event carries a \
2597                         name and what established it, or neither"
2598                    .to_owned(),
2599            }),
2600        }
2601    }
2602
2603    async fn sweep_unclaimed(
2604        &self,
2605        older_than: Timestamp,
2606        reason: &str,
2607    ) -> Result<usize, StoreError> {
2608        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2609        // `<=`, not `<`: a zero grace window must retire everything already
2610        // buffered, and with second-granularity stamps `<` silently spares
2611        // anything received this second — the same boundary the redb sweep
2612        // states, pinned for both by the conformance battery.
2613        let n = client
2614            .execute(
2615                "UPDATE inbound_events SET dead = TRUE, dead_reason = $2
2616                  WHERE tenant = $3 AND claimed_by IS NULL AND NOT dead
2617                    AND received_at <= $1",
2618                &[
2619                    &older_than.unix_timestamp(),
2620                    &reason.to_owned(),
2621                    &self.tenant_name(),
2622                ],
2623            )
2624            .await
2625            .map_err(|e| be(&e))?;
2626        Ok(usize::try_from(n).unwrap_or(0))
2627    }
2628
2629    async fn dead_letters(&self, limit: usize) -> Result<Vec<DeadLetter>, StoreError> {
2630        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2631        let rows = client
2632            .query(
2633                "SELECT bare_id, kind, payload, received_at, dead_reason, source, event_id,
2634                        by_actor, by_basis
2635                   FROM inbound_events WHERE tenant = $2 AND dead
2636                  ORDER BY received_at DESC LIMIT $1",
2637                &[
2638                    &i64::try_from(limit).unwrap_or(i64::MAX),
2639                    &self.tenant_name(),
2640                ],
2641            )
2642            .await
2643            .map_err(|e| be(&e))?;
2644        let mut out = Vec::with_capacity(rows.len());
2645        for row in rows {
2646            let payload: String = row.get(2);
2647            // The correlation keys are read back, as the redb backend reads
2648            // them: the first question about an unclaimed message is *what was
2649            // it correlated by*, and a reconstructed event with no keys is a
2650            // valid-looking message silently stripped of the field that
2651            // explains it.
2652            let event_id: String = row.get(6);
2653            let corr = client
2654                .query(
2655                    "SELECT namespace, value FROM inbound_correlation
2656                      WHERE tenant = $2 AND event_id = $1",
2657                    &[&event_id, &self.tenant_name()],
2658                )
2659                .await
2660                .map_err(|e| be(&e))?;
2661            out.push(DeadLetter {
2662                event: InboundEvent {
2663                    source: row.get(5),
2664                    id: row.get(0),
2665                    kind: row.get(1),
2666                    by: match (
2667                        row.get::<_, Option<String>>(7),
2668                        row.get::<_, Option<String>>(8),
2669                    ) {
2670                        (Some(actor), Some(basis)) => {
2671                            Some(super::decode_operator(&actor, &basis, "inbound_events")?)
2672                        }
2673                        _ => None,
2674                    },
2675                    correlation: corr
2676                        .iter()
2677                        .map(|r| CorrelationKey::new(r.get::<_, String>(0), r.get::<_, String>(1)))
2678                        .collect(),
2679                    payload: serde_json::from_str(&payload)?,
2680                },
2681                received_at: Timestamp::from_unix_timestamp(row.get::<_, i64>(3))
2682                    .map_err(|e| corrupt("unrepresentable received_at", e))?,
2683                // One wording across backends: a retirement that recorded no
2684                // reason reads "unclaimed", exactly as the redb backend
2685                // answers, rather than an empty string here and a word there —
2686                // an operator scripting on the reason field must not have to
2687                // know which store the plane journals to.
2688                reason: row
2689                    .get::<_, Option<String>>(4)
2690                    .filter(|reason| !reason.is_empty())
2691                    .unwrap_or_else(|| "unclaimed".to_owned()),
2692            });
2693        }
2694        Ok(out)
2695    }
2696
2697    async fn waiting(&self, limit: usize) -> Result<Vec<Subscription>, StoreError> {
2698        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2699        let rows = client
2700            .query(
2701                "SELECT run_id, effect_key, case_id, step, phase, event_kind, namespace, value,
2702                        from_source
2703                   FROM subscriptions WHERE tenant = $2
2704                  ORDER BY created_at ASC LIMIT $1",
2705                &[
2706                    &i64::try_from(limit).unwrap_or(i64::MAX),
2707                    &self.tenant_name(),
2708                ],
2709            )
2710            .await
2711            .map_err(|e| be(&e))?;
2712        let mut out = Vec::with_capacity(rows.len());
2713        for row in rows {
2714            let run: String = row.get(0);
2715            let effect: String = row.get(1);
2716            let case: Option<String> = row.get(2);
2717            let step: i64 = row.get(3);
2718            let phase: String = row.get(4);
2719            out.push(Subscription {
2720                run: RunId::parse(&run).map_err(|e| corrupt("bad run id", e))?,
2721                case: case
2722                    .map(|c| CaseId::parse(&c))
2723                    .transpose()
2724                    .map_err(|e| corrupt("bad case id", e))?,
2725                effect: EffectKey::from_hex(&effect).map_err(|e| corrupt("bad effect key", e))?,
2726                step: crate::core::StepId(u32::try_from(step).unwrap_or(0)),
2727                phase: phase_from(&phase)?,
2728                kind: row.get(5),
2729                correlation: vec![CorrelationKey::new(
2730                    row.get::<_, String>(6),
2731                    row.get::<_, String>(7),
2732                )],
2733                from: row.get(8),
2734            });
2735        }
2736        Ok(out)
2737    }
2738}
2739
2740async fn buffered_from(
2741    row: &tokio_postgres::Row,
2742    client: &deadpool_postgres::Client,
2743    tenant: &str,
2744) -> Result<BufferedEvent, StoreError> {
2745    let id: String = row.get(0);
2746    let payload: String = row.get(2);
2747    // Correlation is keyed by the dedup key, which is what `buffer` wrote — not
2748    // by `bare_id`, the producer's own id. Reading it back under the wrong one
2749    // returns nothing, and a claimed event would arrive with no keys at all:
2750    // valid-looking, silently stripped of what it was routed on.
2751    let event_id: String = row.get(5);
2752    let corr = client
2753        .query(
2754            "SELECT namespace, value FROM inbound_correlation
2755              WHERE tenant = $2 AND event_id = $1",
2756            &[&event_id, &tenant],
2757        )
2758        .await
2759        .map_err(|e| be(&e))?;
2760    Ok(BufferedEvent {
2761        event: InboundEvent {
2762            source: row.get(4),
2763            id: id.clone(),
2764            kind: row.get(1),
2765            correlation: corr
2766                .iter()
2767                .map(|r| CorrelationKey::new(r.get::<_, String>(0), r.get::<_, String>(1)))
2768                .collect(),
2769            payload: serde_json::from_str(&payload)?,
2770            // Both columns or neither: a name with no basis is an
2771            // attribution this build cannot state, and inventing one would
2772            // claim the strong form for a row that never held it.
2773            by: match (
2774                row.get::<_, Option<String>>(6),
2775                row.get::<_, Option<String>>(7),
2776            ) {
2777                (Some(actor), Some(basis)) => {
2778                    Some(super::decode_operator(&actor, &basis, "inbound_events")?)
2779                }
2780                (None, None) => None,
2781                _ => {
2782                    return Err(StoreError::Corrupt {
2783                        seq: 0,
2784                        detail: "inbound_events holds half an operator: a minted event \
2785                                 carries a name and what established it, or neither"
2786                            .to_owned(),
2787                    });
2788                }
2789            },
2790        },
2791        received_at: Timestamp::from_unix_timestamp(row.get::<_, i64>(3))
2792            .map_err(|e| corrupt("unrepresentable received_at", e))?,
2793    })
2794}
2795
2796/// How long a timer claim holds before another sweep may take it.
2797const CLAIM_LEASE: i64 = 60;
2798
2799#[async_trait]
2800impl TimerStore for PostgresStore {
2801    fn tenant(&self) -> &str {
2802        crate::journal::JournalStore::tenant(self)
2803    }
2804
2805    async fn arm(&self, timer: &Timer) -> Result<(), StoreError> {
2806        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2807        client
2808            .execute(
2809                "INSERT INTO timers
2810                   (run_id, effect_key, case_id, step, phase, fire_at, tenant)
2811                 VALUES ($1, $2, $3, $4, $5, $6, $7)
2812                 ON CONFLICT (tenant, run_id, effect_key) DO NOTHING",
2813                &[
2814                    &timer.run.to_string(),
2815                    &timer.effect.to_hex(),
2816                    &timer.case.map(|c| c.to_string()),
2817                    &i64::from(timer.step.0),
2818                    &crate::core::Phase::as_str(timer.phase),
2819                    &timer.fire_at.unix_timestamp(),
2820                    &self.tenant_name(),
2821                ],
2822            )
2823            .await
2824            .map_err(|e| be(&e))?;
2825        Ok(())
2826    }
2827
2828    async fn claim_due(&self, now: Timestamp, limit: usize) -> Result<Vec<Timer>, StoreError> {
2829        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2830        // A sealed run can record no wake. Its due timers are retired rather
2831        // than claimed, or each is claimed and fails once per lease period for
2832        // ever.
2833        client
2834            .execute(
2835                "DELETE FROM timers t
2836                  WHERE t.tenant = $2 AND t.fire_at <= $1
2837                    AND EXISTS (SELECT 1 FROM run_seal s
2838                                 WHERE s.tenant = t.tenant AND s.run_id = t.run_id)",
2839                &[&now.unix_timestamp(), &self.tenant_name()],
2840            )
2841            .await
2842            .map_err(|e| be(&e))?;
2843        // Claimed and selected in one statement, with `SKIP LOCKED` so a second
2844        // sweeper takes different rows rather than blocking on the first's.
2845        let rows = client
2846            .query(
2847                "UPDATE timers SET claimed_at = $1
2848                  WHERE tenant = $4 AND (run_id, effect_key) IN (
2849                      SELECT run_id, effect_key FROM timers
2850                       WHERE tenant = $4 AND fire_at <= $1
2851                         AND (claimed_at IS NULL OR claimed_at <= $2)
2852                       ORDER BY fire_at ASC
2853                       FOR UPDATE SKIP LOCKED
2854                       LIMIT $3)
2855              RETURNING run_id, effect_key, case_id, step, phase, fire_at",
2856                &[
2857                    &now.unix_timestamp(),
2858                    &(now.unix_timestamp() - CLAIM_LEASE),
2859                    &i64::try_from(limit).unwrap_or(i64::MAX),
2860                    &self.tenant_name(),
2861                ],
2862            )
2863            .await
2864            .map_err(|e| be(&e))?;
2865        rows.iter().map(timer_from).collect()
2866    }
2867
2868    async fn pending_count(&self) -> Result<u64, StoreError> {
2869        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2870        let n: i64 = client
2871            .query_one(
2872                "SELECT COUNT(*) FROM timers WHERE tenant = $1",
2873                &[&self.tenant_name()],
2874            )
2875            .await
2876            .map_err(|e| be(&e))?
2877            .get(0);
2878        Ok(u64::try_from(n).unwrap_or(0))
2879    }
2880
2881    async fn disarm(&self, run: RunId, effect: EffectKey) -> Result<(), StoreError> {
2882        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2883        client
2884            .execute(
2885                "DELETE FROM timers WHERE tenant = $3 AND run_id = $1 AND effect_key = $2",
2886                &[&run.to_string(), &effect.to_hex(), &self.tenant_name()],
2887            )
2888            .await
2889            .map_err(|e| be(&e))?;
2890        Ok(())
2891    }
2892
2893    async fn disarm_run(&self, run: RunId) -> Result<usize, StoreError> {
2894        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2895        let n = client
2896            .execute(
2897                "DELETE FROM timers WHERE tenant = $2 AND run_id = $1",
2898                &[&run.to_string(), &self.tenant_name()],
2899            )
2900            .await
2901            .map_err(|e| be(&e))?;
2902        Ok(usize::try_from(n).unwrap_or(usize::MAX))
2903    }
2904}
2905
2906fn timer_from(row: &tokio_postgres::Row) -> Result<Timer, StoreError> {
2907    let run: String = row.get(0);
2908    let effect: String = row.get(1);
2909    let case: Option<String> = row.get(2);
2910    let step: i64 = row.get(3);
2911    let phase: String = row.get(4);
2912    Ok(Timer {
2913        run: RunId::parse(&run).map_err(|e| corrupt("bad run id", e))?,
2914        case: case
2915            .map(|c| CaseId::parse(&c))
2916            .transpose()
2917            .map_err(|e| corrupt("bad case id", e))?,
2918        effect: EffectKey::from_hex(&effect).map_err(|e| corrupt("bad effect key", e))?,
2919        step: crate::core::StepId(u32::try_from(step).unwrap_or(0)),
2920        phase: phase_from(&phase)?,
2921        fire_at: Timestamp::from_unix_timestamp(row.get::<_, i64>(5))
2922            .map_err(|e| corrupt("unrepresentable fire_at", e))?,
2923    })
2924}
2925
2926/// One task's row, read on a connection the caller already holds.
2927///
2928/// This exists because of a deadlock, not for tidiness. `claim` held its
2929/// pooled connection across a `Self::task` call, and `task` acquired a second
2930/// connection from the same pool — so every in-flight claim needed two
2931/// connections while holding one. Sixteen reviewers racing one task on a small
2932/// pool each held a connection and waited for another that only a waiter could
2933/// release: a deadlock that reproduces exactly under the concurrency the claim
2934/// verb exists to survive, and only where the pool is small enough to exhaust
2935/// — a large development machine passes over the defect a CI runner hangs on,
2936/// which is how it shipped. Every read made while a connection is held goes
2937/// through this instead of re-entering the pool.
2938async fn task_on(
2939    client: &tokio_postgres::Client,
2940    tenant: &str,
2941    id: TaskId,
2942) -> Result<Option<Task>, StoreError> {
2943    let row = client
2944        .query_opt(
2945            &format!("SELECT {TASK_COLS} FROM tasks WHERE task_id = $1 AND tenant = $2"),
2946            &[&id.to_hex(), &tenant],
2947        )
2948        .await
2949        .map_err(|e| be(&e))?;
2950    row.as_ref().map(task_from).transpose()
2951}
2952
2953#[async_trait]
2954impl TaskStore for PostgresStore {
2955    fn tenant(&self) -> &str {
2956        self.tenant_str()
2957    }
2958
2959    async fn open(&self, task: &Task) -> Result<Task, StoreError> {
2960        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2961        client
2962            .execute(
2963                "INSERT INTO tasks (task_id, run_id, case_id, kind, justification,
2964                                    candidate_roles, escalate_to, excluded_actors, assignee,
2965                                    priority, state, on_expiry, priority_rank, created_at,
2966                                    due_at, tenant, withheld)
2967                 VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17)
2968                 ON CONFLICT (tenant, task_id) DO NOTHING",
2969                &[
2970                    &task.id.to_hex(),
2971                    &task.run.to_string(),
2972                    &task.case.map(|c| c.to_string()),
2973                    &task.kind,
2974                    &serde_json::to_string(&task.justification)?,
2975                    &task.candidate_roles,
2976                    &task.escalate_to,
2977                    &task.excluded_actors,
2978                    &task.assignee,
2979                    &task.priority.as_str(),
2980                    &task.state.as_str(),
2981                    &OnExpiry::as_str(task.on_expiry),
2982                    &i16::from(task.priority.rank()),
2983                    &task.created_at.unix_timestamp(),
2984                    &task.due_at.map(Timestamp::unix_timestamp),
2985                    &self.tenant_name(),
2986                    &task.withheld.map(crate::core::Withheld::as_str),
2987                ],
2988            )
2989            .await
2990            .map_err(|e| be(&e))?;
2991        Ok(task_on(&client, &self.tenant_name(), task.id)
2992            .await?
2993            .unwrap_or_else(|| task.clone()))
2994    }
2995
2996    async fn task(&self, id: TaskId) -> Result<Option<Task>, StoreError> {
2997        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
2998        task_on(&client, &self.tenant_name(), id).await
2999    }
3000
3001    async fn claim(&self, id: TaskId, actor: &str, roles: &[String]) -> Result<Task, ClaimError> {
3002        // One connection for the whole verb — every read below goes through
3003        // `task_on` rather than `Self::task`, which would take a second
3004        // connection while this one is held. See `task_on` for the deadlock
3005        // that shape produced under exactly the concurrency this verb exists
3006        // to survive.
3007        let client = self
3008            .pool()
3009            .get()
3010            .await
3011            .map_err(|e| ClaimError::Store(pool_err(&e)))?;
3012        let tenant = self.tenant_name();
3013        let Some(task) = task_on(&client, &tenant, id)
3014            .await
3015            .map_err(ClaimError::Store)?
3016        else {
3017            return Err(ClaimError::NotFound(id));
3018        };
3019        // Eligibility before availability, and the order is load-bearing — see
3020        // `TaskStore::claim`.
3021        //
3022        // Four eyes: whoever proposed the action does not approve it.
3023        if task.excluded_actors.iter().any(|a| a == actor) {
3024            return Err(ClaimError::Excluded {
3025                actor: actor.to_owned(),
3026            });
3027        }
3028        if !task.candidate_roles.is_empty()
3029            && !task.candidate_roles.iter().any(|r| roles.contains(r))
3030        {
3031            return Err(ClaimError::WrongRole {
3032                actor: actor.to_owned(),
3033            });
3034        }
3035        if !task.state.is_pending() {
3036            return Err(ClaimError::NotPending {
3037                task: id,
3038                state: task.state,
3039            });
3040        }
3041
3042        // The reservation itself is one statement, guarded on the row still
3043        // being unheld **and still pending**. Checking above and writing here
3044        // would leave a window two reviewers both pass through — and the state
3045        // predicate is not redundant with the eligibility read: a sweep can
3046        // expire the task between that read and this write, and an UPDATE
3047        // keyed on the assignee alone would resurrect an expired task into
3048        // `claimed`, un-deciding the expiry policy that already fired on it.
3049        // `take_over` and `release` carry the same predicate for the same
3050        // reason.
3051        //
3052        // The second arm is the **same-holder re-claim**, which is idempotent
3053        // success — the contract on `TaskStore::claim`, and the one the redb
3054        // backend already honoured while this one refused with
3055        // `AlreadyClaimed { holder: yourself }`. A claim whose acknowledgement
3056        // was lost is retried by an honest client, and the retry must converge
3057        // on "you hold it" rather than bounce off its own success. It does not
3058        // widen the race: a task claimed by anybody *else* still fails every
3059        // arm, and an expired task matches neither state predicate.
3060        let updated = client
3061            .execute(
3062                "UPDATE tasks SET assignee = $2, state = $5
3063                  WHERE task_id = $1 AND tenant = $3
3064                    AND ((state = ANY($4::text[])
3065                          AND (assignee IS NULL OR assignee = $2))
3066                         OR (state = $5 AND assignee = $2))",
3067                &[
3068                    &id.to_hex(),
3069                    &actor.to_owned(),
3070                    &self.tenant_name(),
3071                    &task_states(TaskState::is_queued),
3072                    &TaskState::Claimed.as_str(),
3073                ],
3074            )
3075            .await
3076            .map_err(|e| ClaimError::Store(be(&e)))?;
3077        if updated == 0 {
3078            // Which guard refused: state, or holder. Re-read and say so — a
3079            // barred claim reported as "held by nobody" sends the reviewer
3080            // back to a queue that will refuse them again.
3081            let current = task_on(&client, &tenant, id)
3082                .await
3083                .map_err(ClaimError::Store)?;
3084            return Err(match current {
3085                None => ClaimError::NotFound(id),
3086                Some(t) if !t.state.is_pending() => ClaimError::NotPending {
3087                    task: id,
3088                    state: t.state,
3089                },
3090                Some(t) => ClaimError::AlreadyClaimed {
3091                    task: id,
3092                    holder: t.assignee.unwrap_or_default(),
3093                },
3094            });
3095        }
3096        task_on(&client, &tenant, id)
3097            .await
3098            .map_err(ClaimError::Store)?
3099            .ok_or(ClaimError::NotFound(id))
3100    }
3101
3102    async fn take_over(
3103        &self,
3104        id: TaskId,
3105        from: &str,
3106        actor: &str,
3107        roles: &[String],
3108    ) -> Result<Task, ClaimError> {
3109        // One connection for the whole verb, exactly as `claim` — see `task_on`.
3110        let client = self
3111            .pool()
3112            .get()
3113            .await
3114            .map_err(|e| ClaimError::Store(pool_err(&e)))?;
3115        let tenant = self.tenant_name();
3116        let Some(task) = task_on(&client, &tenant, id)
3117            .await
3118            .map_err(ClaimError::Store)?
3119        else {
3120            return Err(ClaimError::NotFound(id));
3121        };
3122        // Claim's eligibility-first order, unchanged: a take-over is a claim,
3123        // and four-eyes exclusion does not thin because the previous reviewer
3124        // left.
3125        if task.excluded_actors.iter().any(|a| a == actor) {
3126            return Err(ClaimError::Excluded {
3127                actor: actor.to_owned(),
3128            });
3129        }
3130        if !task.candidate_roles.is_empty()
3131            && !task.candidate_roles.iter().any(|r| roles.contains(r))
3132        {
3133            return Err(ClaimError::WrongRole {
3134                actor: actor.to_owned(),
3135            });
3136        }
3137        if !task.state.is_pending() {
3138            return Err(ClaimError::NotPending {
3139                task: id,
3140                state: task.state,
3141            });
3142        }
3143
3144        // The displacement is one statement, guarded on the holder still being
3145        // the one the caller named — the compare-and-swap that keeps a
3146        // take-over decided from a stale view from displacing whoever holds
3147        // the task *now*.
3148        let updated = client
3149            .execute(
3150                "UPDATE tasks SET assignee = $2, state = 'claimed'
3151                  WHERE task_id = $1 AND tenant = $4
3152                    AND assignee = $3 AND state = 'claimed'",
3153                &[
3154                    &id.to_hex(),
3155                    &actor.to_owned(),
3156                    &from.to_owned(),
3157                    &self.tenant_name(),
3158                ],
3159            )
3160            .await
3161            .map_err(|e| ClaimError::Store(be(&e)))?;
3162        if updated == 0 {
3163            return Err(ClaimError::NotHeld {
3164                task: id,
3165                actor: from.to_owned(),
3166            });
3167        }
3168        task_on(&client, &tenant, id)
3169            .await
3170            .map_err(ClaimError::Store)?
3171            .ok_or(ClaimError::NotFound(id))
3172    }
3173
3174    async fn release(&self, id: TaskId, actor: &str) -> Result<(), ClaimError> {
3175        let client = self
3176            .pool()
3177            .get()
3178            .await
3179            .map_err(|e| ClaimError::Store(pool_err(&e)))?;
3180        let freed = client
3181            .execute(
3182                "UPDATE tasks SET assignee = NULL, state = 'open'
3183                  WHERE task_id = $1 AND tenant = $3
3184                    AND assignee = $2 AND state = 'claimed'",
3185                &[&id.to_hex(), &actor.to_owned(), &self.tenant_name()],
3186            )
3187            .await
3188            .map_err(|e| ClaimError::Store(be(&e)))?;
3189        // The predicate did the work; the row count is what says whether it
3190        // matched. Discarding it reported success for a release that freed
3191        // nothing — caught by the conformance battery, not by this backend's
3192        // own tests.
3193        if freed == 0 {
3194            return Err(ClaimError::NotHeld {
3195                task: id,
3196                actor: actor.to_owned(),
3197            });
3198        }
3199        Ok(())
3200    }
3201
3202    async fn set_state(&self, id: TaskId, state: TaskState) -> Result<bool, StoreError> {
3203        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3204        // A compare-and-set from the pending states: the settlement that lost
3205        // the race to a decision or an expiry must not overwrite the winner's.
3206        let n = client
3207            .execute(
3208                "UPDATE tasks SET state = $2
3209                  WHERE task_id = $1 AND tenant = $3 AND state = ANY($4::text[])",
3210                &[
3211                    &id.to_hex(),
3212                    &state.as_str(),
3213                    &self.tenant_name(),
3214                    &task_states(TaskState::is_pending),
3215                ],
3216            )
3217            .await
3218            .map_err(|e| be(&e))?;
3219        if n > 0 {
3220            return Ok(true);
3221        }
3222        // Nothing moved: settled already, or no such task — and the caller
3223        // needs to know which, as the embedded backend tells it.
3224        match task_on(&client, &self.tenant_name(), id).await? {
3225            Some(_) => Ok(false),
3226            None => Err(StoreError::NotFound(id.to_string())),
3227        }
3228    }
3229
3230    async fn withdraw_run(&self, run: RunId, awaited: &[TaskId]) -> Result<usize, StoreError> {
3231        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3232        let awaited: Vec<String> = awaited.iter().map(|id| id.to_hex()).collect();
3233        let n = client
3234            .execute(
3235                "UPDATE tasks SET state = $3
3236                  WHERE tenant = $1 AND run_id = $2 AND state = ANY($4::text[])
3237                    AND task_id = ANY($5::text[])",
3238                &[
3239                    &self.tenant_name(),
3240                    &run.to_string(),
3241                    &TaskState::Withdrawn.as_str(),
3242                    &task_states(TaskState::is_pending),
3243                    &awaited,
3244                ],
3245            )
3246            .await
3247            .map_err(|e| be(&e))?;
3248        Ok(usize::try_from(n).unwrap_or(usize::MAX))
3249    }
3250
3251    async fn escalate(&self, id: TaskId) -> Result<Task, StoreError> {
3252        // One connection for the whole verb, exactly as `claim` — see `task_on`.
3253        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3254        let tenant = self.tenant_name();
3255        let Some(task) = task_on(&client, &tenant, id).await? else {
3256            return Err(StoreError::NotFound(id.to_string()));
3257        };
3258        // A decided task stays decided — the sweep that escalates races the
3259        // reviewer it is escalating past, and the decision winning that race
3260        // is the outcome everybody wanted.
3261        if !task.state.is_pending() {
3262            return Ok(task);
3263        }
3264        // `Task::escalate` is the one implementation of what escalating means;
3265        // this backend only persists its result. The state predicate is the
3266        // same compare-and-swap `claim` carries: a decision landing between
3267        // the read above and this write must win, and an unguarded UPDATE
3268        // would resurrect it into `escalated`.
3269        let mut updated = task;
3270        updated.escalate();
3271        let n = client
3272            .execute(
3273                "UPDATE tasks SET state = $5, assignee = NULL, candidate_roles = $2
3274                  WHERE task_id = $1 AND tenant = $3
3275                    AND state = ANY($4::text[])",
3276                &[
3277                    &id.to_hex(),
3278                    &updated.candidate_roles,
3279                    &tenant,
3280                    &task_states(TaskState::is_pending),
3281                    &TaskState::Escalated.as_str(),
3282                ],
3283            )
3284            .await
3285            .map_err(|e| be(&e))?;
3286        if n == 0 {
3287            // The race fired: somebody decided it. Return what stands.
3288            return task_on(&client, &tenant, id)
3289                .await?
3290                .ok_or_else(|| StoreError::NotFound(id.to_string()));
3291        }
3292        Ok(updated)
3293    }
3294
3295    async fn queue(&self, roles: &[String], limit: usize) -> Result<Vec<Task>, StoreError> {
3296        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3297        // Most urgent first, oldest within a rank — the trait's contract, which
3298        // the redb backend serves from a rank-keyed index. Ordering by age
3299        // alone made the limit a filter on the wrong axis: an urgent task
3300        // behind a page of older normal ones was simply absent from the page,
3301        // which for a worklist means the most important decision is the one
3302        // nobody is shown.
3303        let rows = client
3304            .query(
3305                &format!(
3306                    "SELECT {TASK_COLS} FROM tasks
3307                      WHERE tenant = $2 AND state = ANY($3::text[])
3308                      ORDER BY priority_rank ASC, created_at ASC
3309                      LIMIT $1"
3310                ),
3311                &[
3312                    &i64::try_from(limit).unwrap_or(i64::MAX),
3313                    &self.tenant_name(),
3314                    &task_states(TaskState::is_queued),
3315                ],
3316            )
3317            .await
3318            .map_err(|e| be(&e))?;
3319        let all: Result<Vec<Task>, StoreError> = rows.iter().map(task_from).collect();
3320        Ok(all?
3321            .into_iter()
3322            .filter(|t| {
3323                t.candidate_roles.is_empty() || t.candidate_roles.iter().any(|r| roles.contains(r))
3324            })
3325            .collect())
3326    }
3327
3328    async fn for_case(&self, case: CaseId) -> Result<Vec<Task>, StoreError> {
3329        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3330        let rows = client
3331            .query(
3332                &format!(
3333                    "SELECT {TASK_COLS} FROM tasks
3334                      WHERE case_id = $1 AND tenant = $2 ORDER BY created_at ASC"
3335                ),
3336                &[&case.to_string(), &self.tenant_name()],
3337            )
3338            .await
3339            .map_err(|e| be(&e))?;
3340        rows.iter().map(task_from).collect()
3341    }
3342
3343    async fn open_count(&self) -> Result<u64, StoreError> {
3344        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3345        let n: i64 = client
3346            .query_one(
3347                "SELECT COUNT(*) FROM tasks
3348                  WHERE tenant = $1 AND state = ANY($2::text[])",
3349                &[&self.tenant_name(), &task_states(TaskState::is_pending)],
3350            )
3351            .await
3352            .map_err(|e| be(&e))?
3353            .get(0);
3354        Ok(u64::try_from(n).unwrap_or(0))
3355    }
3356
3357    async fn overdue(&self, now: Timestamp, limit: usize) -> Result<Vec<Task>, StoreError> {
3358        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3359        // Not `escalated`, although escalated tasks are pending and past due:
3360        // this scan drives the expiry sweep, and escalation is that policy
3361        // having fired — see `TaskStore::overdue` for what including them
3362        // starves.
3363        let rows = client
3364            .query(
3365                &format!(
3366                    "SELECT {TASK_COLS} FROM tasks
3367                      WHERE tenant = $3 AND state = ANY($4::text[])
3368                        AND due_at IS NOT NULL AND due_at <= $1
3369                      ORDER BY due_at ASC LIMIT $2"
3370                ),
3371                &[
3372                    &now.unix_timestamp(),
3373                    &i64::try_from(limit).unwrap_or(i64::MAX),
3374                    &self.tenant_name(),
3375                    &task_states(TaskState::awaits_expiry),
3376                ],
3377            )
3378            .await
3379            .map_err(|e| be(&e))?;
3380        rows.iter().map(task_from).collect()
3381    }
3382}
3383
3384const TASK_COLS: &str = "task_id, run_id, case_id, kind, justification, candidate_roles, \
3385                         escalate_to, excluded_actors, assignee, priority, state, on_expiry, \
3386                         created_at, due_at, withheld";
3387
3388fn task_from(row: &tokio_postgres::Row) -> Result<Task, StoreError> {
3389    let id: String = row.get(0);
3390    let run: String = row.get(1);
3391    let case: Option<String> = row.get(2);
3392    let justification: String = row.get(4);
3393    let priority: String = row.get(9);
3394    let state: String = row.get(10);
3395    let on_expiry: String = row.get(11);
3396    let due: Option<i64> = row.get(13);
3397    let withheld: Option<String> = row.get(14);
3398
3399    Ok(Task {
3400        id: TaskId::parse(&id).map_err(|e| corrupt("bad task id", e))?,
3401        run: RunId::parse(&run).map_err(|e| corrupt("bad run id", e))?,
3402        case: case
3403            .map(|c| CaseId::parse(&c))
3404            .transpose()
3405            .map_err(|e| corrupt("bad case id", e))?,
3406        kind: row.get(3),
3407        justification: serde_json::from_str(&justification)?,
3408        candidate_roles: row.get(5),
3409        escalate_to: row.get(6),
3410        excluded_actors: row.get(7),
3411        assignee: row.get(8),
3412        priority: priority_from(&priority)?,
3413        state: task_state_from(&state)?,
3414        on_expiry: expiry_from(&on_expiry)?,
3415        created_at: Timestamp::from_unix_timestamp(row.get::<_, i64>(12))
3416            .map_err(|e| corrupt("unrepresentable created_at", e))?,
3417        due_at: due
3418            .map(Timestamp::from_unix_timestamp)
3419            .transpose()
3420            .map_err(|e| corrupt("unrepresentable due_at", e))?,
3421        withheld: withheld
3422            .map(|w| decoded("withheld reason", &w, crate::core::Withheld::parse(&w)))
3423            .transpose()?,
3424    })
3425}
3426
3427#[async_trait]
3428impl BatchStore for PostgresStore {
3429    fn tenant(&self) -> &str {
3430        crate::journal::JournalStore::tenant(self)
3431    }
3432
3433    async fn open(&self, id: BatchId, plan_digest: &str) -> Result<(), StoreError> {
3434        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3435        client
3436            .execute(
3437                "INSERT INTO batches (batch_id, plan_digest, tenant) VALUES ($1, $2, $3)
3438                 ON CONFLICT (tenant, batch_id) DO NOTHING",
3439                &[
3440                    &id.to_string(),
3441                    &plan_digest.to_owned(),
3442                    &self.tenant_name(),
3443                ],
3444            )
3445            .await
3446            .map_err(|e| be(&e))?;
3447        // Read back what stands: whichever open won the insert, its digest is
3448        // the batch's, and a resume offering another one is refused rather
3449        // than run — one batch runs one frozen plan, and this row is the only
3450        // witness to which.
3451        let stored: String = client
3452            .query_one(
3453                "SELECT plan_digest FROM batches WHERE batch_id = $1 AND tenant = $2",
3454                &[&id.to_string(), &self.tenant_name()],
3455            )
3456            .await
3457            .map_err(|e| be(&e))?
3458            .get(0);
3459        if stored != plan_digest {
3460            return Err(StoreError::BatchPlanChanged {
3461                batch: id.to_string(),
3462                stored,
3463                offered: plan_digest.to_owned(),
3464            });
3465        }
3466        Ok(())
3467    }
3468
3469    async fn plan_digest(&self, id: BatchId) -> Result<Option<String>, StoreError> {
3470        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3471        Ok(client
3472            .query_opt(
3473                "SELECT plan_digest FROM batches WHERE batch_id = $1 AND tenant = $2",
3474                &[&id.to_string(), &self.tenant_name()],
3475            )
3476            .await
3477            .map_err(|e| be(&e))?
3478            .map(|r| r.get(0)))
3479    }
3480
3481    async fn mark_exhausted(&self, id: BatchId) -> Result<(), StoreError> {
3482        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3483        let updated = client
3484            .execute(
3485                "UPDATE batches SET exhausted = TRUE WHERE batch_id = $1 AND tenant = $2",
3486                &[&id.to_string(), &self.tenant_name()],
3487            )
3488            .await
3489            .map_err(|e| be(&e))?;
3490        // The row count says whether the mark landed — see `record`. This is
3491        // the one bit that lets a census read as *finished*, and `Ok` over
3492        // nothing is the quietest way to lose it.
3493        if updated == 0 {
3494            return Err(StoreError::NotFound(id.to_string()));
3495        }
3496        Ok(())
3497    }
3498
3499    async fn is_exhausted(&self, id: BatchId) -> Result<bool, StoreError> {
3500        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3501        Ok(client
3502            .query_opt(
3503                "SELECT exhausted FROM batches WHERE batch_id = $1 AND tenant = $2",
3504                &[&id.to_string(), &self.tenant_name()],
3505            )
3506            .await
3507            .map_err(|e| be(&e))?
3508            .is_some_and(|r| r.get::<_, bool>(0)))
3509    }
3510
3511    async fn reserve(
3512        &self,
3513        batch: BatchId,
3514        key: &str,
3515        run: RunId,
3516    ) -> Result<ItemRecord, StoreError> {
3517        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3518        // `DO NOTHING` then read back: an item already reserved must hand back
3519        // the *original* run id, or the journal holding its effects is orphaned
3520        // and they are performed again.
3521        client
3522            .execute(
3523                "INSERT INTO batch_items (batch_id, item_key, run_id, tenant)
3524                 VALUES ($1, $2, $3, $4)
3525                 ON CONFLICT (tenant, batch_id, item_key) DO NOTHING",
3526                &[
3527                    &batch.to_string(),
3528                    &key.to_owned(),
3529                    &run.to_string(),
3530                    &self.tenant_name(),
3531                ],
3532            )
3533            .await
3534            .map_err(|e| be(&e))?;
3535        let row = client
3536            .query_one(
3537                "SELECT run_id, outcome, detail, tokens, minor FROM batch_items
3538                  WHERE batch_id = $1 AND item_key = $2 AND tenant = $3",
3539                &[&batch.to_string(), &key.to_owned(), &self.tenant_name()],
3540            )
3541            .await
3542            .map_err(|e| be(&e))?;
3543        item_from(&row, key)
3544    }
3545
3546    async fn record(
3547        &self,
3548        batch: BatchId,
3549        key: &str,
3550        outcome: &ItemOutcome,
3551        spend: Spend,
3552    ) -> Result<(), StoreError> {
3553        let detail = match outcome {
3554            ItemOutcome::Succeeded => None,
3555            ItemOutcome::Failed(d)
3556            | ItemOutcome::Quarantined(d)
3557            | ItemOutcome::Suspended(d)
3558            | ItemOutcome::Exhausted(d)
3559            | ItemOutcome::Withheld(d) => Some(d.clone()),
3560        };
3561        let state = outcome.as_str();
3562        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3563        let updated = client
3564            .execute(
3565                "UPDATE batch_items SET outcome = $3, detail = $4, tokens = $5, minor = $6
3566                  WHERE batch_id = $1 AND item_key = $2 AND tenant = $7",
3567                &[
3568                    &batch.to_string(),
3569                    &key.to_owned(),
3570                    &state,
3571                    &detail,
3572                    &sql_amount(spend.tokens),
3573                    &sql_amount(spend.minor_units),
3574                    &self.tenant_name(),
3575                ],
3576            )
3577            .await
3578            .map_err(|e| be(&e))?;
3579        // The predicate did the work; the row count says whether it matched.
3580        // Discarding it reported success for a record that wrote nothing —
3581        // the same lie a release that freed nothing tells.
3582        if updated == 0 {
3583            return Err(StoreError::NotFound(format!("{batch}/{key}")));
3584        }
3585        Ok(())
3586    }
3587
3588    async fn cursor(&self, batch: BatchId) -> Result<Option<String>, StoreError> {
3589        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3590        // The contiguous terminal prefix: an item still running or suspended
3591        // holds the cursor behind it, or a resume steps over work outstanding.
3592        let first_open: Option<String> = client
3593            .query_one(
3594                "SELECT MIN(item_key) FROM batch_items
3595                  WHERE batch_id = $1 AND tenant = $2
3596                    AND (outcome IS NULL OR NOT (outcome = ANY($3::text[])))",
3597                &[
3598                    &batch.to_string(),
3599                    &self.tenant_name(),
3600                    &ItemOutcome::terminal_tags(),
3601                ],
3602            )
3603            .await
3604            .map_err(|e| be(&e))?
3605            .get(0);
3606
3607        let row = match first_open {
3608            Some(open) => client
3609                .query_one(
3610                    "SELECT MAX(item_key) FROM batch_items
3611                      WHERE batch_id = $1 AND item_key < $2 AND tenant = $3",
3612                    &[&batch.to_string(), &open, &self.tenant_name()],
3613                )
3614                .await
3615                .map_err(|e| be(&e))?,
3616            None => client
3617                .query_one(
3618                    "SELECT MAX(item_key) FROM batch_items
3619                      WHERE batch_id = $1 AND tenant = $2",
3620                    &[&batch.to_string(), &self.tenant_name()],
3621                )
3622                .await
3623                .map_err(|e| be(&e))?,
3624        };
3625        Ok(row.get(0))
3626    }
3627
3628    async fn census(&self, batch: BatchId) -> Result<BatchCensus, StoreError> {
3629        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3630        let rows = client
3631            .query(
3632                "SELECT outcome, COUNT(*), COALESCE(SUM(tokens),0), COALESCE(SUM(minor),0)
3633                   FROM batch_items WHERE batch_id = $1 AND tenant = $2
3634                  GROUP BY outcome",
3635                &[&batch.to_string(), &self.tenant_name()],
3636            )
3637            .await
3638            .map_err(|e| be(&e))?;
3639
3640        let mut c = BatchCensus::default();
3641        for row in rows {
3642            let outcome: Option<String> = row.get(0);
3643            let n: i64 = row.get(1);
3644            let tokens: i64 = row.get(2);
3645            let minor: i64 = row.get(3);
3646            let n = amount_of(n);
3647            match outcome.as_deref() {
3648                // Bucketed off the *parsed* outcome, so an outcome added later
3649                // is a compiler error here rather than a row in no bucket.
3650                Some(tag) => {
3651                    match decoded("item outcome", tag, ItemOutcome::parse(tag, String::new()))? {
3652                        ItemOutcome::Succeeded => c.succeeded = n,
3653                        ItemOutcome::Failed(_) => c.failed = n,
3654                        ItemOutcome::Quarantined(_) => c.quarantined = n,
3655                        ItemOutcome::Suspended(_) => c.suspended = n,
3656                        ItemOutcome::Exhausted(_) => c.exhausted = n,
3657                        ItemOutcome::Withheld(_) => c.withheld = n,
3658                    }
3659                }
3660                // Reserved, no outcome recorded.
3661                None => c.in_flight = n,
3662            }
3663            c.spend += Spend {
3664                tokens: amount_of(tokens),
3665                minor_units: amount_of(minor),
3666            };
3667        }
3668        Ok(c)
3669    }
3670
3671    async fn items(&self, batch: BatchId, limit: usize) -> Result<Vec<ItemRecord>, StoreError> {
3672        self.query_items(
3673            batch,
3674            limit,
3675            "SELECT item_key, run_id, outcome, detail, tokens, minor FROM batch_items
3676              WHERE batch_id = $1 AND tenant = $3 ORDER BY item_key ASC LIMIT $2",
3677            &[],
3678        )
3679        .await
3680    }
3681
3682    async fn items_needing_attention(
3683        &self,
3684        batch: BatchId,
3685        limit: usize,
3686    ) -> Result<Vec<ItemRecord>, StoreError> {
3687        // `<> ALL` rather than `<> 'succeeded'`, and the set comes from the
3688        // type: an outcome added later lands in this listing without anyone
3689        // remembering to widen a literal here.
3690        //
3691        // `IS NULL OR` is load-bearing: a reserved item has no outcome at all,
3692        // and `NULL <> ALL (...)` is NULL, which `WHERE` treats as false — so
3693        // without it the interrupted items, the ones a crash left mid-flight,
3694        // would be exactly the ones this listing dropped.
3695        self.query_items(
3696            batch,
3697            limit,
3698            "SELECT item_key, run_id, outcome, detail, tokens, minor FROM batch_items
3699              WHERE batch_id = $1 AND tenant = $3
3700                AND (outcome IS NULL OR outcome <> ALL($4))
3701              ORDER BY item_key ASC LIMIT $2",
3702            &ItemOutcome::settled_tags(),
3703        )
3704        .await
3705    }
3706}
3707
3708impl PostgresStore {
3709    /// One shape of item query, so the two listings cannot drift in how they
3710    /// decode a row or order it.
3711    async fn query_items(
3712        &self,
3713        batch: BatchId,
3714        limit: usize,
3715        sql: &str,
3716        settled: &[&str],
3717    ) -> Result<Vec<ItemRecord>, StoreError> {
3718        let client = self.pool().get().await.map_err(|e| pool_err(&e))?;
3719        let batch_key = batch.to_string();
3720        let capped = i64::try_from(limit).unwrap_or(i64::MAX);
3721        let tenant = self.tenant_name();
3722        let settled: Vec<String> = settled.iter().map(|s| (*s).to_owned()).collect();
3723        let mut params: Vec<&(dyn tokio_postgres::types::ToSql + Sync)> =
3724            vec![&batch_key, &capped, &tenant];
3725        if !settled.is_empty() {
3726            params.push(&settled);
3727        }
3728        let rows = client.query(sql, &params).await.map_err(|e| be(&e))?;
3729        let mut out = Vec::with_capacity(rows.len());
3730        for row in rows {
3731            let key: String = row.get(0);
3732            out.push(ItemRecord {
3733                key: key.clone(),
3734                run: RunId::parse(&row.get::<_, String>(1))
3735                    .map_err(|e| corrupt("bad run id", e))?,
3736                outcome: outcome_from(row.get::<_, Option<String>>(2), row.get(3))?,
3737                spend: Spend {
3738                    tokens: amount_of(row.get::<_, i64>(4)),
3739                    minor_units: amount_of(row.get::<_, i64>(5)),
3740                },
3741            });
3742        }
3743        Ok(out)
3744    }
3745}
3746
3747/// `None` means *no outcome yet* — and only that. An outcome string this
3748/// store cannot read is damage, not absence: decoded to `None` it would say
3749/// the item never ran, and the census would carry the row as in-flight
3750/// forever, keeping the batch `Running` over damage nobody is told about.
3751fn outcome_from(
3752    state: Option<String>,
3753    detail: Option<String>,
3754) -> Result<Option<ItemOutcome>, StoreError> {
3755    let d = detail.unwrap_or_default();
3756    let Some(state) = state else {
3757        return Ok(None);
3758    };
3759    decoded("item outcome", &state, ItemOutcome::parse(&state, d)).map(Some)
3760}
3761
3762fn item_from(row: &tokio_postgres::Row, key: &str) -> Result<ItemRecord, StoreError> {
3763    let run: String = row.get(0);
3764    Ok(ItemRecord {
3765        key: key.to_owned(),
3766        run: RunId::parse(&run).map_err(|e| corrupt("bad run id", e))?,
3767        outcome: outcome_from(row.get::<_, Option<String>>(1), row.get(2))?,
3768        spend: Spend {
3769            tokens: amount_of(row.get::<_, i64>(3)),
3770            minor_units: amount_of(row.get::<_, i64>(4)),
3771        },
3772    })
3773}
3774
3775#[cfg(test)]
3776mod codec_tests {
3777    use super::*;
3778
3779    /// Every string this store writes is one it reads back as the same value.
3780    ///
3781    /// The encoder is exhaustive over the type, so the pair cannot drift while
3782    /// this passes: a variant added to `Phase`, `OnExpiry` or `Priority`
3783    /// changes the encoder's output and this fails on the decoder's refusal.
3784    #[test]
3785    fn every_written_string_decodes_to_the_value_that_wrote_it() {
3786        use crate::core::Phase;
3787        for phase in [Phase::Forward, Phase::Compensating] {
3788            assert_eq!(
3789                phase_from(crate::core::Phase::as_str(phase)).expect("round trip"),
3790                phase
3791            );
3792        }
3793        for expiry in [OnExpiry::Deny, OnExpiry::Escalate, OnExpiry::Proceed] {
3794            assert_eq!(
3795                expiry_from(OnExpiry::as_str(expiry)).expect("round trip"),
3796                expiry
3797            );
3798        }
3799        for priority in [
3800            Priority::Low,
3801            Priority::Normal,
3802            Priority::High,
3803            Priority::Urgent,
3804        ] {
3805            assert_eq!(
3806                priority_from(priority.as_str()).expect("round trip"),
3807                priority
3808            );
3809        }
3810    }
3811
3812    /// **A row this store cannot read is not a row with a default value.**
3813    ///
3814    /// A decoder that answers with a fallback cannot report damage, so the
3815    /// damage arrives as a decision. `phase` is the one with teeth: it tells a
3816    /// step's forward pass from its compensating one, so an unreadable value
3817    /// answered `Forward` hands the unwind logic a compensating record wearing
3818    /// the wrong half of the saga. `Deny` and `Normal` are *safe* values,
3819    /// which is exactly what makes them the wrong answer — a fail-closed
3820    /// default is still a fact invented about a row nobody could read.
3821    #[test]
3822    fn an_unreadable_column_is_refused_rather_than_defaulted() {
3823        for (what, err) in [
3824            ("phase", phase_from("").err()),
3825            ("phase", phase_from("Compensating").err()),
3826            ("expiry", expiry_from("escalate_later").err()),
3827            ("priority", priority_from("normal ").err()),
3828            (
3829                "item outcome",
3830                outcome_from(Some("done".into()), None).err(),
3831            ),
3832        ] {
3833            assert!(
3834                matches!(err, Some(StoreError::Corrupt { .. })),
3835                "an unrecognised {what} decoded to a value instead of refusing"
3836            );
3837        }
3838    }
3839}