agentplane/case/mod.rs
1//! Case storage: correlation, state, obligations, and inbound events.
2
3mod events;
4mod tasks;
5mod timers;
6
7pub use events::{
8 ADDRESSEE_CONCLUDED_REASON, BufferedEvent, ERASED_REASON, EventStore, Minter, Retired,
9 TargetedDelivery,
10};
11pub use tasks::{ClaimError, TaskStore};
12pub use timers::TimerStore;
13
14use std::fmt::Debug;
15
16use async_trait::async_trait;
17use serde_json::Value;
18
19use crate::core::{
20 BreachNote, Case, CaseId, CaseStatus, CaseVersion, CorrelationKey, Deadline, DeadlineState,
21 Digest, LegalHold, RunId, StoreError, Timestamp,
22};
23
24/// What admission did with an inbound trigger.
25#[derive(Debug, Clone, Copy, PartialEq, Eq)]
26pub enum Correlation {
27 /// No open case matched; a new one was created.
28 Opened(CaseId),
29 /// An open case matched and this run joined it.
30 Attached(CaseId),
31}
32
33impl Correlation {
34 #[must_use]
35 pub fn case_id(self) -> CaseId {
36 match self {
37 Self::Opened(id) | Self::Attached(id) => id,
38 }
39 }
40
41 #[must_use]
42 pub fn is_new(self) -> bool {
43 matches!(self, Self::Opened(_))
44 }
45}
46
47/// Long-lived case state, correlated by business key.
48///
49/// Correlation is a deterministic lookup — never a model call. It runs at
50/// admission, before planning, because which case a message belongs to is a
51/// question of fact, not of judgement.
52#[async_trait]
53pub trait CaseStore: Send + Sync + Debug {
54 /// Whose rows this handle can reach.
55 ///
56 /// Defaults to [`TenantId::DEFAULT`](crate::core::TenantId::DEFAULT), the
57 /// tenant a store serves until told otherwise. Override it with the tenant
58 /// the handle is actually scoped to.
59 ///
60 /// This exists so a mismatch with the plane's tenant is a **startup
61 /// refusal**. When a key ring is wired, `build()` seals case state under
62 /// the plane's tenant while the store writes rows under its own; the two
63 /// disagreeing is not a leak — the scopes simply differ — but it puts case
64 /// state under a scope `erase_case` will never destroy. That is an erasure
65 /// that reports success and misses, which is the one failure a deletion
66 /// guarantee cannot have.
67 fn tenant(&self) -> &str {
68 crate::core::TenantId::DEFAULT
69 }
70
71 /// Find an **open** case matching any of these keys.
72 ///
73 /// Closed cases are not matched: a new message about a settled matter opens
74 /// a new case rather than reanimating one that was concluded and audited.
75 async fn correlate(&self, keys: &[CorrelationKey]) -> Result<Option<CaseId>, StoreError>;
76
77 /// Correlate, or open a new case if nothing matched.
78 ///
79 /// Implementations must make this atomic. Two messages for the same new case
80 /// arriving concurrently must produce one case, not two — otherwise a
81 /// process fragments across cases and its obligations are tracked in neither.
82 async fn correlate_or_open(
83 &self,
84 kind: &str,
85 keys: &[CorrelationKey],
86 at: Timestamp,
87 ) -> Result<Correlation, StoreError>;
88
89 /// Fetch one case.
90 ///
91 /// Named for what it returns rather than `get`: a single store commonly
92 /// implements several of these traits, and same-named methods make every
93 /// call site ambiguous.
94 async fn case(&self, id: CaseId) -> Result<Option<Case>, StoreError>;
95
96 /// Every case, one bounded page at a time, in stable id order.
97 ///
98 /// This is the export's read. `by_status` cannot serve it: a bounded list
99 /// with no cursor enumerates a prefix and calls it everything — the silent
100 /// truncation this project refuses elsewhere, at the one boundary whose
101 /// whole job is completeness. `after` is the last id a caller saw, and
102 /// paging resumes strictly beyond it.
103 ///
104 /// Ordered by id rather than by anything business-shaped, because the
105 /// order's only job is that two pages never overlap and never gap.
106 ///
107 /// Returns state **as stored**. A sealing decorator does not open it here,
108 /// unlike [`case`](Self::case): this read exists for the export, and an
109 /// export of plaintext would quietly undo erasure — the key destroyed
110 /// tomorrow would no longer reach the copy taken today.
111 async fn cases(&self, after: Option<CaseId>, limit: usize) -> Result<Vec<Case>, StoreError>;
112
113 /// Write one complete case, exactly as given — an **import authority**, not
114 /// a runtime path.
115 ///
116 /// The ordinary write paths refuse to say some of what a restore must:
117 /// `put_state` allocates versions one at a time, `correlate_or_open` mints
118 /// a fresh id, and neither can reproduce a case at version 4 000 with the
119 /// id every exported record already names. This is the same seam governed
120 /// memory keeps for the same reason — direct complete-item writes are a
121 /// deployment/import authority at the store boundary, never something a
122 /// skill reaches.
123 ///
124 /// Implementations must leave the imported case reachable by **every**
125 /// read path — `case`, `correlate`, `by_status`, `due`, `blobs_of` — which
126 /// is what the conformance battery holds them to: the read paths are the
127 /// check on this method's index maintenance, because an import that
128 /// rebuilds five indexes out of six reads perfectly until somebody queries
129 /// the sixth.
130 ///
131 /// # Errors
132 ///
133 /// A store that already holds this case id refuses: a restore rebuilds a
134 /// case layer, it does not merge one.
135 async fn import_case(
136 &self,
137 case: &Case,
138 deadlines: &[Deadline],
139 blobs: &[Digest],
140 ) -> Result<(), StoreError>;
141
142 /// Record that a run touched this case.
143 async fn attach_run(&self, case: CaseId, run: RunId) -> Result<(), StoreError>;
144
145 /// Record that a case produced a blob.
146 ///
147 /// The case is what an erasure request actually names — nobody asks to
148 /// forget a digest — so something has to know which bytes belong to which
149 /// matter. That association cannot live in the blob store, which is
150 /// content-addressed on purpose and has no idea what a case is, and it
151 /// cannot be recomputed later: a digest is deliberately not reversible.
152 ///
153 /// Recording the same blob twice is the same record. Two runs on one case
154 /// storing identical bytes land on one digest by construction, and that is
155 /// one artifact, not two.
156 ///
157 /// # Errors
158 ///
159 /// If the store rejects the write.
160 async fn link_blob(
161 &self,
162 case: CaseId,
163 digest: Digest,
164 at: Timestamp,
165 ) -> Result<(), StoreError>;
166
167 /// Every blob this case produced, oldest first.
168 ///
169 /// # Errors
170 ///
171 /// If the store cannot be read.
172 async fn blobs_of(&self, case: CaseId) -> Result<Vec<Digest>, StoreError>;
173
174 /// Replace a case's state, if it is still at `expected`.
175 ///
176 /// # Why this takes a version
177 ///
178 /// A case is shared by every run correlated to it, and the window between
179 /// reading its state and writing it back contains a model call — which is
180 /// unbounded. Two runs on one case *will* overlap, and a blind write in that
181 /// window silently discards whichever one lost. See [`CaseVersion`].
182 ///
183 /// Implementations must make the check part of the write itself — a single
184 /// `UPDATE ... WHERE version = ?` — and not a read followed by a write. A
185 /// check performed in application code is a check the next caller can race.
186 ///
187 /// # Errors
188 ///
189 /// * [`StoreError::CaseConflict`] if the case has moved past `expected`. The
190 /// caller re-reads and decides again; retrying the same write is the lost
191 /// update this exists to prevent.
192 /// * [`StoreError::NotFound`] if there is no such case. Implementations must
193 /// tell these apart — reporting a missing case as a conflict sends the
194 /// caller into a re-read loop against something that will never exist.
195 async fn put_state(
196 &self,
197 case: CaseId,
198 expected: CaseVersion,
199 state: Value,
200 ) -> Result<CaseVersion, StoreError>;
201
202 /// Move a case to a status.
203 ///
204 /// `Closed` routes through [`close`](Self::close), because closure releases
205 /// the correlation keys as well as writing the column, and the two spellings
206 /// of *closed* must not drift.
207 ///
208 /// **Leaving `Closed` re-claims them**, which is the same rule read
209 /// backwards. Without it a matter reopened by any route — a run calling
210 /// `set_case_status`, or the sweep escalating over an expired task — comes
211 /// back as a case no inbound message can ever correlate to again, which
212 /// looks exactly like a live matter and is not one. A key another case has
213 /// since claimed stays with that case: the identifier belongs to whichever
214 /// matter is open for it now, and reopening must not take one back.
215 async fn set_status(&self, case: CaseId, status: CaseStatus) -> Result<(), StoreError>;
216
217 /// Close a case.
218 ///
219 /// **Fails while an obligation is open.** A case with an unmet deadline
220 /// cannot be closed silently: that is precisely how a missed regulatory
221 /// window becomes invisible.
222 async fn close(&self, case: CaseId) -> Result<(), StoreError>;
223
224 /// Register an obligation. The resolved instant is stored as given and never
225 /// recomputed.
226 ///
227 /// # Errors
228 ///
229 /// [`StoreError::CaseClosed`] if the matter is closed. [`close`](Self::close)
230 /// refuses while an obligation is outstanding, and that check is worth
231 /// nothing on its own: it holds at one instant, and this is the write that
232 /// walks past it afterwards. Both halves are needed for *a closed case owes
233 /// nothing* to be a property of the store rather than of the order two
234 /// callers happened to run in.
235 ///
236 /// Implementations must make the two decide one at a time. The redb backend
237 /// gets that from its single write transaction; a SQL backend has to take
238 /// the case row's lock, because two snapshots each reading the other's
239 /// pre-state is how both writes commit and the matter ends up closed and
240 /// owing.
241 async fn register_deadline(&self, deadline: &Deadline) -> Result<(), StoreError>;
242
243 async fn deadlines(&self, case: CaseId) -> Result<Vec<Deadline>, StoreError>;
244
245 async fn set_deadline_state(
246 &self,
247 case: CaseId,
248 name: &str,
249 state: DeadlineState,
250 ) -> Result<(), StoreError>;
251
252 /// Breach an obligation and escalate its case, in one transaction, if it
253 /// is still outstanding and due at `now`. Returns whether it applied.
254 ///
255 /// One verb because the sweep decides from a read that is already stale
256 /// when it acts: a run can meet the obligation and close the case in
257 /// between. Spelled as separate writes, the sweep then escalates a closed
258 /// case — reopening it and re-claiming its correlation keys — and breaches
259 /// an obligation that was met. Here the check and both writes are one
260 /// decision, and `false` means somebody got there first.
261 ///
262 /// The case is escalated only while it is not
263 /// [`Closed`](CaseStatus::Closed).
264 ///
265 /// An applied breach is also marked as **owing its account**, in the same
266 /// transaction: the sweep acts first and writes its notes after, so a
267 /// crash between the two leaves a breach [`due`](Self::due) no longer
268 /// lists. [`breaches_to_note`](Self::breaches_to_note) is where the next
269 /// tick finds it, and [`mark_breach_noted`](Self::mark_breach_noted) is
270 /// what the notes landing retires.
271 ///
272 /// # Errors
273 ///
274 /// [`StoreError::NotFound`] if the case or the obligation does not exist.
275 async fn breach_deadline(
276 &self,
277 case: CaseId,
278 name: &str,
279 now: Timestamp,
280 ) -> Result<bool, StoreError>;
281
282 /// Breaches [`breach_deadline`](Self::breach_deadline) applied whose
283 /// account the sweep has not yet written, longest-overdue first.
284 async fn breaches_to_note(&self, limit: usize) -> Result<Vec<Deadline>, StoreError>;
285
286 /// Record that a breach's account is on the journal. Idempotent.
287 ///
288 /// # Errors
289 ///
290 /// [`StoreError::NotFound`] if the obligation does not exist.
291 async fn mark_breach_noted(&self, case: CaseId, name: &str) -> Result<(), StoreError>;
292
293 /// Obligations that are due or approaching, oldest first.
294 ///
295 /// The sweep that turns a passing instant into an escalation. Without it a
296 /// deadline is a stored number that nobody reads.
297 async fn due(&self, now: Timestamp, limit: usize) -> Result<Vec<Deadline>, StoreError>;
298
299 /// Missed obligations nobody has accounted for, longest-overdue first.
300 ///
301 /// Reads the obligation's own row, which outlives every status its case
302 /// passes through — the escalation a breach causes is retired by
303 /// [`close`](Self::close), and closure is when people stop looking.
304 ///
305 /// Not derivable from [`due`](Self::due), which answers *what still needs
306 /// attention*. A breach is past attention; it needs an account.
307 ///
308 /// **Acknowledged breaches are excluded**, which is what makes this a backlog
309 /// rather than a level that only rises: taken oldest-first and bounded by a
310 /// page, a listing nothing ever leaves shows the same head forever, and the
311 /// entries an operator could still act on are the ones they never see.
312 /// [`acknowledge_breach`](Self::acknowledge_breach) is the verb that empties
313 /// it; the obligation stays [`Breached`](crate::core::DeadlineState::Breached)
314 /// either way, because what ends is the question, not the fact.
315 async fn breached(&self, limit: usize) -> Result<Vec<Deadline>, StoreError>;
316
317 /// Record that somebody has accounted for a breach.
318 ///
319 /// Takes it off [`breached`](Self::breached) and leaves the obligation
320 /// `Breached`, with the account readable through
321 /// [`deadlines`](Self::deadlines) — the fact and the answer to it are
322 /// different things and both are kept.
323 ///
324 /// **Idempotent, first account wins.** Returns whether this call was the one
325 /// that recorded it, so a retry cannot rewrite who looked or when — the same
326 /// rule [`KeyRing::destroy`] follows for the same reason.
327 ///
328 /// # Errors
329 ///
330 /// [`StoreError::NotFound`] if the case or the obligation does not exist,
331 /// and [`StoreError::NotBreached`] if it exists and has not been breached —
332 /// accepting an account before the breach would take an obligation off the
333 /// listing while it was still going to be missed.
334 ///
335 /// [`KeyRing::destroy`]: crate::keyring::KeyRing::destroy
336 async fn acknowledge_breach(
337 &self,
338 case: CaseId,
339 name: &str,
340 note: &BreachNote,
341 ) -> Result<bool, StoreError>;
342
343 /// Preserve this matter against every erasure verb until somebody lifts it.
344 ///
345 /// The only thing in this crate that makes an erasure **fail** rather than
346 /// succeed, and it exists because a retention pass is automatic: it runs on
347 /// a window nobody re-reads, and a matter under a preservation order looks
348 /// exactly like every other closed case old enough to sweep.
349 ///
350 /// **Idempotent, first placement wins**, returning whether this call placed
351 /// it — [`acknowledge_breach`](Self::acknowledge_breach)'s rule, for its
352 /// reason: a retry must not rewrite when the hold began or why.
353 ///
354 /// # Errors
355 ///
356 /// [`StoreError::NotFound`] if the case does not exist — a hold on a matter
357 /// that is not there preserves nothing and reads as effective in the
358 /// listing.
359 async fn place_hold(&self, case: CaseId, hold: &LegalHold) -> Result<bool, StoreError>;
360
361 /// Lift a hold, returning whether one was there to lift.
362 ///
363 /// The verb that empties [`holds`](Self::holds); without it that listing is
364 /// a level that only rises, which [`breached`](Self::breached) refuses for
365 /// the same reason.
366 ///
367 /// **The register answers *what is preserved now*, not *what ever was*.**
368 /// The register keeps no row for a released hold, deliberately: a kept row
369 /// would make its free-text reason outlive the matter it preserved. The
370 /// history is the journal's —
371 /// [`Runtime::release_hold`](crate::runtime::Runtime::release_hold)
372 /// records who released the hold, when, and who had placed it, stamped
373 /// with the case and without the reason, before removing the hold through
374 /// [`release_hold_if`](Self::release_hold_if); both of the plane's doors
375 /// release through that. This is the raw register verb: a
376 /// direct call is the embedder's own act, and nothing records it.
377 ///
378 /// # Errors
379 ///
380 /// [`StoreError::NotFound`] if the case does not exist.
381 async fn release_hold(&self, case: CaseId) -> Result<bool, StoreError>;
382
383 /// Lift the hold on `case` only if it is still `standing` — same instant,
384 /// same reason, same operator.
385 ///
386 /// Answers whether that hold was removed. `false` when the matter is not
387 /// held or holds a different hold; a different hold is left standing.
388 ///
389 /// The compare and the removal MUST be one write. A release followed by a
390 /// re-place puts a new hold where the old one was, so an unconditional
391 /// removal by a releaser that read the old one would delete the new hold
392 /// under a record that names the old.
393 /// [`Runtime::release_hold`](crate::runtime::Runtime::release_hold)
394 /// removes through this, with the hold it journaled.
395 ///
396 /// # Errors
397 ///
398 /// [`StoreError::NotFound`] if the case does not exist; otherwise if the
399 /// store cannot be reached, or holds a row it cannot read.
400 async fn release_hold_if(&self, case: CaseId, standing: &LegalHold)
401 -> Result<bool, StoreError>;
402
403 /// The hold on one matter, if it is held.
404 ///
405 /// The question every erasure verb asks before it destroys anything.
406 async fn hold(&self, case: CaseId) -> Result<Option<LegalHold>, StoreError>;
407
408 /// Every matter under hold, **oldest hold first**.
409 ///
410 /// The half that makes this a control: a store answering only
411 /// [`hold`](Self::hold) can be queried solely by somebody who already knows
412 /// which case to ask about, which delivers nothing to the person whose job
413 /// is to find out what is still preserved and why.
414 ///
415 /// Ascending legitimately — [`release_hold`](Self::release_hold) removes
416 /// entries, so the head is not permanent, and the longest-standing unlifted
417 /// hold is the one that needs a question asked about it.
418 async fn holds(
419 &self,
420 after: Option<CaseId>,
421 limit: usize,
422 ) -> Result<Vec<(CaseId, LegalHold)>, StoreError>;
423
424 /// Cases matching a status, newest first.
425 async fn by_status(&self, status: CaseStatus, limit: usize) -> Result<Vec<Case>, StoreError>;
426
427 /// How much is open right now, for the gauges in `runtime::metrics`.
428 ///
429 /// Deliberately not expressible as `by_status(...).len()`: that is bounded
430 /// by a `limit`, and a gauge computed from a truncated list is a number that
431 /// stops rising exactly when the backlog becomes worth knowing about. It is
432 /// also the only consumer of a case's `opened_at` — a count alone cannot
433 /// distinguish ten cases open for an hour from ten open for a month.
434 ///
435 /// `now` is passed in so the reading is testable against arbitrary ageing
436 /// and needs no escape from the determinism gate.
437 async fn census(&self, now: Timestamp) -> Result<CaseCensus, StoreError>;
438
439 /// Record what the last recovery rehearsal found.
440 ///
441 /// **One row, replaced.** A history of drills is a different artifact and
442 /// a larger promise; what an audit asks is *when did you last rehearse and
443 /// did it pass*, and a row that the latest write replaces answers exactly
444 /// that without becoming a table nobody prunes.
445 ///
446 /// # Errors
447 ///
448 /// [`StoreError`] if the store is unreachable. A rehearsal whose verdict
449 /// could not be written is reported to its caller rather than swallowed:
450 /// the next audit would otherwise read the *previous* drill's date as the
451 /// most recent one.
452 async fn record_drill(&self, record: &DrillRecord) -> Result<(), StoreError>;
453
454 /// The last recovery rehearsal, or `None` if this plane has never run one.
455 ///
456 /// `None` is the honest answer and a meaningful one: *nobody has
457 /// rehearsed* is what an auditor most wants to know and is exactly what a
458 /// missing CI log cannot distinguish from *the log rotated*.
459 ///
460 /// # Errors
461 ///
462 /// [`StoreError`] if the store is unreachable.
463 async fn last_drill(&self) -> Result<Option<DrillRecord>, StoreError>;
464}
465
466/// What the last recovery rehearsal found, and when.
467///
468/// # Why a row rather than a journal record
469///
470/// A drill is about the **plane**, not a run: it walks the case layer against
471/// the blob store and the key ring it is actually wired to. It has no run id
472/// and no chain to append to, which is the position a standing halt is in and
473/// gets the same answer — a row that the latest write replaces.
474///
475/// # Why it is stored at all
476///
477/// What an audit asks is not *can you rehearse* but **when you last did, and
478/// whether it passed**. Left in the process that ran the verb, that answer is
479/// a CI log or a wiki page — the artifact this project argues against relying
480/// on everywhere else. `serve --drill-every` makes the rehearsal a scheduled
481/// job; this makes its verdict a fact the plane can be asked for.
482///
483/// # What it does not carry
484///
485/// The findings themselves. A drill's report names each unrecoverable
486/// reference and is unbounded in a way a single row must not be; what belongs
487/// here is the verdict and enough to re-run for detail. The counts are kept
488/// **separately from `sound`** because a pass over nothing is not a pass: a
489/// drill that could not check the blob store reports no findings and has
490/// established almost nothing, which is the same rule
491/// [`AuditReport::not_checked`](crate::audit::AuditReport::not_checked) is
492/// held to.
493#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
494pub struct DrillRecord {
495 /// When the rehearsal ran. The caller's clock: a drill is an operational
496 /// act in the outside world, not a journaled observation.
497 #[serde(with = "time::serde::rfc3339")]
498 pub at: Timestamp,
499 /// Whether it found nothing unrecoverable.
500 pub sound: bool,
501 /// Cases walked.
502 pub cases: u64,
503 /// Unrecoverable references found.
504 pub findings: u64,
505 /// Questions this pass could not answer — an absent blob store, a key
506 /// ring it was not given. Non-zero with `sound` true means *incomplete*,
507 /// not *clean*.
508 pub not_checked: u64,
509 /// The log this plane was serving when it ran, and how far along it was.
510 ///
511 /// **Which store, answered so a reader cannot be fooled by the obvious
512 /// mistake**: a drill against a restored copy proves that copy
513 /// recoverable and says nothing about production. An origin and a size
514 /// pin both, and both come from the plane's own checkpoint rather than
515 /// from a string somebody typed.
516 pub origin: String,
517 /// The checkpoint size at the time of the rehearsal.
518 pub size: u64,
519}
520
521/// What the case store is currently holding.
522#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
523pub struct CaseCensus {
524 /// Cases in any status other than `Closed`.
525 pub open: u64,
526 /// The longest-open case's age in seconds, or `None` if none are open.
527 pub oldest_age_secs: Option<u64>,
528 /// Obligations at or past their instant, still `Pending` or `Warned`.
529 pub due: u64,
530 /// Breaches nobody has accounted for.
531 ///
532 /// A gauge rather than a counter, and the distinction is the point: this
533 /// number falls when somebody acknowledges one, so an alert on it says
534 /// *there is unattended work* rather than *this deployment has ever missed
535 /// something*. A monotonic instrument cannot express the first, which is
536 /// the only one anybody can act on.
537 pub breached: u64,
538}