1use 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
39fn 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
525fn 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
565fn 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
581fn 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 ¬e.at.unix_timestamp(),
1494 ¬e.by.actor(),
1495 ¬e.note,
1496 ¬e.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 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 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 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 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
1830fn 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 &[
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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
2796const 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 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 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
2926async 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 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 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 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 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 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 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 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 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 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 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 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 if !task.state.is_pending() {
3262 return Ok(task);
3263 }
3264 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 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 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 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 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 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 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 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 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 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 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 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 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, ¶ms).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
3747fn 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 #[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 #[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}