agentplane/export.rs
1//! Getting the record out, in a form nothing here has to be present to read.
2//!
3//! # Why this is a deliverable and not a `serde` derive
4//!
5//! [`audit`](crate::audit) exists because the party under examination must not
6//! also be the only party able to examine. That argument has a second half this
7//! crate did not have: an auditor who can *check* the history but cannot
8//! *obtain* it is still dependent on the operator, and a regulator asking a
9//! financial entity to demonstrate an exit is asking about obtaining, not
10//! checking. A store nobody can get data out of is a concentration risk with a
11//! hash chain on top.
12//!
13//! So the export is a first-class operation with three properties, and each one
14//! is a refusal of an easier design:
15//!
16//! * **Streaming, one JSON object per line.** A whole-journal `Vec` is a
17//! memory ceiling disguised as an API, and the export that matters most is
18//! the one taken from the largest store. JSON Lines also means an interrupted
19//! export is a *prefix* rather than a corrupt document — which is the failure
20//! an operator actually hits.
21//! * **Self-describing.** The first line is a header naming the log, its
22//! checkpoint, and the canonicalization rule the digests were computed under.
23//! Without that, an export is bytes an auditor has to be told how to read,
24//! and being told is the dependency this module exists to remove.
25//! * **It says what it did not export.** The trailer carries the counts and any
26//! run that could not be read. A truncated export shaped exactly like a
27//! complete one is the failure this project refuses everywhere else, and it
28//! is worst here: the missing run is the interesting one.
29//!
30//! # What it deliberately does not do
31//!
32//! It does not decrypt. With a key ring configured the journal commits to
33//! ciphertext, and an export of plaintext would quietly undo
34//! [erasure](crate::keyring) — destroying the key would no longer reach the
35//! copy somebody exported last month. The export carries what the chain
36//! committed to, which is also what verifies.
37//!
38//! It does not re-verify. [`audit`](crate::audit) answers *is this sound*, this
39//! answers *here it is*, and folding them would produce an export that refuses
40//! to emit the very history an auditor wants to examine *because* it is
41//! suspect.
42//!
43//! It is scoped to one tenant, because a [`JournalStore`] handle is. There is
44//! no argument here that could widen it, which is the same reason the rest of
45//! the tenancy story is in keys rather than in filters.
46
47use std::sync::Arc;
48
49use crate::core::{RunId, StoreError};
50use crate::journal::{Append, Checkpoint, JournalStore};
51
52/// The export format's own version — see [`Header::version`].
53///
54/// One constant, because three readers consume it: the writer stamps it, the
55/// verifier refuses what it cannot interpret, and the restore refuses what it
56/// cannot faithfully replay. A version that only the writer knew about would be
57/// a declaration that does nothing — a reader would parse a future format as
58/// far as the lines happened to look familiar, and report findings about a
59/// file it never understood.
60///
61/// Two shapes are load-bearing enough to state with the constant, because
62/// each was once tempting to do the other way. The case layer is mandatory,
63/// never an optional extension: a reader that tolerated its absence could not
64/// tell *this plane has no cases* from *the case layer was dropped from this
65/// file* — and the second is the finding that matters. And every record line
66/// carries `raw`, the **exact bytes the chain hashed**, which is what
67/// verification recomputes over: verifying a re-serialization of the parsed
68/// body would hold only while this build's canonicalization agreed
69/// byte-for-byte with the writer's — the wire-bytes rule the journal itself
70/// refuses to bend, bent by its own export.
71pub const FORMAT_VERSION: u32 = 1;
72
73/// The first line of an export: what this is and how to read it.
74#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
75pub struct Header {
76 /// Always `"agentplane.export"`, so a reader can tell this file from any
77 /// other line-delimited JSON without being told what it is.
78 pub kind: &'static str,
79 /// The export format's own version, which is **not** the crate's.
80 ///
81 /// A reader pins this. Tying it to the crate version would make every
82 /// release look like a format change to anyone parsing defensively.
83 pub version: u32,
84 /// The log this came from, and its commitment at the moment of export.
85 pub checkpoint: Checkpoint,
86 /// Which canonicalization rule produced the digests in these records.
87 ///
88 /// A digest is meaningless without the rule that computed it, and an
89 /// export outlives the build that wrote it — so the rule travels with the
90 /// digests rather than being whatever the reader happens to implement.
91 pub canon: u16,
92}
93
94/// The first line of a disclosure package: an export of chosen runs, which
95/// says so.
96pub const DISCLOSURE_KIND: &str = "agentplane.disclosure";
97
98/// What a disclosure package was asked for: whole cases, single runs, or both.
99///
100/// A package carries the runs these resolve to, each sealed one with its path
101/// against the header's checkpoint, and only the cases those runs belong to.
102/// Nothing about the log's other leaves is disclosed, and no reader may take
103/// the package for a whole export: its header has its own `kind`, and the
104/// run blocks carry `proof`, a member a whole export never has.
105#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
106#[serde(deny_unknown_fields)]
107pub struct Selection {
108 pub cases: Vec<crate::core::CaseId>,
109 pub runs: Vec<RunId>,
110}
111
112/// A disclosure package's first line: the export header plus its selection.
113#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
114#[serde(deny_unknown_fields)]
115pub struct PackageHeader {
116 /// Always [`DISCLOSURE_KIND`].
117 pub kind: String,
118 /// [`FORMAT_VERSION`]: a package is a framing of the export format and
119 /// moves with it.
120 pub version: u32,
121 /// The checkpoint every path in the file is against.
122 pub checkpoint: Checkpoint,
123 pub canon: u16,
124 pub selection: Selection,
125}
126
127impl PackageHeader {
128 /// Read a package's first line, refusing another kind, another version or
129 /// an unknown member by name.
130 ///
131 /// # Errors
132 ///
133 /// When the line is not a package header this build reads.
134 pub fn parse(line: &str) -> Result<Self, String> {
135 let value: serde_json::Value =
136 serde_json::from_str(line).map_err(|e| format!("not JSON: {e}"))?;
137 let kind = value.get("kind").and_then(serde_json::Value::as_str);
138 if kind != Some(DISCLOSURE_KIND) {
139 return Err(format!(
140 "the line is a {kind:?}, not a {DISCLOSURE_KIND} header"
141 ));
142 }
143 let version = value.get("version").and_then(serde_json::Value::as_u64);
144 if version != Some(FORMAT_VERSION.into()) {
145 return Err(format!(
146 "the package is at format version {version:?}, and this build reads \
147 {FORMAT_VERSION}"
148 ));
149 }
150 serde_json::from_value(value).map_err(|e| format!("the package header is malformed: {e}"))
151 }
152}
153
154/// Why a reader that rebuilds or re-derives over a whole file refuses a
155/// package, in one sentence every such reader uses.
156const PACKAGE_REFUSED: &str = "the file is a disclosure package — the runs of one matter with \
157 their inclusion paths, not the plane's whole log — so nothing can be rebuilt or derived \
158 from it as if it were whole; verify it with `agentplane verify`";
159
160/// One journal record, as an export line.
161///
162/// Written out explicitly rather than by deriving `Serialize` on
163/// [`Record`](crate::journal::Record), and the reason is that this is a
164/// **durable format**. A derive makes the wire shape a side effect of the
165/// struct's field list, so adding a private field or renaming a public one
166/// silently changes what every downstream reader parses. Naming the four parts
167/// here means the format changes when somebody edits *this*, which is the only
168/// arrangement in which [`Header::version`] can mean anything.
169///
170/// The chain links travel with the body because an export without them is not
171/// checkable: `prev_hash` and `hash` are what let a reader re-walk the chain
172/// offline, which is the whole point of taking the record away.
173#[derive(Debug, Clone, PartialEq, serde::Serialize)]
174pub struct ExportedRecord<'r> {
175 pub seq: crate::core::Seq,
176 /// The typed view, for a reader's eyes. Verification never touches it —
177 /// see `raw` — and the verifier holds the two to each other so this cannot
178 /// quietly say something the hashed bytes do not.
179 ///
180 /// **Parsed from `raw`, never taken from the store's in-memory record.**
181 /// A sealed journal hands reads back *opened* — that is its job for the
182 /// runtime, whose own steps must read what they wrote — so an export that
183 /// copied the record's `body` field would write every sealed payload's
184 /// plaintext into a file, and destroying the key would no longer reach the
185 /// copy somebody exported last month. Deriving the display copy from the
186 /// hashed bytes makes body-matches-wire true by construction and keeps
187 /// sealed payloads sealed, which is the same rule the case layer's export
188 /// read states in prose.
189 pub body: DisplayBody,
190 pub prev_hash: &'r crate::core::Digest,
191 pub hash: &'r crate::core::Digest,
192 /// The plane's workload-key signature over this record's chain hash — who
193 /// wrote the record, not a hardware attestation of where. Present only
194 /// where the plane was configured to sign. `None` is an ordinary state
195 /// and is emitted as such rather than omitted, so a reader can tell
196 /// *unsigned* from *a field this export forgot*.
197 pub signature: Option<&'r crate::core::KeySignature>,
198 /// The exact bytes [`hash`](Self::hash) covers, verbatim.
199 ///
200 /// This is the wire-bytes rule, applied to the export: the chain is over
201 /// history **as written**, and a verifier that re-serialized the parsed
202 /// body was holding the file to *this build's* canonicalization rather
203 /// than to the bytes the store sealed. Canonical record bytes are UTF-8
204 /// JSON, so they travel as a string — escaped, exact, and recoverable
205 /// byte-for-byte.
206 pub raw: std::borrow::Cow<'r, str>,
207}
208
209/// A record line's display copy, which says what the hashed bytes say.
210///
211/// This build's typed view of bytes at the shape it writes, and the bytes' own
212/// JSON for a record at an older one; `verify` holds either to the bytes.
213#[derive(Debug, Clone, PartialEq, serde::Serialize)]
214#[serde(untagged)]
215pub enum DisplayBody {
216 Current(Box<crate::journal::RecordBody>),
217 Written(serde_json::Value),
218}
219
220impl<'r> ExportedRecord<'r> {
221 /// Build an export line from a stored record, deriving the display copy
222 /// from the wire bytes.
223 ///
224 /// Fallible on purpose, with no fallback to the record's opened `body`: a
225 /// record whose hashed bytes do not parse is corrupt, and substituting the
226 /// in-memory view would export exactly the plaintext this constructor
227 /// exists to keep out of the file — a silent fallback on the one value two
228 /// mechanisms must agree about.
229 fn from_stored(r: &'r crate::journal::Record) -> Result<Self, String> {
230 let unparsed =
231 |e: serde_json::Error| format!("record {}'s wire bytes do not parse: {e}", r.seq());
232 // The record was read through the store's upcaster, so its `body` is at
233 // this build's version; bytes at another are an older shape, and only
234 // their own JSON says what they say.
235 let body = match serde_json::from_slice::<crate::journal::RecordBody>(r.raw()) {
236 Ok(body) if body.v == r.body.v => DisplayBody::Current(Box::new(body)),
237 _ => DisplayBody::Written(serde_json::from_slice(r.raw()).map_err(unparsed)?),
238 };
239 Ok(Self {
240 seq: r.seq(),
241 body,
242 prev_hash: &r.prev_hash,
243 hash: &r.hash,
244 signature: r.signature.as_ref(),
245 raw: String::from_utf8_lossy(r.raw()),
246 })
247 }
248}
249
250/// A run's header line, emitted before its records.
251///
252/// It carries the one thing the record stream cannot: **where this run sits in
253/// the Merkle log**. That order is store state — a monotonic index assigned at
254/// seal time — and it appears in no record, so an export without it can be
255/// walked but cannot be checked against the checkpoint in its own header. The
256/// difference is between a transcript and evidence: a reader could confirm each
257/// chain links to itself and still not know whether a run had been dropped from
258/// the middle of the log.
259///
260/// `index` and `seal` are absent for a run that is still open. An unsealed run
261/// is not in the log and has no leaf, which is a state rather than a gap.
262#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
263pub struct RunBlock {
264 /// Always `"agentplane.export.run"`.
265 pub kind: &'static str,
266 pub run: RunId,
267 /// Position in the Merkle log, in seal order.
268 #[serde(skip_serializing_if = "Option::is_none")]
269 pub index: Option<u64>,
270 /// The leaf value: this run's terminal chain hash.
271 #[serde(skip_serializing_if = "Option::is_none")]
272 pub seal: Option<crate::core::Digest>,
273 /// In a disclosure package only: the sibling hashes, leaf-upwards, that
274 /// prove `seal` at `index` against the header's checkpoint. A whole export
275 /// never carries it, because it carries every leaf and rebuilds the root.
276 #[serde(skip_serializing_if = "Option::is_none")]
277 pub proof: Option<Vec<crate::core::Digest>>,
278}
279
280/// One case, as an export line — the case layer's whole account of one matter.
281///
282/// Emitted after the run blocks because the two halves answer different
283/// questions: the journal is *what happened*, the case is *what it happened
284/// to*. A restore of the journal alone rebuilds every index the journal owns
285/// and none of these rows, because case state is not derivable from records —
286/// which is exactly why the export has to carry it.
287///
288/// `state` travels **as stored**: sealed on a sealed plane. Exporting
289/// plaintext would quietly undo erasure — see the module docs, which make the
290/// same refusal for record payloads.
291///
292/// `blobs` carries digests, never bytes. Presence and integrity of the bytes
293/// are a question about a live blob store, which an offline file cannot
294/// answer and honestly reports as unchecked.
295///
296/// `hold` is always written, `null` for a matter nobody ordered preserved. A
297/// restore that brought a held matter back without its hold would hand the
298/// next retention pass a closed, old, unheld case to erase.
299#[derive(Debug, Clone, PartialEq, serde::Serialize)]
300pub struct CaseBlock {
301 /// Always `"agentplane.export.case"`.
302 pub kind: &'static str,
303 pub case: crate::core::Case,
304 pub deadlines: Vec<crate::core::Deadline>,
305 pub blobs: Vec<crate::core::Digest>,
306 /// The legal hold on this matter: instant, reason and operator.
307 pub hold: Option<crate::core::LegalHold>,
308}
309
310/// The last line of an export: what it contains, and what it does not.
311#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
312pub struct Trailer {
313 /// Always `"agentplane.export.end"`. Its **absence** is the signal that
314 /// matters: an export cut short by a crash, a full disk or a killed pipe
315 /// ends without one, so a reader can tell a prefix from a whole file
316 /// without comparing counts against a source it does not have.
317 pub kind: &'static str,
318 /// How many runs were asked for.
319 pub runs_requested: usize,
320 /// How many were read in full.
321 pub runs_exported: usize,
322 /// How many records were written.
323 pub records: usize,
324 /// How many cases the case layer contributed.
325 ///
326 /// Every case the case store holds is exported, so a record stamped with a
327 /// case this file does not carry is a finding the verifier makes.
328 pub cases: usize,
329 /// Runs that could not be read, and why.
330 ///
331 /// Named rather than counted. A count tells an auditor that something is
332 /// missing and not which case to go and ask about, and the run that fails
333 /// to read is not a random one.
334 pub unreadable: Vec<Unreadable>,
335}
336
337/// A run the export could not read.
338#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
339pub struct Unreadable {
340 pub run: RunId,
341 pub reason: String,
342}
343
344/// Write every record of `runs` as JSON Lines, framed by a header and trailer.
345///
346/// The writer is `std::io::Write` rather than a path so this composes with a
347/// file, a pipe, a socket or a buffer, and so the caller owns where the bytes
348/// land — an export function that chose the destination would be one an
349/// operator has to work around.
350///
351/// `cases` is the plane's case store, and every case it holds is written: a
352/// file whose records name a matter it does not carry is a finding to every
353/// reader of the format.
354///
355/// A run that cannot be read is recorded in the trailer and the export
356/// continues. Aborting instead would make one damaged run withhold every
357/// healthy one, which is the opposite of what an export is for; the trailer is
358/// what keeps that from being silent.
359///
360/// # Errors
361///
362/// Only for a failure to *write*. A failure to *read* a run is data — it lands
363/// in [`Trailer::unreadable`] — because the export still succeeded at the job
364/// it was given, and an auditor needs the part that survived.
365pub async fn to_jsonl<W: std::io::Write>(
366 store: &Arc<dyn JournalStore>,
367 cases: &Arc<dyn crate::case::CaseStore>,
368 runs: &[RunId],
369 out: W,
370) -> Result<Trailer, std::io::Error> {
371 write_file(store, cases, runs, None, out)
372 .await
373 .map(|written| written.trailer)
374}
375
376/// What a disclosure package carried, for the act that records it.
377#[derive(Debug, Clone, PartialEq, Eq)]
378pub struct Package {
379 /// The checkpoint every path in the file is against.
380 pub checkpoint: Checkpoint,
381 /// The runs the selection resolved to, in file order.
382 pub runs: Vec<RunId>,
383 /// The cases whose blocks the file carries.
384 pub cases: Vec<crate::core::CaseId>,
385 /// Whether any record carried a sealed payload.
386 pub sealed: bool,
387 pub trailer: Trailer,
388}
389
390/// Write a disclosure package: the runs `selection` names, each sealed one
391/// with its inclusion path against the header's checkpoint, and only the cases
392/// those runs belong to.
393///
394/// A case contributes the runs [`Case::runs`](crate::core::Case::runs) holds at
395/// the moment of the call. The case layer is the selection's cases plus every
396/// case a carried record is stamped with, so the coverage rule holds over the
397/// package as over a whole export.
398///
399/// This writes bytes and records nothing: a package that leaves the plane is a
400/// disclosure, and [`crate::disclosure::disclose`] is what records it before
401/// any byte reaches its destination.
402///
403/// # Errors
404///
405/// `NotFound` for a named case or run the plane does not hold — never a
406/// fall-back to the whole plane — an empty selection, a failure to write or
407/// to prove a sealed run at the header's size, or a run whose conclusion seals
408/// and which is still not in the log after the file was written against
409/// three successive checkpoints. The file is built in memory and written
410/// to `out` whole, so a refusal writes nothing.
411pub async fn package_to_jsonl<W: std::io::Write>(
412 store: &Arc<dyn JournalStore>,
413 cases: &Arc<dyn crate::case::CaseStore>,
414 selection: &Selection,
415 mut out: W,
416) -> Result<Package, std::io::Error> {
417 let runs = resolve(store, cases, selection).await?;
418 // A run concluded and sealed after the header's checkpoint was taken
419 // would travel open with its conclusion, which a reader of a package
420 // must take for a stripped leaf. So the file is written again against a
421 // later checkpoint until every sealing conclusion it carries is placed.
422 let mut attempts = 0;
423 let (bytes, written) = loop {
424 let mut bytes = Vec::new();
425 let written = write_file(store, cases, &runs, Some(selection), &mut bytes).await?;
426 attempts += 1;
427 match written.unplaced.first() {
428 None => break (bytes, written),
429 Some(run) if attempts >= PACKAGE_ATTEMPTS => {
430 return Err(std::io::Error::other(format!(
431 "run {run} concluded under an outcome that seals and is not in the log — \
432 the seal follows the conclusion, and recovery completes one a crash \
433 interrupted; no package was written, so retry once it is sealed"
434 )));
435 }
436 Some(_) => {}
437 }
438 };
439 out.write_all(&bytes)?;
440 out.flush()?;
441 Ok(Package {
442 checkpoint: written.checkpoint,
443 runs,
444 cases: written.cases,
445 sealed: written.sealed,
446 trailer: written.trailer,
447 })
448}
449
450/// How many checkpoints a package is written against before a sealing
451/// conclusion with no leaf is refused rather than raced.
452const PACKAGE_ATTEMPTS: usize = 3;
453
454/// The runs a selection names, each once, in the order named: a case's runs
455/// first, then the runs named alone.
456async fn resolve(
457 store: &Arc<dyn JournalStore>,
458 cases: &Arc<dyn crate::case::CaseStore>,
459 selection: &Selection,
460) -> Result<Vec<RunId>, std::io::Error> {
461 use std::io::{Error, ErrorKind};
462 if selection.cases.is_empty() && selection.runs.is_empty() {
463 return Err(Error::new(
464 ErrorKind::InvalidInput,
465 "a disclosure package names at least one case or run",
466 ));
467 }
468 let mut runs: Vec<RunId> = Vec::new();
469 for &id in &selection.cases {
470 let case = cases
471 .case(id)
472 .await
473 .map_err(|e| as_io(&e))?
474 .ok_or_else(|| {
475 Error::new(ErrorKind::NotFound, format!("no case {id} on this plane"))
476 })?;
477 for run in case.runs {
478 if !runs.contains(&run) {
479 runs.push(run);
480 }
481 }
482 }
483 for &run in &selection.runs {
484 if store.read(run, 1).await.map_err(|e| as_io(&e))?.is_empty() {
485 return Err(Error::new(
486 ErrorKind::NotFound,
487 format!("no run {run} on this plane"),
488 ));
489 }
490 if !runs.contains(&run) {
491 runs.push(run);
492 }
493 }
494 Ok(runs)
495}
496
497/// What one pass of the writer produced.
498struct Written {
499 checkpoint: Checkpoint,
500 cases: Vec<crate::core::CaseId>,
501 sealed: bool,
502 trailer: Trailer,
503 /// A package's runs that concluded under an outcome that seals and have
504 /// no leaf below the header's size.
505 unplaced: Vec<RunId>,
506}
507
508/// The one writer behind a whole export and a package; `package` is the
509/// selection for a package and `None` for a whole export.
510#[allow(clippy::too_many_lines)]
511async fn write_file<W: std::io::Write>(
512 store: &Arc<dyn JournalStore>,
513 cases: &Arc<dyn crate::case::CaseStore>,
514 runs: &[RunId],
515 package: Option<&Selection>,
516 mut out: W,
517) -> Result<Written, std::io::Error> {
518 let checkpoint = store.checkpoint().await.map_err(|e| as_io(&e))?;
519 let taken = checkpoint.clone();
520 // Held out of the header for the one comparison below: the header owns the
521 // checkpoint from here on, and the log size is the half of it every run
522 // block is checked against.
523 let log_size = checkpoint.size;
524 // Read from the build, never from the caller. The rule that computed the
525 // digests is a fact about the store's own writes, and a parameter here was
526 // a header any embedder could make lie — every caller passed
527 // `canon::VERSION` verbatim, which is what a fact looks like when it is
528 // asked for as an argument.
529 match package {
530 None => {
531 let header = Header {
532 kind: "agentplane.export",
533 version: FORMAT_VERSION,
534 checkpoint,
535 canon: crate::core::canon::VERSION,
536 };
537 writeln!(out, "{}", to_line(&header)?)?;
538 }
539 Some(selection) => {
540 let header = PackageHeader {
541 kind: DISCLOSURE_KIND.to_owned(),
542 version: FORMAT_VERSION,
543 checkpoint,
544 canon: crate::core::canon::VERSION,
545 selection: selection.clone(),
546 };
547 writeln!(out, "{}", to_line(&header)?)?;
548 }
549 }
550 let mut sealed = false;
551 let mut stamped: std::collections::BTreeSet<crate::core::CaseId> = package
552 .map(|s| s.cases.iter().copied().collect())
553 .unwrap_or_default();
554
555 let mut records = 0usize;
556 let mut exported = 0usize;
557 let mut unreadable = Vec::new();
558 let mut unplaced = Vec::new();
559
560 // Asked before the records so each block heads them, and asked at all
561 // because the log position is the half of the evidence the records do not
562 // carry. A store that cannot answer leaves the runs unsealed rather than
563 // failing the export: the position is missing, and the verifier says so,
564 // which is better than no export.
565 let positions = store
566 .log_positions(runs)
567 .await
568 .unwrap_or_else(|_| vec![None; runs.len()]);
569 for (&run, placed) in runs.iter().zip(positions) {
570 // A run sealed *after* the header's checkpoint was taken is not in that
571 // checkpoint. Stamping its position anyway would make the export
572 // disagree with its own first line: the verifier rebuilds a tree one
573 // leaf larger than the root it compares against, and reports tampering
574 // where there was only time. Such a run is exported as still open —
575 // true relative to the moment this export describes — and the next
576 // export carries it sealed.
577 let placed = placed.filter(|&(index, _)| index < log_size);
578 // A package proves each leaf on its own, against the header's size: a
579 // path against the live log would fail against the header's root for
580 // every run sealed while the file is written.
581 let proof = match (package, placed) {
582 (Some(_), Some(_)) => store
583 .inclusion_proof_at(run, log_size)
584 .await
585 .map_err(|e| as_io(&e))?
586 .map(|inclusion| inclusion.proof),
587 _ => None,
588 };
589 writeln!(
590 out,
591 "{}",
592 to_line(&RunBlock {
593 kind: "agentplane.export.run",
594 run,
595 index: placed.map(|(index, _)| index),
596 seal: placed.map(|(_, seal)| seal),
597 proof,
598 })?
599 )?;
600
601 match store.read(run, 1).await {
602 // A run the store holds nothing for is filed as unreadable, not
603 // exported as an empty block. Both backends answer an unknown run
604 // with an empty read rather than an error, so without this arm a
605 // mistyped run id produced a block with no records under it — a
606 // shape the verifier must otherwise treat as records removed after
607 // the fact. Naming it here keeps the trailer's accounting honest:
608 // an empty block in a file whose trailer does not declare the run
609 // unreadable is tampering, and only because no honest writer
610 // produces one.
611 Ok(found) if found.is_empty() => unreadable.push(Unreadable {
612 run,
613 reason: "the store holds no records for this run".to_owned(),
614 }),
615 Ok(found) => {
616 // Every line is derived from its wire bytes before any is
617 // written, so a record that cannot be derived files the whole
618 // run as unreadable instead of leaving a half-written block
619 // shaped like a complete one.
620 match found
621 .iter()
622 .map(ExportedRecord::from_stored)
623 .collect::<Result<Vec<_>, _>>()
624 {
625 Ok(lines) => {
626 for line in &lines {
627 writeln!(out, "{}", to_line(line)?)?;
628 records += 1;
629 }
630 exported += 1;
631 if package.is_some() {
632 if placed.is_none() && crate::audit::has_sealing_conclusion(&found) {
633 unplaced.push(run);
634 }
635 stamped.extend(found.iter().filter_map(|r| r.body.case));
636 sealed |= found.iter().any(|r| carries_sealed(r.raw()));
637 }
638 }
639 Err(reason) => unreadable.push(Unreadable { run, reason }),
640 }
641 }
642 Err(e) => unreadable.push(Unreadable {
643 run,
644 reason: e.to_string(),
645 }),
646 }
647 }
648
649 // The case layer, after the runs and before the trailer. Every case, not
650 // the cases these runs touch: a case is the unit an erasure request or a
651 // regulator names, and a subset chosen by run membership would silently
652 // drop the matter whose runs happened not to be asked for. A package is
653 // the one file that is a subset on purpose, and says so in its header, so
654 // it carries the cases it was asked for and every case its records name.
655 let mut written_cases = Vec::new();
656 if package.is_some() {
657 for id in stamped {
658 if let Some(case) = cases.case(id).await.map_err(|e| as_io(&e))? {
659 write_case(cases, case, &mut out).await?;
660 written_cases.push(id);
661 }
662 }
663 } else {
664 let mut after: Option<crate::core::CaseId> = None;
665 loop {
666 let page = cases.cases(after, CASE_PAGE).await.map_err(|e| as_io(&e))?;
667 let Some(last) = page.last() else { break };
668 after = Some(last.id);
669 let full = page.len() >= CASE_PAGE;
670 for case in page {
671 written_cases.push(case.id);
672 write_case(cases, case, &mut out).await?;
673 }
674 if !full {
675 break;
676 }
677 }
678 }
679 let case_count = written_cases.len();
680
681 let trailer = Trailer {
682 kind: "agentplane.export.end",
683 runs_requested: runs.len(),
684 runs_exported: exported,
685 records,
686 cases: case_count,
687 unreadable,
688 };
689 writeln!(out, "{}", to_line(&trailer)?)?;
690 out.flush()?;
691 Ok(Written {
692 checkpoint: taken,
693 cases: written_cases,
694 sealed,
695 trailer,
696 unplaced,
697 })
698}
699
700/// Whether a record's wire bytes carry a sealed payload anywhere in them.
701fn carries_sealed(raw: &[u8]) -> bool {
702 fn walk(value: &serde_json::Value) -> bool {
703 match value {
704 serde_json::Value::String(text) => crate::journal::payload::is_sealed_text(text),
705 serde_json::Value::Array(items) => items.iter().any(walk),
706 serde_json::Value::Object(members) => {
707 crate::journal::payload::is_sealed(value) || members.values().any(walk)
708 }
709 _ => false,
710 }
711 }
712 serde_json::from_slice::<serde_json::Value>(raw).is_ok_and(|value| walk(&value))
713}
714
715/// One case block: the case, its deadlines, its blob digests and its hold.
716async fn write_case<W: std::io::Write>(
717 cases: &Arc<dyn crate::case::CaseStore>,
718 case: crate::core::Case,
719 out: &mut W,
720) -> Result<(), std::io::Error> {
721 let deadlines = cases.deadlines(case.id).await.map_err(|e| as_io(&e))?;
722 let blobs = cases.blobs_of(case.id).await.map_err(|e| as_io(&e))?;
723 let hold = cases.hold(case.id).await.map_err(|e| as_io(&e))?;
724 writeln!(
725 out,
726 "{}",
727 to_line(&CaseBlock {
728 kind: "agentplane.export.case",
729 case,
730 deadlines,
731 blobs,
732 hold,
733 })?
734 )
735}
736
737/// How many cases one enumeration page holds — shared with the live drill
738/// ([`crate::drill`]), which walks the same case layer with the same paging.
739/// Interior to the crate either way: the stream out is unbounded, and the
740/// page only bounds memory. One constant, because two walks that paged
741/// differently would be two subtly different definitions of "every case".
742pub(crate) const CASE_PAGE: usize = 256;
743
744/// One value as one line, refusing to write a line that is not valid JSON.
745fn to_line<T: serde::Serialize>(value: &T) -> Result<String, std::io::Error> {
746 serde_json::to_string(value).map_err(|e| std::io::Error::other(e.to_string()))
747}
748
749fn as_io(e: &StoreError) -> std::io::Error {
750 std::io::Error::other(e.to_string())
751}
752
753// ── Reading one back ────────────────────────────────────────────────────────
754
755/// What a verification pass concluded, and what it could not look at.
756///
757/// The same shape as [`AuditReport`](crate::audit::AuditReport) and for the same
758/// reason: a pass that reports only failures tells you about its coverage by
759/// omission. An export verified without a public key has not established
760/// authorship, and saying so is the difference between *this is sound* and
761/// *nothing I checked was wrong*.
762#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
763pub struct VerifyReport {
764 /// The checkpoint the export claims to be a copy of.
765 pub checkpoint: Checkpoint,
766 /// Runs whose chain recomputed exactly.
767 pub sound: Vec<RunId>,
768 /// What went wrong, in the order found.
769 pub findings: Vec<String>,
770 /// Checks that were not performed, and why.
771 pub not_checked: Vec<String>,
772 /// How many records were read.
773 pub records: usize,
774 /// How many case blocks were read.
775 pub cases: usize,
776 /// Whether the file ended with its trailer.
777 ///
778 /// A truncated export is otherwise a valid prefix: every line parses, every
779 /// chain link joins, and the only thing wrong is what is missing.
780 pub complete: bool,
781 /// Why this reader could check nothing past the header, when it could not.
782 ///
783 /// Set for a `canon` this build does not implement: the rule names the
784 /// digest algorithm too, so every hash in the file is one this reader
785 /// cannot recompute. Neither a finding nor corruption — another build may
786 /// verify the file.
787 pub unverifiable: Option<String>,
788 /// What the file selected, when it is a disclosure package: the runs of
789 /// one matter, each proved by its own path, and nothing about the log's
790 /// other leaves. `None` for a whole export.
791 pub selection: Option<Selection>,
792}
793
794impl VerifyReport {
795 /// Whether every check that ran, passed, and the checks could run. See
796 /// [`Self::not_checked`] and [`Self::unverifiable`].
797 #[must_use]
798 pub fn is_sound(&self) -> bool {
799 self.findings.is_empty() && self.complete && self.unverifiable.is_none()
800 }
801}
802
803/// Recompute an export from its own bytes, and check it against its checkpoint.
804///
805/// **This is the restore drill, and it is the half that makes a restore worth
806/// having.** Putting records back into a store proves that bytes moved; it does
807/// not prove they are the bytes that were taken, in the order they were taken,
808/// with nothing dropped from the middle. That is what this establishes, and it
809/// establishes it without the runtime that wrote the data and without the store
810/// it came from — an export and this function are the whole dependency.
811///
812/// Four properties, each checkable only because the export was designed to
813/// carry the evidence for it:
814///
815/// * **Every record's hash is recomputed**, from its body and its predecessor's
816/// hash, using the canonicalization rule the header names. A record whose
817/// stored hash disagrees was edited after sealing. This is not a comparison of
818/// the file against itself: `Record::seal` is the same function the store
819/// sealed through, so agreement means the bytes are the ones that were
820/// written.
821/// * **Chains join, and sequences are contiguous.** A removed record breaks a
822/// link; a removed *tail* does not, which is why the sequence is checked too.
823/// * **The Merkle root is rebuilt** from the per-run log positions and compared
824/// with `expected` — the checkpoint the reader was given by somebody other
825/// than whoever wrote this file. This is the one that catches a whole run
826/// dropped from the middle of the export: the per-run chains all still
827/// verify, and only the tree notices.
828///
829/// Without `expected` it can only be compared with the file's **own header**,
830/// which reads like the same check and is not: an editor who drops a run and
831/// rewrites the header's size and root produces a file that agrees with
832/// itself perfectly. This crate makes the same argument one level down — a
833/// record's `prev_hash` is checked by rehashing the wire bytes, never against
834/// the previous line, because "the file agrees with itself" is what a
835/// competent editor achieves. So a pass with no `expected` reports the root
836/// check under [`VerifyReport::not_checked`]: internal consistency
837/// established, deletion not.
838/// * **The file is framed.** A missing trailer means the export was cut short,
839/// and every line before the cut is still perfectly valid.
840/// * **The trailer's own accounting holds.** Its run and record counts are
841/// compared against what was actually read, and a run it declares unreadable
842/// is reported as *unchecked* rather than as tampering — the writer said at
843/// export time that the run's records are not here, which is the opposite of
844/// hiding it. An empty run block the trailer does **not** declare unreadable
845/// is the tamper case: no honest writer produces one.
846///
847/// Signatures are checked when a verifier is supplied and reported as unchecked
848/// when not.
849///
850/// # Errors
851///
852/// Only for a failure to read the input. A malformed or dishonest export is a
853/// *finding*, not an error — the whole point is to produce a report about it.
854pub fn verify<R: std::io::BufRead>(
855 input: R,
856 verifier: Option<&dyn crate::core::Verifier>,
857 anchors: &[crate::journal::Anchor],
858) -> Result<VerifyReport, std::io::Error> {
859 let upcaster = crate::journal::current_upcaster();
860 verify_with(input, verifier, anchors, upcaster.as_ref())
861}
862
863/// [`verify`], reading each record through `upcaster` rather than the one this
864/// build ships.
865///
866/// The hash is held to the bytes as written either way; the upcaster decides
867/// only which record versions this reader can read, and a version it cannot
868/// reach is a build skew, never an edit.
869///
870/// # Errors
871///
872/// As [`verify`].
873pub fn verify_with<R: std::io::BufRead>(
874 input: R,
875 verifier: Option<&dyn crate::core::Verifier>,
876 anchors: &[crate::journal::Anchor],
877 upcaster: &dyn crate::journal::Upcaster,
878) -> Result<VerifyReport, std::io::Error> {
879 verify_observed(input, verifier, anchors, upcaster, &mut |_| {})
880}
881
882/// One run block as the verification pass closed it.
883pub(crate) struct ClosedRun<'a> {
884 pub(crate) run: RunId,
885 /// Whether the block declared a log position and a seal.
886 pub(crate) sealed: bool,
887 /// The records as the pass recomputed them. Whether they are sound is
888 /// [`VerifyReport::sound`] once the pass returns.
889 pub(crate) records: &'a [crate::journal::Record],
890}
891
892/// [`verify`], handing each run block to `observe` as it closes.
893#[allow(clippy::too_many_lines)]
894pub(crate) fn verify_observed<R: std::io::BufRead>(
895 input: R,
896 verifier: Option<&dyn crate::core::Verifier>,
897 anchors: &[crate::journal::Anchor],
898 upcaster: &dyn crate::journal::Upcaster,
899 observe: &mut dyn FnMut(ClosedRun<'_>),
900) -> Result<VerifyReport, std::io::Error> {
901 use crate::core::Digest;
902 use serde_json::Value;
903
904 let mut report = VerifyReport {
905 checkpoint: Checkpoint {
906 origin: String::new(),
907 size: 0,
908 root: Digest::ZERO,
909 },
910 sound: Vec::new(),
911 findings: Vec::new(),
912 not_checked: Vec::new(),
913 records: 0,
914 cases: 0,
915 complete: false,
916 unverifiable: None,
917 selection: None,
918 };
919 unanswerable(&mut report, verifier.is_some());
920
921 let mut header_seen = false;
922 // (index, leaf) for every sealed run, so the tree can be rebuilt in log
923 // order rather than in the order the export happened to walk.
924 let mut leaves: Vec<(u64, crate::core::merkle::LeafHash)> = Vec::new();
925 let mut pass: Option<RunPass> = None;
926 // The reader's own tally, held against the trailer's at the end: every run
927 // block seen, every block that carried at least one record, and every block
928 // that carried none. The trailer adjudicates the empty ones — an export
929 // that declared the run unreadable was honest about it, and one that did
930 // not has had records removed — which is why they are collected rather
931 // than judged on the spot: the trailer is the last line, and an
932 // intermediate block closes before it is read.
933 let mut run_blocks = 0usize;
934 let mut read_runs = 0usize;
935 let mut empty_blocks: Vec<RunId> = Vec::new();
936 let mut claims = TrailerClaims::default();
937 // The two halves of the case cross-check: what the records name, and what
938 // the case layer carries. Settled at the end, because either side can
939 // arrive first in the file.
940 let mut stamped: std::collections::BTreeSet<crate::core::CaseId> =
941 std::collections::BTreeSet::new();
942 let mut carried: std::collections::BTreeSet<crate::core::CaseId> =
943 std::collections::BTreeSet::new();
944 let mut blob_digests = 0usize;
945 // Every run a block names and every position a block claims, so a run
946 // carried twice or a position claimed twice is named, whatever the mode.
947 let mut blocks_of: std::collections::BTreeSet<RunId> = std::collections::BTreeSet::new();
948 let mut positions: std::collections::BTreeSet<u64> = std::collections::BTreeSet::new();
949
950 // Only the first non-empty line may be a header: it alone fixes which
951 // checkpoint every path is held to and which rules the file is read under.
952 let mut first = true;
953 for line in input.lines() {
954 let line = line?;
955 if line.trim().is_empty() {
956 continue;
957 }
958 let is_first = std::mem::replace(&mut first, false);
959 let Ok(value) = serde_json::from_str::<Value>(&line) else {
960 report
961 .findings
962 .push("a line is not valid JSON, so the export is unreadable from there on".into());
963 break;
964 };
965 if let Some(kind) = value.get("kind").and_then(Value::as_str) {
966 note_unknown_members(kind, report.selection.is_some(), &value, &mut report);
967 }
968 match value.get("kind").and_then(Value::as_str) {
969 // A package is read by the same pass, with its own settlement: each
970 // leaf is proved by its path, and nothing counts the log.
971 Some(kind @ ("agentplane.export" | DISCLOSURE_KIND)) => {
972 if !is_first {
973 report.findings.push(format!(
974 "a '{kind}' header appears past the first line and was ignored — the \
975 first line alone names the checkpoint and the rules this file is \
976 read under, and a later one would re-choose them for every run \
977 after it"
978 ));
979 continue;
980 }
981 header_seen = true;
982 read_header(&value, &mut report);
983 if report.unverifiable.is_some() {
984 return Ok(report);
985 }
986 if kind == DISCLOSURE_KIND {
987 read_selection(&value, &mut report);
988 }
989 }
990 Some("agentplane.export.case") => {
991 read_case_block(&value, &mut report, &mut carried, &mut blob_digests);
992 }
993 Some("agentplane.export.run") => {
994 run_blocks += 1;
995 finish_run(
996 &mut report,
997 pass.take(),
998 verifier,
999 &mut read_runs,
1000 &mut empty_blocks,
1001 observe,
1002 );
1003 pass = open_run_block(&value, &mut leaves, &mut report);
1004 if let Some(opened) = pass.as_mut() {
1005 claim_once(opened, &mut blocks_of, &mut positions, &mut report);
1006 }
1007 }
1008 Some("agentplane.export.end") => read_trailer(&value, &mut report, &mut claims),
1009 _ => {
1010 report.records += 1;
1011 let Some(pass) = pass.as_mut() else {
1012 report.findings.push(
1013 "a record appears before any run block, so nothing says which run it \
1014 belongs to"
1015 .into(),
1016 );
1017 continue;
1018 };
1019 read_record(&value, pass, &mut report, &mut stamped, upcaster);
1020 }
1021 }
1022 }
1023 finish_run(
1024 &mut report,
1025 pass,
1026 verifier,
1027 &mut read_runs,
1028 &mut empty_blocks,
1029 observe,
1030 );
1031
1032 // Open runs are the blocks that contributed no leaf. Computed here because
1033 // this is the one place that holds both numbers, and passed on rather than
1034 // recounted.
1035 let open_runs = run_blocks.saturating_sub(leaves.len());
1036 settle(&mut report, header_seen, leaves, anchors);
1037 settle_trailer(
1038 &mut report,
1039 &claims,
1040 run_blocks,
1041 read_runs,
1042 open_runs,
1043 &empty_blocks,
1044 );
1045 settle_cases(&mut report, &stamped, &carried, blob_digests);
1046 Ok(report)
1047}
1048
1049/// What the trailer claims about the file, held for the settlement.
1050///
1051/// Collected rather than compared on the spot, for two reasons that are the
1052/// same reason: the trailer is the last line, so the totals it must be held
1053/// against only exist once the whole file has been read — and the per-run
1054/// verdicts it adjudicates (is an empty block an honestly-declared unreadable
1055/// run, or records removed after the fact?) close *before* it is read, because
1056/// each run block is finished when the next one starts.
1057#[derive(Default)]
1058struct TrailerClaims {
1059 runs_requested: Option<u64>,
1060 runs_exported: Option<u64>,
1061 records: Option<u64>,
1062 /// Runs the export itself declared unreadable, with the writer's reason.
1063 unreadable: Vec<(RunId, String)>,
1064}
1065
1066/// Read the trailer: the file is complete, and its case count holds.
1067///
1068/// The case-count comparison is what catches the case layer stripped *whole*:
1069/// with every block gone the coverage cross-check has nothing to compare, and
1070/// the file would read as an export of a plane that simply had no cases —
1071/// while its trailer still says otherwise. The run and record counts are
1072/// collected here and compared in [`settle_trailer`], where the totals exist.
1073fn read_trailer(value: &serde_json::Value, report: &mut VerifyReport, claims: &mut TrailerClaims) {
1074 report.complete = true;
1075 if let Some(declared) = value.get("cases").and_then(serde_json::Value::as_u64)
1076 && declared != report.cases as u64
1077 {
1078 report.findings.push(format!(
1079 "the trailer says {declared} case(s) were exported and this file \
1080 carries {} — the case layer was cut after the export was taken",
1081 report.cases
1082 ));
1083 }
1084 claims.runs_requested = value
1085 .get("runs_requested")
1086 .and_then(serde_json::Value::as_u64);
1087 claims.runs_exported = value
1088 .get("runs_exported")
1089 .and_then(serde_json::Value::as_u64);
1090 claims.records = value.get("records").and_then(serde_json::Value::as_u64);
1091 if let Some(list) = value
1092 .get("unreadable")
1093 .and_then(serde_json::Value::as_array)
1094 {
1095 for entry in list {
1096 let Some(run) = entry
1097 .get("run")
1098 .and_then(serde_json::Value::as_str)
1099 .and_then(|s| RunId::parse(s).ok())
1100 else {
1101 continue;
1102 };
1103 let reason = entry
1104 .get("reason")
1105 .and_then(serde_json::Value::as_str)
1106 .unwrap_or("no reason recorded")
1107 .to_owned();
1108 claims.unreadable.push((run, reason));
1109 }
1110 }
1111}
1112
1113/// Hold the trailer's own accounting to what was actually read.
1114///
1115/// Before this settlement existed, only the trailer's `cases` count was ever
1116/// consulted — `runs_requested`, `runs_exported`, `records` and `unreadable`
1117/// were fields the writer stamped and no reader read, so deleting an open
1118/// run's tail records while keeping the trailer verified clean: an open run
1119/// has no leaf to pin its tail, a chain prefix verifies, and the only witness
1120/// left is the count.
1121///
1122/// The empty blocks are adjudicated here too, and the trailer is what decides
1123/// which way each one goes. A run the export *declares* unreadable is
1124/// unchecked, not tampering: the writer said at export time that this run's
1125/// records are not in the file, which is the opposite of hiding it, and
1126/// reporting it as "records removed after sealing" would teach an operator
1127/// that findings are noise. An empty block the trailer does **not** declare is
1128/// the tamper case — no honest writer produces one, because an unreadable or
1129/// empty read files the run in the trailer instead.
1130///
1131/// What this does NOT cover: a trailer rewritten to match an edited file. The
1132/// counts are the file's claim about itself, and holding a file to itself
1133/// never catches an editor who updates both halves — that is the chain, leaf
1134/// and Merkle-root checks' job, which tie the surviving bytes to history. Nor
1135/// does it cover an open run's tail cut *before* the export was taken: the
1136/// store served the shortened history, the writer counted what it served, and
1137/// no offline file can see past its own writer.
1138fn settle_trailer(
1139 report: &mut VerifyReport,
1140 claims: &TrailerClaims,
1141 run_blocks: usize,
1142 read_runs: usize,
1143 open_runs: usize,
1144 empty_blocks: &[RunId],
1145) {
1146 // Said here because `audit` says it about a live store, and these two
1147 // answer one question about one history: an open run has no Merkle leaf,
1148 // so nothing pins its tail and a truncation is undetectable until the run
1149 // seals. The offline reader is the one an independent auditor holds, so it
1150 // is the worse of the two to leave silent.
1151 //
1152 // Once per file rather than per run. A reader deciding what a clean report
1153 // is worth needs the count and the reason; a line per run buries both.
1154 if open_runs > 0 {
1155 report.not_checked.push(format!(
1156 "{open_runs} open run(s): a run that has not concluded has no position in the \
1157 Merkle log, so the root proves nothing about it — its chain and signatures \
1158 were verified, and records cut from its tail before the export was taken are \
1159 undetectable from this file"
1160 ));
1161 }
1162 for (run, reason) in &claims.unreadable {
1163 report.not_checked.push(format!(
1164 "run {run}: the export declares it unreadable ({reason}), so its records are \
1165 not in this file and nothing about it was verified"
1166 ));
1167 }
1168 for run in empty_blocks {
1169 if claims.unreadable.iter().any(|(u, _)| u == run) {
1170 continue;
1171 }
1172 report.findings.push(format!(
1173 "run {run}: its block carries no records and the export does not declare it \
1174 unreadable — either the records were removed after the export was taken, or \
1175 the file was cut short before them"
1176 ));
1177 }
1178 // The counts exist only on a framed file; a missing trailer is already the
1179 // truncation finding in `settle`, and comparing against nothing would
1180 // manufacture a second finding about the same cut.
1181 if !report.complete {
1182 return;
1183 }
1184 match (claims.runs_requested, claims.runs_exported, claims.records) {
1185 (Some(requested), Some(exported), Some(records)) => {
1186 if requested != run_blocks as u64 {
1187 report.findings.push(format!(
1188 "the trailer says {requested} run(s) were requested and this file carries \
1189 {run_blocks} run block(s) — whole runs were removed or added after the \
1190 export was taken"
1191 ));
1192 }
1193 if exported != read_runs as u64 {
1194 report.findings.push(format!(
1195 "the trailer says {exported} run(s) were exported in full and this file \
1196 carries records for {read_runs} — a run's records were removed after the \
1197 export was taken"
1198 ));
1199 }
1200 if records != report.records as u64 {
1201 report.findings.push(format!(
1202 "the trailer says {records} record(s) were written and this file carries \
1203 {} — record lines were removed or added after the export was taken",
1204 report.records
1205 ));
1206 }
1207 }
1208 _ => report.findings.push(
1209 "the trailer is missing counts this format always writes (runs_requested, \
1210 runs_exported, records) — a reader cannot hold the file to its own accounting"
1211 .to_owned(),
1212 ),
1213 }
1214}
1215
1216/// Read one case block: count it, collect its id for the coverage settlement,
1217/// and flag the malformations a reader would otherwise trip over silently.
1218fn read_case_block(
1219 value: &serde_json::Value,
1220 report: &mut VerifyReport,
1221 carried: &mut std::collections::BTreeSet<crate::core::CaseId>,
1222 blob_digests: &mut usize,
1223) {
1224 use serde_json::Value;
1225
1226 report.cases += 1;
1227 match serde_json::from_value::<crate::core::Case>(
1228 value.get("case").cloned().unwrap_or(Value::Null),
1229 ) {
1230 Ok(case) => {
1231 carried.insert(case.id);
1232 }
1233 Err(e) => report
1234 .findings
1235 .push(format!("a case block is malformed: {e}")),
1236 }
1237 if value
1238 .get("deadlines")
1239 .is_none_or(|d| serde_json::from_value::<Vec<crate::core::Deadline>>(d.clone()).is_err())
1240 {
1241 report
1242 .findings
1243 .push("a case block's deadlines are malformed".to_owned());
1244 }
1245 if let Err(e) = case_hold(value) {
1246 report.findings.push(e);
1247 }
1248 *blob_digests += value
1249 .get("blobs")
1250 .and_then(Value::as_array)
1251 .map_or(0, Vec::len);
1252}
1253
1254/// What this pass cannot answer, whatever the file turns out to contain.
1255///
1256/// The twin of the audit's own missing-evidence list, and it exists for the same
1257/// reason: a pass that quietly skips a question and then reports clean is the
1258/// reassuring-but-empty artifact this module is built to avoid. Two kinds sit
1259/// here — a check this invocation lacked an input for, and state the format
1260/// never carries at all.
1261fn unanswerable(report: &mut VerifyReport, verifier_supplied: bool) {
1262 if !verifier_supplied {
1263 report.not_checked.push(
1264 "signatures — no public key was supplied, so this pass cannot say who wrote \
1265 anything"
1266 .to_owned(),
1267 );
1268 }
1269 report
1270 .not_checked
1271 .extend(UNCARRIED.iter().map(|limit| (*limit).to_owned()));
1272}
1273
1274/// Operational state this format never carries, named on every pass.
1275///
1276/// A restored plane's runs and cases come back; the rows beside them do not.
1277/// Said unconditionally rather than where something references them, because
1278/// these are properties of the *format* — a reader meeting a minimal artifact is
1279/// the one most likely to assume a clean report means a total restore.
1280///
1281/// Cursors and unconsumed decisions degrade safely, which is the argument for
1282/// leaving them out: a cursor costs repetition against receivers that already
1283/// deduplicate, and a decision is taken again under the same four-eyes and
1284/// expiry. The disclosure register is left out for its recipients' sake, and
1285/// its loss is that a restored plane's erasures name no earlier copy. Stating
1286/// each argument is what stops it from being a silence.
1287const UNCARRIED: [&str; 3] = [
1288 "webhook delivery cursors — this file carries no push registrations, so a restored plane \
1289 re-delivers from the start of each subscriber's history rather than from where it got \
1290 to. Receivers deduplicate on the event's own identity, so the cost is repetition rather \
1291 than loss",
1292 "worklist decisions no run has consumed — a decision recorded against a task and not yet \
1293 read back by the run it answers is a store row, not a record, so it does not survive \
1294 here. The task re-opens and is decided again under the same four-eyes and expiry",
1295 "the disclosure register — which matters left the plane, and to whom, is an operator \
1296 row, not a record, so a restored plane names no earlier disclosure in its erasures. \
1297 Carrying it would tell every recipient of an export who else received what",
1298];
1299
1300/// The case layer's own settlement: coverage, and what a file cannot check.
1301///
1302/// The coverage rule has a deliberate asymmetry. A record stamped with a case
1303/// the file does not carry is a **finding** — this plane had a case layer (the
1304/// stamp proves it) and the export is missing a matter the journal names,
1305/// whether one case block is missing or all of them. The reverse is not: a
1306/// case whose runs are absent is the ordinary result of exporting a subset of
1307/// runs, and every case travels regardless of which runs were asked for.
1308fn settle_cases(
1309 report: &mut VerifyReport,
1310 stamped: &std::collections::BTreeSet<crate::core::CaseId>,
1311 carried: &std::collections::BTreeSet<crate::core::CaseId>,
1312 blob_digests: usize,
1313) {
1314 for case in stamped.difference(carried) {
1315 report.findings.push(format!(
1316 "case {case} is stamped on exported records and missing from the case layer — \
1317 the journal names a matter this file does not carry"
1318 ));
1319 }
1320 // A package proves the inclusion of what it carries, never completeness:
1321 // a matter its own header names and its case layer omits is the one
1322 // omission the file itself can show.
1323 let selected: Vec<crate::core::CaseId> = report
1324 .selection
1325 .as_ref()
1326 .map(|s| s.cases.clone())
1327 .unwrap_or_default();
1328 for case in selected {
1329 if !carried.contains(&case) {
1330 report.findings.push(format!(
1331 "case {case} is named by the package's selection and missing from its case \
1332 layer — the package does not carry a matter it says it discloses"
1333 ));
1334 }
1335 }
1336 if carried.is_empty() {
1337 return;
1338 }
1339 if blob_digests > 0 {
1340 report.not_checked.push(format!(
1341 "blob bytes — the case layer references {blob_digests} blob digest(s) and this \
1342 file carries digests, not bytes; presence and integrity are a question about a \
1343 live blob store"
1344 ));
1345 }
1346 report.not_checked.push(
1347 "sealed-state keys — whether sealed case state can still be opened is a question \
1348 about a live key ring, which an offline file cannot answer"
1349 .to_owned(),
1350 );
1351}
1352
1353/// Open a run block: fresh per-run state, and the block's leaf collected for
1354/// the tree rebuild. Returns `None` for a block whose run id does not parse —
1355/// the records under it are then flagged as belonging to no run, which is the
1356/// honest reading of a block nothing can be looked up by.
1357///
1358/// A block is placed only when it carries both an integer `index` and a
1359/// `seal` that parses; one carrying either alone, or a seal that does not
1360/// parse, claims a position nothing can check, and is a finding rather than
1361/// an open run.
1362fn open_run_block(
1363 value: &serde_json::Value,
1364 leaves: &mut Vec<(u64, crate::core::merkle::LeafHash)>,
1365 report: &mut VerifyReport,
1366) -> Option<RunPass> {
1367 use crate::core::{Digest, merkle};
1368 use serde_json::Value;
1369
1370 let run = value
1371 .get("run")
1372 .and_then(Value::as_str)
1373 .and_then(|s| RunId::parse(s).ok())?;
1374 let index = value.get("index").and_then(Value::as_u64);
1375 let declared_seal = value
1376 .get("seal")
1377 .and_then(|s| serde_json::from_value::<Digest>(s.clone()).ok());
1378 let placed = index.zip(declared_seal);
1379 let claims_place = value.get("index").is_some() || value.get("seal").is_some();
1380 let malformed = claims_place && placed.is_none();
1381 if malformed {
1382 report.findings.push(format!(
1383 "run {run}: its block claims a log position without both an integer index and a \
1384 seal that parses, so nothing can place it"
1385 ));
1386 }
1387 if let Some((index, seal)) = placed {
1388 leaves.push((index, merkle::leaf_hash(&seal)));
1389 }
1390 Some(RunPass {
1391 run,
1392 declared_seal: placed.map(|(_, seal)| seal),
1393 sealed: placed.is_some(),
1394 index,
1395 proof: value
1396 .get("proof")
1397 .and_then(|p| serde_json::from_value::<Vec<Digest>>(p.clone()).ok()),
1398 prev: Digest::ZERO,
1399 last_seq: 0,
1400 records: 0,
1401 resealed: Vec::new(),
1402 clean: !malformed,
1403 })
1404}
1405
1406/// Hold a run block to the ones before it: a run is carried once and a log
1407/// position is claimed by one run. Either repeated leaves the block unsound,
1408/// and a run carried twice is withdrawn from the sound list its first block
1409/// earned, since neither block is the run's one history.
1410fn claim_once(
1411 pass: &mut RunPass,
1412 blocks_of: &mut std::collections::BTreeSet<RunId>,
1413 positions: &mut std::collections::BTreeSet<u64>,
1414 report: &mut VerifyReport,
1415) {
1416 if !blocks_of.insert(pass.run) {
1417 report.findings.push(format!(
1418 "run {}: the file carries it in two blocks, so neither is the run's one history",
1419 pass.run
1420 ));
1421 report.sound.retain(|run| *run != pass.run);
1422 pass.clean = false;
1423 }
1424 if pass.sealed
1425 && let Some(index) = pass.index
1426 && !positions.insert(index)
1427 {
1428 report.findings.push(format!(
1429 "log index {index} is claimed by two runs — one position in the log holds one leaf"
1430 ));
1431 pass.clean = false;
1432 }
1433}
1434
1435/// The verifier's working state for the run block it is inside.
1436///
1437/// One struct rather than five parallel locals, because they reset together —
1438/// a new run block replaces all of them at once, and a field that survived the
1439/// boundary would carry one run's evidence into another's verdict.
1440struct RunPass {
1441 run: RunId,
1442 declared_seal: Option<crate::core::Digest>,
1443 /// Whether the block declared both a log position and a seal.
1444 sealed: bool,
1445 index: Option<u64>,
1446 /// A package's path for this leaf against the header's checkpoint.
1447 proof: Option<Vec<crate::core::Digest>>,
1448 prev: crate::core::Digest,
1449 last_seq: u64,
1450 /// How many record lines this block carried. Zero is a state the trailer
1451 /// must explain: see [`settle_trailer`].
1452 records: usize,
1453 resealed: Vec<crate::journal::Record>,
1454 /// Whether every record in this block checked out so far.
1455 ///
1456 /// [`VerifyReport::sound`] promises *chain recomputed exactly*, and the
1457 /// leaf comparison alone cannot hold that promise for an **open** run —
1458 /// there is no leaf, so without this flag an edited record in an unsealed
1459 /// run produced a finding *and* left the run listed sound.
1460 clean: bool,
1461}
1462
1463/// Hold the export's header to every anchor, and return one it matched.
1464///
1465/// Three answers, because the size relation decides which applies: an anchor
1466/// *above* the file, or of another log, names a history the file cannot be part
1467/// of; an anchor *at* the file's size either matches it or names a second
1468/// history of that size; and an anchor *below* it is checked against the
1469/// file's own first leaves, which the file carries — a mismatch is a finding,
1470/// and a match anchors the prefix and leaves the rest to the header, which is
1471/// said.
1472fn compare_anchors(
1473 report: &mut VerifyReport,
1474 header_seen: bool,
1475 anchors: &[crate::journal::Anchor],
1476 leaves: &[(u64, crate::core::merkle::LeafHash)],
1477) -> Option<Checkpoint> {
1478 let mut matched = None;
1479 if !header_seen {
1480 return matched;
1481 }
1482 for anchor in anchors {
1483 let given = &anchor.checkpoint;
1484 if given.origin != report.checkpoint.origin || given.size > report.checkpoint.size {
1485 report.findings.push(format!(
1486 "the export's header names log '{}' at size {} with root {}, and the \
1487 checkpoint held by {} names '{}' at size {} with root {} — the file \
1488 describes a different history than the one it is being checked against",
1489 report.checkpoint.origin,
1490 report.checkpoint.size,
1491 report.checkpoint.root.to_hex(),
1492 anchor.obtained_from,
1493 given.origin,
1494 given.size,
1495 given.root.to_hex(),
1496 ));
1497 } else if given.size == report.checkpoint.size {
1498 if given.root == report.checkpoint.root {
1499 if matched.is_none() {
1500 matched = Some(given.clone());
1501 }
1502 } else {
1503 report.findings.push(format!(
1504 "the export's header names log '{}' at size {} with root {}, and the \
1505 checkpoint held by {} holds that same size with root {} — one tree of \
1506 a given size has one root, so these are two histories",
1507 report.checkpoint.origin,
1508 report.checkpoint.size,
1509 report.checkpoint.root.to_hex(),
1510 anchor.obtained_from,
1511 given.root.to_hex(),
1512 ));
1513 }
1514 } else if prefix_root(leaves, given.size) != Some(given.root) {
1515 report.findings.push(format!(
1516 "the checkpoint held by {} commits to the first {} run(s) of log '{}' with \
1517 root {}, and this export's first {} run(s) do not rebuild to it — a run \
1518 inside that prefix was removed, replaced or moved",
1519 anchor.obtained_from,
1520 given.size,
1521 given.origin,
1522 given.root.to_hex(),
1523 given.size,
1524 ));
1525 } else {
1526 report.not_checked.push(format!(
1527 "the checkpoint held by {} is at size {} and matches this export's first {} \
1528 run(s); the {} after it are held only to the file's own header",
1529 anchor.obtained_from,
1530 given.size,
1531 given.size,
1532 report.checkpoint.size - given.size
1533 ));
1534 }
1535 }
1536 matched
1537}
1538
1539/// The root of the tree of a file's first `size` leaves, when the file carries
1540/// exactly the positions `0..size` among them.
1541///
1542/// An export holds every leaf from position 0, so a checkpoint smaller than
1543/// the file is a tree the file can rebuild — no consistency proof is needed
1544/// for a reader who holds the leaves themselves. `leaves` is sorted by
1545/// position; a missing or duplicated position below `size` answers `None`.
1546fn prefix_root(
1547 leaves: &[(u64, crate::core::merkle::LeafHash)],
1548 size: u64,
1549) -> Option<crate::core::Digest> {
1550 let size = usize::try_from(size).ok()?;
1551 let prefix = leaves.get(..size)?;
1552 prefix
1553 .iter()
1554 .enumerate()
1555 .all(|(at, (index, _))| u64::try_from(at) == Ok(*index))
1556 .then(|| {
1557 crate::core::merkle::root(&prefix.iter().map(|(_, leaf)| *leaf).collect::<Vec<_>>())
1558 })
1559}
1560
1561/// The checks that can only be made once the whole file has been read.
1562///
1563/// Separated because they answer a different question from the per-record pass:
1564/// that one asks *is each record what it says it is*, and every one of these
1565/// asks *is anything missing* — which no single line can reveal.
1566fn settle(
1567 report: &mut VerifyReport,
1568 header_seen: bool,
1569 mut leaves: Vec<(u64, crate::core::merkle::LeafHash)>,
1570 anchors: &[crate::journal::Anchor],
1571) {
1572 use crate::core::merkle;
1573
1574 if anchors.iter().any(|a| !a.witnessed.is_empty()) {
1575 report.not_checked.push(
1576 "freshness — the anchors carry witness times and verify does not judge them; \
1577 `agentplane audit --max-checkpoint-age` does"
1578 .to_owned(),
1579 );
1580 }
1581 if report.selection.is_some() {
1582 settle_package(report, header_seen, leaves.len(), anchors);
1583 return;
1584 }
1585
1586 // Which checkpoint the rebuild is held to, and everything below turns on
1587 // it. The header's own is a claim by whoever wrote the file; an anchor is
1588 // one the reader was given by somebody else — printed by an earlier audit,
1589 // cosigned by a witness, pasted into a ticket. Only the second makes the
1590 // Merkle rebuild evidence about *deletion*; against the header it is
1591 // evidence that the file is self-consistent, which an editor who dropped a
1592 // run and rewrote the header also achieves.
1593 //
1594 // **Every anchor is consulted, and what cannot be is said.** Three answers,
1595 // because the size relation decides which one applies: an anchor *above*
1596 // the file is a log that shrank, an anchor *at* the file's size either
1597 // matches it or names a different history, and an anchor *below* it is
1598 // rebuilt from the file's own first leaves. Collapsing any of them is what
1599 // lets an operator pick whichever observer their export happens to satisfy.
1600 leaves.sort_by_key(|(index, _)| *index);
1601 let matched = compare_anchors(report, header_seen, anchors, &leaves);
1602 let against = if let Some(checkpoint) = matched {
1603 checkpoint
1604 } else {
1605 if anchors.is_empty() {
1606 report.not_checked.push(
1607 "deletion — no checkpoint was supplied, so the Merkle root could only be \
1608 rebuilt and compared against this file's own header. That proves the \
1609 file is internally consistent, which is also what an editor who dropped \
1610 a run and rewrote the header achieves. Pass the checkpoint an earlier \
1611 audit printed, or one a witness cosigned"
1612 .to_owned(),
1613 );
1614 }
1615 report.checkpoint.clone()
1616 };
1617
1618 if !header_seen {
1619 report
1620 .findings
1621 .push("the export has no header, so nothing says which log it came from".into());
1622 }
1623
1624 // The tree, rebuilt in log order. This is what notices a whole run dropped
1625 // from the middle: every per-run chain above still verified, because a chain
1626 // links records within a run and knows nothing about its neighbours.
1627 let size = u64::try_from(leaves.len()).unwrap_or(u64::MAX);
1628 if size == against.size {
1629 // The positions are part of the claim, not bookkeeping: a checkpoint of
1630 // size N commits to leaves 0..N, so a duplicated or out-of-range
1631 // position is a relabelled log. Named here rather than left to surface
1632 // as a root mismatch, because "the root differs" tells an auditor that
1633 // something is wrong and not that two runs claim one place in history —
1634 // and a tree built over duplicated positions would compare garbage
1635 // against the root and report the wrong defect.
1636 let contiguous = leaves
1637 .iter()
1638 .enumerate()
1639 .all(|(at, (index, _))| u64::try_from(at) == Ok(*index));
1640 if contiguous {
1641 let rebuilt =
1642 merkle::root(&leaves.into_iter().map(|(_, leaf)| leaf).collect::<Vec<_>>());
1643 if rebuilt != against.root {
1644 report.findings.push(
1645 "the Merkle root rebuilt from this export does not match the checkpoint it \
1646 claims to be a copy of"
1647 .to_owned(),
1648 );
1649 }
1650 } else {
1651 report.findings.push(format!(
1652 "the run blocks' log positions are not the contiguous 0..{} the checkpoint \
1653 commits to — a position is duplicated or missing, so this file describes a \
1654 different log than the one it names",
1655 against.size
1656 ));
1657 }
1658 } else {
1659 report.findings.push(format!(
1660 "the export carries {size} sealed run(s) and its checkpoint commits to {} — the \
1661 difference is runs that were in the log and are not in this file",
1662 against.size
1663 ));
1664 }
1665
1666 if !report.complete {
1667 report.findings.push(
1668 "the export has no trailer, so it was cut short — every line in it is still valid, \
1669 which is why the frame is the signal"
1670 .to_owned(),
1671 );
1672 }
1673}
1674
1675/// A package's settlement: its header against each outside checkpoint, and
1676/// what the file does not disclose.
1677///
1678/// Each leaf was proved by its own path in [`finish_run`], so no tree is
1679/// rebuilt and no count is held to the checkpoint's size — the leaves a
1680/// package leaves out are the point of it, not a deletion. An outside
1681/// checkpoint of the header's size is compared by root; one of another size
1682/// would need a consistency proof the package does not carry, and is said to
1683/// be uncompared rather than judged.
1684fn settle_package(
1685 report: &mut VerifyReport,
1686 header_seen: bool,
1687 disclosed: usize,
1688 anchors: &[crate::journal::Anchor],
1689) {
1690 if !header_seen {
1691 report
1692 .findings
1693 .push("the export has no header, so nothing says which log it came from".into());
1694 }
1695 let header = report.checkpoint.clone();
1696 for anchor in anchors {
1697 let given = &anchor.checkpoint;
1698 if given.origin != header.origin {
1699 report.findings.push(format!(
1700 "the package's header names log '{}', and the checkpoint held by {} names \
1701 '{}' — the file describes a different history than the one it is being \
1702 checked against",
1703 header.origin, anchor.obtained_from, given.origin,
1704 ));
1705 } else if given.size != header.size {
1706 report.not_checked.push(format!(
1707 "the checkpoint held by {} is at size {} and the package's header at {}; a \
1708 package carries no consistency proof, so the two were not compared",
1709 anchor.obtained_from, given.size, header.size,
1710 ));
1711 } else if given.root != header.root {
1712 report.findings.push(format!(
1713 "the package's header names log '{}' at size {} with root {}, and the \
1714 checkpoint held by {} holds that same size with root {} — one tree of a \
1715 given size has one root, so these are two histories",
1716 header.origin,
1717 header.size,
1718 header.root.to_hex(),
1719 anchor.obtained_from,
1720 given.root.to_hex(),
1721 ));
1722 }
1723 }
1724 if anchors.is_empty() {
1725 report.not_checked.push(
1726 "the header's checkpoint — no outside checkpoint was supplied, so each path was \
1727 proved against the file's own header, which whoever wrote the file chose. Pass \
1728 the checkpoint the plane published or a witness cosigned"
1729 .to_owned(),
1730 );
1731 }
1732 report.not_checked.push(format!(
1733 "the rest of the log — this file is a disclosure package: the log holds {} sealed \
1734 run(s) and the package proves {disclosed} of them; nothing about the others is in \
1735 the file or was checked",
1736 header.size
1737 ));
1738 if !report.complete {
1739 report.findings.push(
1740 "the export has no trailer, so it was cut short — every line in it is still valid, \
1741 which is why the frame is the signal"
1742 .to_owned(),
1743 );
1744 }
1745}
1746
1747/// Why no digest in an export written under `canon` can be recomputed here,
1748/// when this build does not implement that rule.
1749///
1750/// The rule names the digest algorithm as well as the canonical form, so under
1751/// another one no hash in the file is one this build can check or rebuild.
1752fn canon_unverifiable(canon: Option<u64>) -> Option<String> {
1753 (canon != Some(u64::from(crate::core::canon::VERSION))).then(|| {
1754 format!(
1755 "unknown canon — the export was written under rule {canon:?} and this build \
1756 implements {}, so no digest in it can be recomputed here. Not a finding: a \
1757 build implementing that rule can verify and restore it",
1758 crate::core::canon::VERSION
1759 )
1760 })
1761}
1762
1763/// Why this build cannot hold an export to its digests, when its header line
1764/// names a canon this build does not implement.
1765///
1766/// What `verify` reports as `unverifiable` and `restore` refuses before
1767/// writing, for a caller that must tell that answer apart from damage or an
1768/// outage before reading the whole file. Checked in `verify`'s order: a line
1769/// that is not this format's header, or names another format version, answers
1770/// `None`, and refusing it is the parser's.
1771#[must_use]
1772pub fn foreign_canon(header: &str) -> Option<String> {
1773 let value: serde_json::Value = serde_json::from_str(header).ok()?;
1774 let ours = value.get("kind").and_then(serde_json::Value::as_str) == Some("agentplane.export")
1775 && value.get("version").and_then(serde_json::Value::as_u64)
1776 == Some(u64::from(FORMAT_VERSION));
1777 if !ours {
1778 return None;
1779 }
1780 canon_unverifiable(value.get("canon").and_then(serde_json::Value::as_u64))
1781}
1782
1783/// Read a package header's selection; one naming nothing is a finding.
1784fn read_selection(value: &serde_json::Value, report: &mut VerifyReport) {
1785 let selection = value
1786 .get("selection")
1787 .and_then(|s| serde_json::from_value::<Selection>(s.clone()).ok())
1788 .filter(|s| !(s.cases.is_empty() && s.runs.is_empty()));
1789 if selection.is_none() {
1790 report.findings.push(
1791 "the package header names no case and no run, so nothing says what it was a \
1792 disclosure of"
1793 .to_owned(),
1794 );
1795 }
1796 report.selection = Some(selection.unwrap_or_default());
1797}
1798
1799/// Read the header line: which format, which log, at what size, under which rule.
1800fn read_header(value: &serde_json::Value, report: &mut VerifyReport) {
1801 let version = value.get("version").and_then(serde_json::Value::as_u64);
1802 if version != Some(u64::from(FORMAT_VERSION)) {
1803 report.findings.push(format!(
1804 "the export claims format version {version:?} and this build reads {FORMAT_VERSION} \
1805 — the findings below describe the lines this build could interpret, which may not \
1806 be all of them"
1807 ));
1808 }
1809 if let Some(why) = canon_unverifiable(value.get("canon").and_then(serde_json::Value::as_u64)) {
1810 report.unverifiable = Some(why);
1811 }
1812 match value
1813 .get("checkpoint")
1814 .and_then(|c| serde_json::from_value::<Checkpoint>(c.clone()).ok())
1815 {
1816 Some(c) => report.checkpoint = c,
1817 None => report
1818 .findings
1819 .push("the header carries no readable checkpoint".to_owned()),
1820 }
1821}
1822
1823/// Rehash one record's wire bytes and hold them to the hash it carries.
1824///
1825/// The rehash is the whole check, and it runs over `raw` — the exact bytes the
1826/// store hashed — never over a re-serialization of the parsed body. That is
1827/// the journal's own wire-bytes rule: re-serializing would hold the file to
1828/// *this build's* canonicalization instead of to what was written, so an
1829/// export from a build whose rule differed would report tampering where there
1830/// was only time, and — worse — an edit that re-serializes identically would
1831/// pass. Comparing the file's `prev_hash` against the previous line's `hash`
1832/// would only prove the file agrees with itself, which an editor who
1833/// recomputed the chain also achieves; rehashing the wire bytes is what makes
1834/// agreement evidence about them.
1835fn read_record(
1836 value: &serde_json::Value,
1837 pass: &mut RunPass,
1838 report: &mut VerifyReport,
1839 stamped: &mut std::collections::BTreeSet<crate::core::CaseId>,
1840 upcaster: &dyn crate::journal::Upcaster,
1841) {
1842 let current = pass.run;
1843 pass.records += 1;
1844 let (Some(raw), Some(claimed)) = (
1845 value.get("raw").and_then(serde_json::Value::as_str),
1846 value
1847 .get("hash")
1848 .and_then(|h| serde_json::from_value::<crate::core::Digest>(h.clone()).ok()),
1849 ) else {
1850 report.findings.push(format!(
1851 "run {current}: a record line carries no wire bytes or no hash — nothing ties \
1852 it to the chain"
1853 ));
1854 pass.clean = false;
1855 return;
1856 };
1857 let raw_bytes = raw.as_bytes();
1858 // The hash first, then the parse. Bytes that do not hash to their claim
1859 // were edited, whatever they parse as; only bytes that do can make a parse
1860 // failure a statement about the reader rather than about the file.
1861 if crate::core::Digest::chain(pass.prev, raw_bytes) != claimed {
1862 return edited_record(raw_bytes, pass, report);
1863 }
1864 // Then the spelling, before a parse failure is believed: bytes that hash
1865 // to their claim and that no writer under this canon produces — a space, a
1866 // member out of order, a member written twice — are a statement about the
1867 // file, and the last would otherwise fail the parse and read as a skew.
1868 // Compared as values, so a record from another shape is held to the rule.
1869 if serde_json::from_slice::<serde_json::Value>(raw_bytes)
1870 .is_ok_and(|wire| crate::core::canon::value_bytes(&wire) != raw_bytes)
1871 {
1872 return uncanonical_record(raw_bytes, pass, report);
1873 }
1874 // The body verification reads is read from the wire bytes — the one
1875 // source the hash actually covers — through the upcaster, which compares
1876 // the record's version before its shape is believed: a record from another
1877 // shape is lifted, or refused as a skew, rather than failing a parse into
1878 // this build's struct and being judged by that.
1879 let signature = value
1880 .get("signature")
1881 .and_then(|a| serde_json::from_value::<Option<crate::core::KeySignature>>(a.clone()).ok())
1882 .flatten();
1883 let record = match crate::journal::Record::from_stored_with(
1884 upcaster,
1885 raw_bytes.to_vec(),
1886 pass.prev,
1887 claimed,
1888 signature,
1889 ) {
1890 Ok(record) => record,
1891 Err(unread) => return unread_record(raw_bytes, &unread, pass, report),
1892 };
1893 let body = &record.body;
1894 // The readable `body` is a courtesy copy, and it is held to the bytes: a
1895 // file whose display half says something its hashed half does not is the
1896 // quiet edit — every hash verifies, and the reader was shown a lie.
1897 let wire: serde_json::Value = serde_json::from_slice(raw_bytes).unwrap_or_default();
1898 if value.get("body") != Some(&wire) {
1899 report.findings.push(format!(
1900 "run {current}: record {}'s readable body does not match its wire bytes — the \
1901 display copy was edited, and every hash still verifies over the real one",
1902 body.seq
1903 ));
1904 pass.clean = false;
1905 }
1906 // Collected for the case-coverage settlement: a stamp is the journal
1907 // naming a matter, and the case layer must carry every matter it names.
1908 if let Some(case) = body.case {
1909 stamped.insert(case);
1910 }
1911 // The record's own body names its run, and it must be the run the block
1912 // claims. Without this comparison an export could relabel a whole history —
1913 // run B's records and B's leaf filed under A's id — and every other check
1914 // would pass, because chain, seal and Merkle all verify B's bytes; only the
1915 // *label* lied, and the label is what the reader looks a run up by.
1916 if body.run != current {
1917 report.findings.push(format!(
1918 "run {current}: a record in this block belongs to run {} — the block was relabelled, \
1919 or spliced from another history",
1920 body.run
1921 ));
1922 pass.clean = false;
1923 }
1924 // A removed record breaks a link; a removed *tail* does not, which is why
1925 // the sequence is checked as well as the chain.
1926 if body.seq != pass.last_seq + 1 {
1927 report.findings.push(format!(
1928 "run {current}: seq {} follows {}, so a record is missing from the middle — \
1929 every chain link either side of the gap still joins",
1930 body.seq, pass.last_seq
1931 ));
1932 pass.clean = false;
1933 }
1934 pass.last_seq = body.seq;
1935
1936 // The sealing record's own claim, held to the chain it sits in — the same
1937 // check the live audit makes. `RunSealed.chain_head` is the head the
1938 // conclusion was drawn over, which is by construction its own record's
1939 // `prev_hash`; `pass.prev` here is that head, recomputed from the wire
1940 // bytes of every line before this one, so agreement is evidence about the
1941 // bytes rather than the file agreeing with itself. A mismatch means the
1942 // conclusion was composed against a different history than the one it was
1943 // appended to, which no honest writer produces. What this does NOT cover:
1944 // a run with no sealing record at all — an open run has made no claim,
1945 // and its absence of one is a state, not a defect.
1946 if let crate::journal::RecordKind::RunConcluded { chain_head, .. } = &body.kind
1947 && *chain_head != pass.prev
1948 {
1949 report.findings.push(format!(
1950 "run {current}: the sealing record claims a chain head that is not the head it \
1951 sits on — the conclusion was drawn over a different history"
1952 ));
1953 pass.clean = false;
1954 }
1955
1956 pass.prev = record.hash;
1957 pass.resealed.push(record);
1958}
1959
1960/// Members of a framing line this build does not know, reported as unchecked.
1961///
1962/// **A verdict is only as wide as the claims the reader understood.** A framing
1963/// line carries claims that are checked — a checkpoint, a leaf, the trailer's
1964/// accounting — so a member added by a later writer may carry one more, and a
1965/// reader that passes over it reports *sound* about a file it read part of.
1966///
1967/// This is `not_checked` rather than a finding, and the difference is the one
1968/// the record path draws for the same question. A record's bytes are hashed, so
1969/// a member nobody knows is refused: the verdict would otherwise be reached over
1970/// evidence the reader did not see. A framing line is not hashed and carries no
1971/// evidence of its own, so an unknown member does not falsify anything already
1972/// checked — it bounds what the check covered, which is what this field is for.
1973fn note_unknown_members(
1974 kind: &str,
1975 package: bool,
1976 value: &serde_json::Value,
1977 report: &mut VerifyReport,
1978) {
1979 let known: &[&str] = match kind {
1980 "agentplane.export" => &["kind", "version", "checkpoint", "canon"],
1981 DISCLOSURE_KIND => &["kind", "version", "checkpoint", "canon", "selection"],
1982 "agentplane.export.run" if package => &["kind", "run", "index", "seal", "proof"],
1983 "agentplane.export.run" => &["kind", "run", "index", "seal"],
1984 "agentplane.export.case" => &["kind", "case", "deadlines", "blobs", "hold"],
1985 "agentplane.export.end" => &[
1986 "kind",
1987 "runs_requested",
1988 "runs_exported",
1989 "records",
1990 "cases",
1991 "unreadable",
1992 ],
1993 // Not a frame: a record line carries no top-level `kind`, and its own
1994 // members are refused rather than noted, one level down.
1995 _ => return,
1996 };
1997 let Some(object) = value.as_object() else {
1998 return;
1999 };
2000 let unknown: Vec<&str> = object
2001 .keys()
2002 .map(String::as_str)
2003 .filter(|member| !known.contains(member))
2004 .collect();
2005 if !unknown.is_empty() {
2006 report.not_checked.push(format!(
2007 "a {kind} line carries {} this build does not know — whatever they claim was \
2008 not checked, and a later build wrote this file",
2009 unknown.join(", ")
2010 ));
2011 }
2012}
2013
2014/// A record line whose bytes do not hash to the hash it carries.
2015///
2016/// Filed before any parse: whatever these bytes parse as, they are not the
2017/// bytes the chain committed to, so a parse failure here says nothing about
2018/// the reader's age.
2019fn edited_record(raw_bytes: &[u8], pass: &mut RunPass, report: &mut VerifyReport) {
2020 let seq = serde_json::from_slice::<serde_json::Value>(raw_bytes)
2021 .ok()
2022 .and_then(|v| v.get("seq").and_then(serde_json::Value::as_u64))
2023 .unwrap_or(pass.last_seq + 1);
2024 report.findings.push(format!(
2025 "run {}: record {seq} does not recompute to the hash it carries — it was edited \
2026 after it was sealed",
2027 pass.run
2028 ));
2029 pass.clean = false;
2030 // The head and the sequence walk forward over the bytes actually present,
2031 // so the leaf comparison at the end of the block speaks about what this
2032 // file carries rather than about the first mismatch.
2033 pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
2034 pass.last_seq = seq;
2035}
2036
2037/// A record whose bytes hash to their claim and are not canonical, filed as a
2038/// finding, with the head and sequence walked forward over the bytes present.
2039fn uncanonical_record(raw_bytes: &[u8], pass: &mut RunPass, report: &mut VerifyReport) {
2040 let seq = serde_json::from_slice::<serde_json::Value>(raw_bytes)
2041 .ok()
2042 .and_then(|v| v.get("seq").and_then(serde_json::Value::as_u64))
2043 .unwrap_or(pass.last_seq + 1);
2044 report.findings.push(format!(
2045 "run {}: record {seq}'s wire bytes are not canonical — they hash to their claim, \
2046 and no writer under the export's canon produces them",
2047 pass.run
2048 ));
2049 pass.clean = false;
2050 pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
2051 pass.last_seq = seq;
2052}
2053
2054/// A record line this reader cannot read, filed under what that means.
2055///
2056/// **A line this reader cannot read is not the same as a line nobody can.**
2057/// The hash was verified before this is reached, so the bytes are the ones the
2058/// chain committed to and a failure is a statement about the reader. A record
2059/// at a version the upcaster cannot reach, or at this build's version in a
2060/// shape it does not parse, is a build skew; answering *edited* or *malformed*
2061/// for either would report an export written by another build as a damaged
2062/// file, record by record, to the one audience that has no other copy. Only
2063/// bytes that are not a record at all are malformed.
2064fn unread_record(
2065 raw_bytes: &[u8],
2066 unread: &crate::core::StoreError,
2067 pass: &mut RunPass,
2068 report: &mut VerifyReport,
2069) {
2070 let current = pass.run;
2071 let at = serde_json::from_slice::<serde_json::Value>(raw_bytes)
2072 .ok()
2073 .and_then(|v| v.get("seq").and_then(serde_json::Value::as_u64));
2074 match unread {
2075 crate::core::StoreError::UnknownRecordVersion { .. } => {
2076 report.findings.push(format!(
2077 "run {current}: record {} is at a version this build does not read, and \
2078 its bytes hash as written — this is a build skew rather than an edit: \
2079 {unread}",
2080 at.unwrap_or(pass.last_seq + 1)
2081 ));
2082 }
2083 crate::core::StoreError::UnreadableRecordShape { .. } => {
2084 report.findings.push(format!(
2085 "run {current}: a record is at a shape this build does not read — a build \
2086 skew rather than a damaged file: {unread}"
2087 ));
2088 }
2089 _ => report
2090 .findings
2091 .push(format!("run {current}: a record line is malformed")),
2092 }
2093 pass.clean = false;
2094 // The head and the sequence walk forward over what the file carries, so the
2095 // records after this one are compared against the history the file actually
2096 // holds. Without it one unreadable line makes every later record in the
2097 // block report a broken link and a gap — a cascade of incident-shaped
2098 // findings from one old reader.
2099 pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
2100 if let Some(seq) = at {
2101 pass.last_seq = seq;
2102 }
2103}
2104
2105/// Close out a run block: its terminal hash must be the leaf the log recorded,
2106/// and its signatures must verify if a key was supplied.
2107fn finish_run(
2108 report: &mut VerifyReport,
2109 pass: Option<RunPass>,
2110 verifier: Option<&dyn crate::core::Verifier>,
2111 read_runs: &mut usize,
2112 empty_blocks: &mut Vec<RunId>,
2113 observe: &mut dyn FnMut(ClosedRun<'_>),
2114) {
2115 let Some(pass) = pass else {
2116 return;
2117 };
2118 let run = pass.run;
2119 // A block with no records is never sound, and it is never judged here:
2120 // whether it is an honestly-declared unreadable run (unchecked) or a run
2121 // emptied after the export was taken (a finding) is written in the
2122 // trailer, which this pass has not necessarily reached — an intermediate
2123 // block closes when the next one starts. Judging it now would also raise a
2124 // false leaf-mismatch for a sealed unreadable run, whose declared leaf is
2125 // genuine and whose records the writer honestly could not read: `prev` is
2126 // still `ZERO`, and ZERO not matching the leaf is a fact about the empty
2127 // walk, not about the history.
2128 if pass.records == 0 {
2129 empty_blocks.push(run);
2130 observe(ClosedRun {
2131 run,
2132 sealed: pass.sealed,
2133 records: &[],
2134 });
2135 return;
2136 }
2137 *read_runs += 1;
2138 let mut ok = pass.clean;
2139
2140 // The one cross-check between the two halves of the export. Without it a
2141 // file could carry a healthy chain and a leaf belonging to some other
2142 // history, and each half would verify on its own.
2143 if let Some(seal) = pass.declared_seal
2144 && seal != pass.prev
2145 {
2146 report.findings.push(format!(
2147 "run {run}: the log's leaf is not this run's terminal hash, so the chain in this \
2148 file is not the chain the checkpoint committed to"
2149 ));
2150 ok = false;
2151 }
2152
2153 // A package places every sealed run it carries, so a conclusion that
2154 // seals under a block with no leaf is a leaf stripped from the file.
2155 if report.selection.is_some()
2156 && !pass.sealed
2157 && crate::audit::has_sealing_conclusion(&pass.resealed)
2158 {
2159 report.findings.push(format!(
2160 "run {run}: it concluded under an outcome that seals and its block carries no \
2161 leaf — a package places every sealed run it carries, so this one's place in the \
2162 log was removed"
2163 ));
2164 ok = false;
2165 }
2166
2167 // In a package no tree is rebuilt, so the path is the whole of the
2168 // evidence that this leaf is in the history the header names.
2169 if report.selection.is_some()
2170 && let Some(seal) = pass.declared_seal
2171 && !leaf_is_proved(seal, pass.index, pass.proof.as_deref(), &report.checkpoint)
2172 {
2173 report.findings.push(format!(
2174 "run {run}: its path does not prove its leaf against the header's checkpoint, so \
2175 nothing ties this run to the history the package names"
2176 ));
2177 ok = false;
2178 }
2179
2180 // One implementation of *is this signed history sound*, and it is the
2181 // crate's own. `require_signature` is true because this is the auditor's
2182 // posture: an unsigned record inside a signed history is the one an
2183 // attacker who cannot sign would add.
2184 if let Some(v) = verifier
2185 && let Err(e) = crate::journal::Record::verify_signed(
2186 &pass.resealed,
2187 crate::core::Digest::ZERO,
2188 v,
2189 true,
2190 )
2191 {
2192 report.findings.push(format!("run {run}: {e}"));
2193 ok = false;
2194 }
2195
2196 if ok {
2197 report.sound.push(run);
2198 }
2199 observe(ClosedRun {
2200 run,
2201 sealed: pass.sealed,
2202 records: &pass.resealed,
2203 });
2204}
2205
2206/// Whether `proof` proves `seal` at `index` in the tree `checkpoint` commits to.
2207fn leaf_is_proved(
2208 seal: crate::core::Digest,
2209 index: Option<u64>,
2210 proof: Option<&[crate::core::Digest]>,
2211 checkpoint: &Checkpoint,
2212) -> bool {
2213 let (Some(index), Some(proof)) = (index, proof) else {
2214 return false;
2215 };
2216 let (Ok(index), Ok(size)) = (usize::try_from(index), usize::try_from(checkpoint.size)) else {
2217 return false;
2218 };
2219 crate::core::merkle::verify_inclusion(
2220 crate::core::merkle::leaf_hash(&seal),
2221 index,
2222 size,
2223 proof,
2224 &checkpoint.root,
2225 )
2226}
2227
2228/// What a restore loses beyond the case layer, as sentences a reader can act
2229/// on.
2230///
2231/// Extracted so the restore's control flow is the *writing* and this is the
2232/// *accounting*. Each entry names one loss: a count of them would tell an
2233/// operator nothing about which one costs them work, and exactly one of these
2234/// does.
2235fn losses(parsed: &Parsed) -> Vec<String> {
2236 let mut out = Vec::new();
2237 if parsed.signed > 0 && !parsed.runs.is_empty() {
2238 out.push(format!(
2239 "{} record(s) carried a signature that this store did not reproduce — `append` \
2240 attests as the restoring store's own signer, so authorship is lost unless it \
2241 holds the original key. Hashes and the Merkle root are unaffected",
2242 parsed.signed
2243 ));
2244 }
2245 out.push(
2246 "activity timestamps — `recent_runs` now orders by restore time rather than by when \
2247 history happened. It is a discovery index for listing, and nothing derives a decision \
2248 from it"
2249 .to_owned(),
2250 );
2251
2252 // Named whether or not this export happens to hold a waiting run. The
2253 // alternative — say it only when `awaiting` is non-empty — makes the
2254 // absence of the sentence mean two different things, and the reader who
2255 // needs it most is the one restoring an export they did not write.
2256 out.push(
2257 "every wait's registration — a timer, a subscription and any worklist row a wait \
2258 opened live in stores this export does not carry, so a restored run that was \
2259 waiting has nothing to wake it: no timer fires, no subscription matches, and its \
2260 lease was released cleanly when it suspended, so recovery does not see it either. \
2261 Resuming each run in `awaiting` re-arms the wait from the journal"
2262 .to_owned(),
2263 );
2264 out.push(
2265 "the worklist, unclaimed inbound events, webhook registrations and their delivery \
2266 cursors, governed memory, and the batch, quota and standing-authority ledgers — \
2267 none of these layers is in the export. A decision a run already consumed survives \
2268 because that run journaled it; one nobody had consumed does not"
2269 .to_owned(),
2270 );
2271
2272 out
2273}
2274
2275/// Which of the restored runs came back waiting.
2276///
2277/// Read with the same function every other surface answers *what does this
2278/// run's history say* with. A fourth copy of the match would be the copy that
2279/// disagrees the day a record kind arrives.
2280async fn awaiting_runs(
2281 store: &Arc<dyn JournalStore>,
2282 runs: &[RestoredRun],
2283) -> Result<Vec<RunId>, StoreError> {
2284 let mut awaiting = Vec::new();
2285 for run in runs {
2286 let records = store.read(run.run, 1).await?;
2287 if matches!(
2288 crate::runtime::observed_status(&records),
2289 Some(crate::runtime::RunStatus::Suspended(_))
2290 ) {
2291 awaiting.push(run.run);
2292 }
2293 }
2294 Ok(awaiting)
2295}
2296
2297/// Every run the outcome indexes cannot name: the ones still in flight.
2298///
2299/// **Selecting what to export is two questions, and this is the second.**
2300/// `runs_by_outcome` indexes *conclusions*, so a run that has not concluded is
2301/// in no outcome, and an export driven by
2302/// [`OUTCOMES_OF_RECORD`](crate::runtime::OUTCOMES_OF_RECORD) alone carries no
2303/// run that is working, sleeping, awaiting a message or waiting on a person —
2304/// which is the work a disaster recovery is for. Nothing downstream can notice:
2305/// the Merkle log commits to **sealed** runs, so a file missing every in-flight
2306/// run restores to an equal root at an equal size and reports itself faithful.
2307///
2308/// **Paged, and bounded by `limit` like every other listing here.** It walks
2309/// every run by id ([`JournalStore::runs_by_id`]) and keeps the runs whose
2310/// history has not concluded, reading **one record** per candidate to decide:
2311/// the head's sequence, then that record. A run whose records cannot be read
2312/// is not silently dropped; it is returned in `unreadable` for the caller to
2313/// report, on the same principle as the export's own trailer.
2314///
2315/// # Errors
2316///
2317/// If the activity index cannot be paged. A single unreadable run is reported
2318/// rather than raised: one damaged run must not cost an operator the export of
2319/// every other.
2320pub async fn runs_in_flight(
2321 store: &Arc<dyn JournalStore>,
2322 limit: usize,
2323) -> Result<InFlight, StoreError> {
2324 let mut found = InFlight::default();
2325 // Pages of every run, not of the answer: most runs in a healthy plane have
2326 // concluded, so the page that yields one in-flight run may have held five
2327 // hundred that had ended. By id, because an id never moves: a run that
2328 // writes during the walk stays where the cursor will reach it.
2329 let mut after: Option<RunId> = None;
2330 while found.runs.len() < limit {
2331 let page = store.runs_by_id(after, CASE_PAGE).await?;
2332 if page.is_empty() {
2333 break;
2334 }
2335 after = page.last().copied();
2336 for run in page {
2337 if found.runs.len() == limit {
2338 found.truncated = true;
2339 return Ok(found);
2340 }
2341 let head = store.head(run).await?;
2342 if head.seq == 0 {
2343 continue;
2344 }
2345 // The last record alone. `read` is inclusive-from, so this is the
2346 // cheapest question the store answers about a run's state, and it
2347 // is the one `observed_status` needs.
2348 match store.read(run, head.seq).await {
2349 Ok(last) => {
2350 // Working (`None` — the last record is neither a
2351 // suspension nor a conclusion) and waiting
2352 // (`Suspended`) are both in flight. Everything else has
2353 // ended, and the wildcard is the arm that matters: a
2354 // conclusion this build cannot interpret is still a
2355 // conclusion, which `observed_status` guarantees by
2356 // failing an unrecognised outcome closed into
2357 // `Quarantined` rather than into `None`.
2358 let in_flight = match crate::runtime::observed_status(&last) {
2359 None | Some(crate::runtime::RunStatus::Suspended(_)) => true,
2360 Some(_) => false,
2361 };
2362 if in_flight {
2363 found.runs.push(run);
2364 }
2365 }
2366 Err(e) => found.unreadable.push((run, e.to_string())),
2367 }
2368 }
2369 }
2370 Ok(found)
2371}
2372
2373/// What [`runs_in_flight`] found, and what it could not read.
2374#[derive(Debug, Clone, Default, PartialEq, Eq)]
2375pub struct InFlight {
2376 /// Runs that had not concluded, by run id.
2377 pub runs: Vec<RunId>,
2378 /// The limit was reached, so this is a page rather than the set.
2379 pub truncated: bool,
2380 /// Runs the activity index names and whose records would not read.
2381 pub unreadable: Vec<(RunId, String)>,
2382}
2383
2384/// The runs an offline reader of one tenant reads, and what it could not reach.
2385#[derive(Debug, Clone, Default, PartialEq, Eq)]
2386pub struct RunsToRead {
2387 /// The runs, concluded ones by outcome first, then those in flight.
2388 pub runs: Vec<RunId>,
2389 /// Each outcome whose listing overflowed `limit`, and `"in-flight runs"`
2390 /// when that listing did.
2391 pub reached: Vec<String>,
2392 /// Runs the activity index names and whose records would not read.
2393 pub unreadable: Vec<(RunId, String)>,
2394 /// How many of `runs` are in flight.
2395 pub in_flight: usize,
2396}
2397
2398/// The runs under `outcomes`, at most `limit` per outcome, plus the runs still
2399/// in flight when `include_in_flight`.
2400///
2401/// Each outcome is asked for one more than `limit`, so a full page and an
2402/// overflowing one are distinguishable and the overflow is named in
2403/// [`RunsToRead::reached`].
2404///
2405/// # Errors
2406///
2407/// If an index cannot be paged.
2408pub async fn runs_to_read(
2409 store: &Arc<dyn JournalStore>,
2410 outcomes: &[String],
2411 include_in_flight: bool,
2412 limit: usize,
2413) -> Result<RunsToRead, StoreError> {
2414 let mut found = RunsToRead::default();
2415 for outcome in outcomes {
2416 let runs = store.runs_by_outcome(outcome, limit + 1).await?;
2417 if runs.len() > limit {
2418 found.reached.push(outcome.clone());
2419 }
2420 found.runs.extend(runs.into_iter().take(limit));
2421 }
2422 if include_in_flight {
2423 let flight = runs_in_flight(store, limit).await?;
2424 if flight.truncated {
2425 found.reached.push("in-flight runs".to_owned());
2426 }
2427 found.in_flight = flight.runs.len();
2428 found.runs.extend(flight.runs);
2429 found.unreadable = flight.unreadable;
2430 }
2431 Ok(found)
2432}
2433
2434/// One export file's runs, parsed from each record's `raw`.
2435#[derive(Debug, Clone, Default)]
2436pub(crate) struct ExportRuns {
2437 pub(crate) runs: Vec<(RunId, Vec<crate::journal::RecordBody>)>,
2438 pub(crate) unreadable: Vec<RunId>,
2439 /// The `agentplane.export.end` line was read.
2440 pub(crate) trailer: bool,
2441}
2442
2443/// Why an export file could not be read.
2444#[derive(Debug)]
2445pub(crate) enum ReadError {
2446 Io(std::io::Error),
2447 NotAnExport(String),
2448}
2449
2450/// Read an export's runs for an offline consumer.
2451///
2452/// Parsed from each record's `raw` — the bytes the chain hashed — never from
2453/// the display copy beside them. Integrity is not checked here; that is
2454/// [`verify`]'s, so a record that does not parse is reported as unreadable and
2455/// never named a build skew, which only hash-verified bytes can be.
2456pub(crate) fn read_runs<R: std::io::BufRead>(input: R) -> Result<ExportRuns, ReadError> {
2457 use serde_json::Value;
2458 let mut out = ExportRuns::default();
2459 let mut header = false;
2460 for (index, line) in input.lines().enumerate() {
2461 let line = line.map_err(ReadError::Io)?;
2462 if line.trim().is_empty() {
2463 continue;
2464 }
2465 let value: Value = serde_json::from_str(&line)
2466 .map_err(|e| ReadError::NotAnExport(format!("line {} is not JSON: {e}", index + 1)))?;
2467 match value.get("kind").and_then(Value::as_str) {
2468 Some("agentplane.export") => {
2469 let version = value.get("version").and_then(Value::as_u64);
2470 if Some(u64::from(FORMAT_VERSION)) != version {
2471 return Err(ReadError::NotAnExport(format!(
2472 "the export is at format version {version:?}, and this build reads \
2473 {FORMAT_VERSION}"
2474 )));
2475 }
2476 header = true;
2477 }
2478 Some(DISCLOSURE_KIND) if !header => {
2479 return Err(ReadError::NotAnExport(PACKAGE_REFUSED.to_owned()));
2480 }
2481 _ if !header => {
2482 return Err(ReadError::NotAnExport(
2483 "the first line is not an agentplane export header".into(),
2484 ));
2485 }
2486 Some("agentplane.export.run") => {
2487 let run = value
2488 .get("run")
2489 .cloned()
2490 .and_then(|r| serde_json::from_value::<RunId>(r).ok())
2491 .ok_or_else(|| {
2492 ReadError::NotAnExport(format!("line {} names no run", index + 1))
2493 })?;
2494 out.runs.push((run, Vec::new()));
2495 }
2496 Some("agentplane.export.end") => {
2497 out.trailer = true;
2498 if let Some(list) = value.get("unreadable").and_then(Value::as_array) {
2499 out.unreadable.extend(list.iter().filter_map(|u| {
2500 u.get("run")
2501 .cloned()
2502 .and_then(|r| serde_json::from_value::<RunId>(r).ok())
2503 }));
2504 }
2505 }
2506 Some(_) => {}
2507 None => {
2508 let raw = value.get("raw").and_then(Value::as_str).ok_or_else(|| {
2509 ReadError::NotAnExport(format!("line {} carries no wire bytes", index + 1))
2510 })?;
2511 let upcaster = crate::journal::current_upcaster();
2512 let body =
2513 crate::journal::RecordBody::read_through(upcaster.as_ref(), raw.as_bytes())
2514 .map_err(|e| {
2515 // No hash is checked here, so a parse failure is
2516 // not evidence that another build wrote the bytes,
2517 // and the skew's own wording would claim it was.
2518 let why = match e {
2519 StoreError::UnreadableRecordShape {
2520 kind,
2521 version,
2522 detail,
2523 } => {
2524 format!(
2525 "its bytes name {kind} at v{version} and do not parse \
2526 as that shape ({detail}); this reader checks no hash, \
2527 so whether another build wrote them or they were \
2528 edited is `verify`'s to say"
2529 )
2530 }
2531 other => other.to_string(),
2532 };
2533 ReadError::NotAnExport(format!(
2534 "line {} holds a record this build does not read: {why}",
2535 index + 1
2536 ))
2537 })?;
2538 let Some((run, records)) = out.runs.last_mut() else {
2539 return Err(ReadError::NotAnExport(format!(
2540 "line {} is a record before any run block",
2541 index + 1
2542 )));
2543 };
2544 if body.run != *run {
2545 return Err(ReadError::NotAnExport(format!(
2546 "line {} belongs to run {}, filed under {run}",
2547 index + 1,
2548 body.run
2549 )));
2550 }
2551 records.push(body);
2552 }
2553 }
2554 }
2555 if !header {
2556 return Err(ReadError::NotAnExport("the input is empty".into()));
2557 }
2558 Ok(out)
2559}
2560
2561// ── Putting one back ────────────────────────────────────────────────────────
2562
2563/// What a restore did, and what it could not carry across.
2564#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
2565pub struct RestoreReport {
2566 /// The checkpoint the export claimed.
2567 pub expected: Checkpoint,
2568 /// The checkpoint the rebuilt store now reports.
2569 ///
2570 /// **These matching is the whole result.** Equal roots at equal size means
2571 /// every record, in every run, in the order the log recorded them, rebuilt
2572 /// to the same commitment — which is a far stronger statement than "the
2573 /// rows loaded".
2574 pub rebuilt: Checkpoint,
2575 pub runs: usize,
2576 pub records: usize,
2577 /// Cases rebuilt from the export's case layer.
2578 pub cases: usize,
2579 /// Runs whose history ends in a wait, and which nothing will now wake.
2580 ///
2581 /// **Ids rather than a count, because the operator has to act on each
2582 /// one.** A wait is journaled but what performs it is not: the timer, the
2583 /// subscription and the task row live in stores this export does not
2584 /// carry. A restored suspended run has no timer to fire, no subscription
2585 /// to match, and released its lease cleanly when it suspended — so the
2586 /// recovery pass does not see it either. Nothing in the system names it,
2587 /// which is the state the runtime otherwise repairs on sight.
2588 ///
2589 /// Resuming each one repairs it: replay reaches the announced wait, finds
2590 /// no terminal record, and re-arms from the journal. That needs the
2591 /// agent's own code, so it is the caller's step and not the restore's.
2592 pub awaiting: Vec<RunId>,
2593 /// What did not survive, named rather than counted.
2594 pub not_carried: Vec<String>,
2595}
2596
2597impl RestoreReport {
2598 /// Whether the rebuilt store commits to exactly the history the export did.
2599 ///
2600 /// **The commitment, not the label.** A checkpoint carries a log *identity*
2601 /// beside its size and root, and a recovery routinely changes that: the
2602 /// realistic restore is into another tenant of a database somebody else is
2603 /// already using, which is the topology `for_tenant` exists for. Comparing
2604 /// the identity made a byte-perfect restore report as a failed one exactly
2605 /// in the case a disaster puts an operator in — and `agentplane restore`
2606 /// exits on this predicate.
2607 ///
2608 /// A relabelling is not silent for being excluded: it is a sentence in
2609 /// [`not_carried`](Self::not_carried), beside every other thing the file
2610 /// could not bring across, and both checkpoints are on the report for a
2611 /// reader who wants to see the names.
2612 #[must_use]
2613 pub fn is_faithful(&self) -> bool {
2614 self.expected.size == self.rebuilt.size && self.expected.root == self.rebuilt.root
2615 }
2616}
2617
2618/// Rebuild a store from an export, then prove it by its own checkpoint.
2619///
2620/// # Why this goes through `append` rather than writing rows
2621///
2622/// The obvious implementation inserts records verbatim and rebuilds each index
2623/// beside them. It is also the one that fails quietly: `append` maintains
2624/// several derived structures — the case index, the exactly-once index, the
2625/// outcome index and its ordering counter, the admission index, both halves of
2626/// the activity index — and a restore that reconstructed all but one of them
2627/// would produce a store that reads perfectly until somebody queries the one it
2628/// missed.
2629///
2630/// Deliberately not a count. A number here is a claim that has to be re-checked
2631/// on every edit and is not, so it goes stale silently and reads as coverage —
2632/// the shape this project catalogues and has been bitten by. Going through
2633/// `append` is what makes the list not need enumerating: whatever `append`
2634/// maintains, a restore maintains.
2635///
2636/// So this replays the ordinary write path, and every constraint the store
2637/// enforces is enforced here too. What travels through it is each record's
2638/// written bytes with its body ([`Append::restored`]): the store indexes the
2639/// body, lifted through the upcaster if the record is from an older shape, and
2640/// stores the bytes as written, so a restore across a shape change rebuilds
2641/// the chain that was exported rather than a re-seal of it. Three properties
2642/// make each record land at the position its bytes name:
2643///
2644/// * **`seq` is re-derived and lands identically**, because a run restored into
2645/// an empty store starts from the same genesis and receives the same records
2646/// in the same order.
2647/// * **`epoch` is carried, not re-derived.** It is a field of the hashed body,
2648/// so a run that ever changed hands — the ones a disaster is most likely to
2649/// involve — would hash differently under a single fresh lease. `append`
2650/// takes the epoch as a parameter and fences only when a lease row *exists*,
2651/// so restoring into a store with no leases writes each record under its own
2652/// original epoch. Records are grouped into runs of equal epoch for exactly
2653/// this reason.
2654/// * **Runs are sealed in log-index order**, so the Merkle log is rebuilt in the
2655/// order the original recorded, which is what makes the roots comparable at
2656/// all.
2657///
2658/// # What does not survive
2659///
2660/// Every loss is a sentence in [`RestoreReport::not_carried`], because the
2661/// reader who needs it is holding the report rather than this page. One of
2662/// them costs work rather than metadata and has its own field:
2663/// [`RestoreReport::awaiting`] names the runs that came back waiting, and says
2664/// there why nothing will wake them and what does.
2665///
2666/// # Errors
2667///
2668/// If the export cannot be read, if the store seals payloads as it writes,
2669/// or if the store already holds any of its runs — each refused before
2670/// anything is written: this rebuilds a history, it does not merge one. Also
2671/// if the store refuses a write, or if a rebuilt record's hash is not the one
2672/// the file claims; those land mid-restore and leave a partial store that a
2673/// retry refuses, so discard it and restore into a fresh one.
2674pub async fn from_jsonl<R: std::io::BufRead>(
2675 store: &Arc<dyn JournalStore>,
2676 cases: Option<&Arc<dyn crate::case::CaseStore>>,
2677 input: R,
2678) -> Result<RestoreReport, StoreError> {
2679 let upcaster = crate::journal::current_upcaster();
2680 from_jsonl_with(store, cases, input, upcaster.as_ref()).await
2681}
2682
2683/// [`from_jsonl`], reading each record through `upcaster` rather than the one
2684/// this build ships.
2685///
2686/// A record at a version `upcaster` cannot reach is refused before anything is
2687/// written.
2688///
2689/// # Errors
2690///
2691/// As [`from_jsonl`].
2692pub async fn from_jsonl_with<R: std::io::BufRead>(
2693 store: &Arc<dyn JournalStore>,
2694 cases: Option<&Arc<dyn crate::case::CaseStore>>,
2695 input: R,
2696 upcaster: &dyn crate::journal::Upcaster,
2697) -> Result<RestoreReport, StoreError> {
2698 let parsed = parse(input, upcaster).map_err(|e| StoreError::Backend(e.to_string()))?;
2699 restore_parsed(store, cases, parsed).await
2700}
2701
2702/// An export, restored into a store that lives only as long as the process.
2703///
2704/// The source a strict replay reads when it is handed a file rather than a
2705/// plane: every run the file holds, rebuilt through [`from_jsonl`] and so
2706/// checked against the file's own checkpoint, with nothing written to disk.
2707#[cfg(feature = "redb")]
2708#[derive(Debug)]
2709pub struct ReplaySource {
2710 pub store: Arc<crate::store::RedbStore>,
2711 /// The runs the file holds, in the order it lists them.
2712 pub runs: Vec<RunId>,
2713 pub report: RestoreReport,
2714}
2715
2716/// Restore an export into memory for replay.
2717///
2718/// # Errors
2719///
2720/// If the file cannot be read or is refused by [`from_jsonl`], or if the
2721/// rebuilt store does not commit to the history the file claims.
2722#[cfg(feature = "redb")]
2723pub async fn open_for_replay<R: std::io::BufRead>(input: R) -> Result<ReplaySource, StoreError> {
2724 let upcaster = crate::journal::current_upcaster();
2725 let parsed = parse(input, upcaster.as_ref()).map_err(|e| StoreError::Backend(e.to_string()))?;
2726 let runs = parsed.runs.iter().map(|r| r.run).collect();
2727 let store = Arc::new(crate::store::RedbStore::open_in_memory()?);
2728 let journal = Arc::clone(&store) as Arc<dyn JournalStore>;
2729 let cases = Arc::clone(&store) as Arc<dyn crate::case::CaseStore>;
2730 let report = restore_parsed(&journal, Some(&cases), parsed).await?;
2731 if !report.is_faithful() {
2732 return Err(StoreError::Backend(format!(
2733 "the export does not rebuild to its own checkpoint: it claims {} records under \
2734 root {}, and restoring it produced {} under {}",
2735 report.expected.size, report.expected.root, report.rebuilt.size, report.rebuilt.root
2736 )));
2737 }
2738 Ok(ReplaySource {
2739 store,
2740 runs,
2741 report,
2742 })
2743}
2744
2745async fn restore_parsed(
2746 store: &Arc<dyn JournalStore>,
2747 cases: Option<&Arc<dyn crate::case::CaseStore>>,
2748 parsed: Parsed,
2749) -> Result<RestoreReport, StoreError> {
2750 refuse_before_writing(store, &parsed).await?;
2751
2752 let mut records = 0usize;
2753 for run in &parsed.runs {
2754 let mut claimed = run.hashes.iter();
2755 // Grouped by epoch, in order. Each group is one `append` carrying that
2756 // group's own epoch, which is what reproduces the hashed bodies of a run
2757 // that changed owner mid-flight.
2758 for batch in run.records.chunk_by(|a, b| a.body.epoch == b.body.epoch) {
2759 let Some(epoch) = batch.first().map(|r| r.body.epoch) else {
2760 continue;
2761 };
2762 let appends: Vec<Append> = batch.iter().cloned().map(Append::restored).collect();
2763 records += appends.len();
2764 for (written, want) in store.append(epoch, appends).await?.iter().zip(&mut claimed) {
2765 if written.hash != *want {
2766 return Err(StoreError::Backend(format!(
2767 "run {} seq {} rebuilt to hash {} and the file claims {want} — the \
2768 store is not reproducing the history it was handed, so the restore \
2769 stopped there",
2770 run.run,
2771 written.seq(),
2772 written.hash
2773 )));
2774 }
2775 }
2776 }
2777 }
2778
2779 // Sealed last, and in the log's own order, because that order *is* the
2780 // Merkle log. Sealing as each run finished would rebuild the tree in
2781 // whatever sequence the file happened to list them, and the roots would
2782 // differ for a history that is otherwise identical.
2783 let mut sealed: Vec<&RestoredRun> = parsed
2784 .runs
2785 .iter()
2786 .filter(|r| r.index.is_some())
2787 .collect::<Vec<_>>();
2788 sealed.sort_by_key(|r| r.index);
2789 for run in sealed {
2790 let (Some(outcome), Some(epoch)) = (
2791 run.outcome.as_deref(),
2792 run.records.last().map(|r| r.body.epoch),
2793 ) else {
2794 continue;
2795 };
2796 store.seal(run.run, epoch, outcome).await?;
2797 }
2798
2799 // The case layer, after the journal. Order matters only for the operator's
2800 // mental model — the two halves share no constraint — but the journal is
2801 // the half whose restore can fail on a constraint, and failing before any
2802 // case row landed leaves the cleaner wreck.
2803 let mut imported = 0usize;
2804 let mut not_carried = Vec::new();
2805 match (cases, parsed.cases.is_empty()) {
2806 (Some(case_store), false) => {
2807 for block in &parsed.cases {
2808 case_store
2809 .import_case(&block.case, &block.deadlines, &block.blobs)
2810 .await?;
2811 if let Some(hold) = &block.hold {
2812 case_store.place_hold(block.case.id, hold).await?;
2813 }
2814 imported += 1;
2815 }
2816 not_carried.push(
2817 "blob link timestamps — the export carries a case's blob digests without \
2818 the instant each link was written, so erasure reachability survives and \
2819 the original ordering does not"
2820 .to_owned(),
2821 );
2822 }
2823 (None, false) => not_carried.push(format!(
2824 "the case layer — the export carries {} case(s) and no case store was supplied, \
2825 so the journal is rebuilt and the matters it names are not",
2826 parsed.cases.len()
2827 )),
2828 (_, true) => {}
2829 }
2830
2831 not_carried.extend(losses(&parsed));
2832
2833 let awaiting = awaiting_runs(store, &parsed.runs).await?;
2834
2835 let rebuilt = store.checkpoint().await?;
2836 if rebuilt.origin != parsed.checkpoint.origin {
2837 not_carried.push(format!(
2838 "the log identity — this history was written by '{}' and is now held by '{}'. \
2839 A legitimate recovery: a restore is pointed at a store, and one tenant's \
2840 history put back under another tenant's name is a different log with the \
2841 same contents. It is named because a checkpoint an auditor holds from \
2842 before the disaster will report the new log as the wrong one",
2843 parsed.checkpoint.origin, rebuilt.origin
2844 ));
2845 }
2846
2847 Ok(RestoreReport {
2848 expected: parsed.checkpoint,
2849 rebuilt,
2850 runs: parsed.runs.len(),
2851 records,
2852 cases: imported,
2853 awaiting,
2854 not_carried,
2855 })
2856}
2857
2858/// Every refusal a restore makes before its first write, so a refused restore
2859/// leaves nothing to clean up.
2860async fn refuse_before_writing(
2861 store: &Arc<dyn JournalStore>,
2862 parsed: &Parsed,
2863) -> Result<(), StoreError> {
2864 // A check `parse` leaves to its caller, because it is not about any one
2865 // line: a format this build does not read cannot be
2866 // *parsed* completely, and `parse` skips what it does not recognise — so
2867 // proceeding would restore whatever subset happened to look familiar and
2868 // report it as the whole file.
2869 if parsed.version != Some(u64::from(FORMAT_VERSION)) {
2870 return Err(StoreError::Backend(format!(
2871 "the export claims format version {:?} and this build reads {FORMAT_VERSION} — \
2872 restoring a format this build cannot fully parse would rebuild an unknowable \
2873 subset and call it a history",
2874 parsed.version
2875 )));
2876 }
2877
2878 if let Some(why) = canon_unverifiable(parsed.canon) {
2879 return Err(StoreError::Backend(why));
2880 }
2881
2882 // The frame is the completeness signal, and the restore is the reader most
2883 // exposed to its absence: a truncated export is a *prefix* in which every
2884 // line is valid, so replaying one rebuilds a partial history shaped
2885 // exactly like a whole one. The quietest cut is the worst — a file cut
2886 // after the last record but before the case layer restores a journal that
2887 // is byte-perfect and `is_faithful`, with every matter it names missing.
2888 // Refused before any write lands, so a refused restore leaves nothing to
2889 // clean up. What this does NOT cover: a file truncated *and* given a
2890 // forged trailer — that is `verify`'s count settlement, and the right
2891 // order is restore, then verify.
2892 if !parsed.complete {
2893 return Err(StoreError::Backend(
2894 "the export has no trailer, so it was cut short — every line in it is a valid \
2895 prefix, and restoring a prefix would rebuild a partial history shaped exactly \
2896 like a whole one. Re-take the export"
2897 .to_owned(),
2898 ));
2899 }
2900
2901 if store.seals() {
2902 return Err(StoreError::Backend(
2903 "the target store seals payloads as it writes them, and a restore must write \
2904 each record exactly as it was recorded — sealed payloads stay sealed under \
2905 their original keys. Restore into the unwrapped store, then open it with the \
2906 keyring"
2907 .to_owned(),
2908 ));
2909 }
2910
2911 refuse_foreign_seals(parsed, store.tenant())?;
2912 // Every run before the first write: the store refuses a held run's records
2913 // on its own, but only once the runs ahead of it in the file are written.
2914 for run in &parsed.runs {
2915 if store.head(run.run).await?.seq != 0 {
2916 return Err(StoreError::Backend(format!(
2917 "this store already holds run {} — a restore rebuilds a history into an \
2918 empty store, it does not merge one, so nothing was written",
2919 run.run
2920 )));
2921 }
2922 }
2923 Ok(())
2924}
2925
2926/// Refuse a file whose sealed payloads or case states name another tenant.
2927///
2928/// An envelope's associated data and erasure scope both name the tenant that
2929/// sealed it, so under another tenant it opens for nobody and `erase_case`
2930/// there destroys a key that wraps none of it. Moving sealed history between
2931/// tenants needs a re-seal, which a restore cannot do.
2932fn refuse_foreign_seals(parsed: &Parsed, tenant: &str) -> Result<(), StoreError> {
2933 use crate::journal::payload::{self, SealedField};
2934
2935 let foreign = |envelope: Option<Vec<u8>>, whose: &dyn std::fmt::Display| match envelope
2936 .as_deref()
2937 .and_then(payload::sealed_tenant)
2938 {
2939 Some(sealer) if sealer == tenant => Ok(()),
2940 sealer => Err(StoreError::Backend(format!(
2941 "{whose} carries a payload sealed for tenant '{}', and this store serves \
2942 '{tenant}' — under another tenant it opens for nobody and an erasure \
2943 there would not reach it, so the file is refused before anything is written",
2944 sealer.as_deref().unwrap_or("(unreadable envelope)")
2945 ))),
2946 };
2947 for run in &parsed.runs {
2948 for record in &run.records {
2949 let mut kind = record.body.kind.clone();
2950 for field in payload::payloads(&mut kind) {
2951 match field {
2952 SealedField::Value(v) if payload::is_sealed(v) => {
2953 foreign(payload::unwrap(v), &format!("run {}", run.run))?;
2954 }
2955 SealedField::Text(t) if payload::is_sealed_text(t) => {
2956 foreign(payload::unwrap_text(t), &format!("run {}", run.run))?;
2957 }
2958 _ => {}
2959 }
2960 }
2961 }
2962 }
2963 for block in &parsed.cases {
2964 if payload::is_sealed(&block.case.state) {
2965 foreign(
2966 payload::unwrap(&block.case.state),
2967 &format!("case {}", block.case.id),
2968 )?;
2969 }
2970 }
2971 Ok(())
2972}
2973
2974/// One run, as an export describes it.
2975struct RestoredRun {
2976 run: RunId,
2977 /// Position in the Merkle log; `None` for a run that was still open.
2978 index: Option<u64>,
2979 outcome: Option<String>,
2980 /// Each record as verified, carrying the bytes the restore writes back.
2981 records: Vec<crate::journal::Record>,
2982 /// Each body's hash as the file claims it, which the rebuilt record must
2983 /// reproduce.
2984 hashes: Vec<crate::core::Digest>,
2985 /// The chain head over the records read so far, which the next record's
2986 /// hash must extend.
2987 prev: crate::core::Digest,
2988}
2989
2990/// One record line, once its bytes are known to be the ones the chain
2991/// committed to and at a version `upcaster` reads.
2992///
2993/// Both are conditions of replay rather than verification. Bytes the claimed
2994/// hash does not cover would rebuild a history the file never committed to,
2995/// and a record at a version no upcaster reaches has no body the store can
2996/// index — either way into a populated store, before the checkpoint comparison
2997/// could say so.
2998fn replayable(
2999 raw: &[u8],
3000 prev: crate::core::Digest,
3001 claimed: crate::core::Digest,
3002 upcaster: &dyn crate::journal::Upcaster,
3003) -> Result<crate::journal::Record, std::io::Error> {
3004 // A store keeps these bytes as they stand, so bytes that hash correctly
3005 // and are not canonical would land a record no writer under this canon
3006 // produces. Compared as values rather than as this build's record shape,
3007 // so a record from another shape is held to the same rule.
3008 if serde_json::from_slice::<serde_json::Value>(raw)
3009 .is_ok_and(|value| crate::core::canon::value_bytes(&value) != raw)
3010 {
3011 return Err(std::io::Error::other(
3012 "a record's wire bytes are not canonical — no writer under the export's canon \
3013 produces them, and a store would keep them as they stand, so the file is \
3014 refused before anything is written",
3015 ));
3016 }
3017 crate::journal::Record::from_stored_with(upcaster, raw.to_vec(), prev, claimed, None).map_err(
3018 |e| {
3019 std::io::Error::other(match e {
3020 StoreError::Corrupt { .. } => format!(
3021 "a record's claimed hash does not cover its wire bytes and the chain \
3022 before it ({e}) — replaying it would rebuild a history the export never \
3023 committed to, so the file is refused before anything is written"
3024 ),
3025 StoreError::UnknownRecordVersion { .. } => format!(
3026 "a record is at a version this build does not restore ({e}) — no \
3027 upcaster reaches it, so the store could not index it as written, and \
3028 the file is refused before anything is written"
3029 ),
3030 other => format!(
3031 "a record line's wire bytes do not parse ({other}) — the record cannot be \
3032 replayed as written, and its display copy is not a substitute"
3033 ),
3034 })
3035 },
3036 )
3037}
3038
3039struct Parsed {
3040 checkpoint: Checkpoint,
3041 /// The format version the header claims, `None` when there was no header.
3042 version: Option<u64>,
3043 /// The canonicalization rule the header names, `None` when absent.
3044 canon: Option<u64>,
3045 /// Whether the file ended with its trailer. A truncated export is a valid
3046 /// prefix, and a restore must refuse it — see [`from_jsonl`].
3047 complete: bool,
3048 runs: Vec<RestoredRun>,
3049 cases: Vec<RestoredCase>,
3050 signed: usize,
3051}
3052
3053/// One case block, as a restore replays it.
3054struct RestoredCase {
3055 case: crate::core::Case,
3056 deadlines: Vec<crate::core::Deadline>,
3057 blobs: Vec<crate::core::Digest>,
3058 hold: Option<crate::core::LegalHold>,
3059}
3060
3061/// A case block's hold: `null` for none, a [`LegalHold`](crate::core::LegalHold)
3062/// otherwise, and an error when the member is missing or unreadable.
3063///
3064/// Shared by the verifier and the restore so the two cannot disagree about
3065/// what a readable hold is.
3066fn case_hold(value: &serde_json::Value) -> Result<Option<crate::core::LegalHold>, String> {
3067 let Some(hold) = value.get("hold") else {
3068 return Err(
3069 "a case block carries no `hold` member, so whether the matter is under \
3070 a legal hold is unknown"
3071 .to_owned(),
3072 );
3073 };
3074 serde_json::from_value::<Option<crate::core::LegalHold>>(hold.clone())
3075 .map_err(|e| format!("a case block's legal hold is malformed: {e}"))
3076}
3077
3078/// One case block as a restore replays it, or `None` for a malformed block.
3079///
3080/// Malformed blocks are skipped and found by `verify`, per [`parse`]'s
3081/// no-checking rule — except an unreadable hold, which refuses the file.
3082fn restored_case(value: &serde_json::Value) -> Result<Option<RestoredCase>, std::io::Error> {
3083 use serde_json::Value;
3084
3085 // A restore that refused the file would refuse the healthy cases too.
3086 let (Ok(case), Some(deadlines), Some(blobs)) = (
3087 serde_json::from_value::<crate::core::Case>(
3088 value.get("case").cloned().unwrap_or(Value::Null),
3089 ),
3090 value
3091 .get("deadlines")
3092 .and_then(|d| serde_json::from_value::<Vec<crate::core::Deadline>>(d.clone()).ok()),
3093 value
3094 .get("blobs")
3095 .and_then(|b| serde_json::from_value::<Vec<crate::core::Digest>>(b.clone()).ok()),
3096 ) else {
3097 return Ok(None);
3098 };
3099 // Refused before anything is written.
3100 let hold = case_hold(value).map_err(|e| {
3101 std::io::Error::other(format!(
3102 "case {}: {e} — restoring the matter without it would let retention \
3103 erase it, so the file is refused",
3104 case.id
3105 ))
3106 })?;
3107 Ok(Some(RestoredCase {
3108 case,
3109 deadlines,
3110 blobs,
3111 hold,
3112 }))
3113}
3114
3115/// Read an export into the shape a restore replays.
3116///
3117/// Checks what replay needs and nothing wider: [`verify`] answers *is this
3118/// sound* and this answers *what does it say*. Folding them would make a
3119/// restore refuse the very history an operator is trying to recover, at the
3120/// moment they most need it — and the right order is restore, then verify the
3121/// result against its own checkpoint, which [`from_jsonl`] reports.
3122///
3123/// One class of line is a hard error rather than a skip: a record line that
3124/// cannot be *replayed as written*. Its wire bytes must be present and parse —
3125/// the one available guess otherwise, the editable display copy, is exactly
3126/// the value the wire-bytes rule exists to keep out of the rebuilt history —
3127/// and they must be the bytes its hash covers, at the version this build
3128/// writes; see [`replayable`]. All of it is decided here, before
3129/// [`from_jsonl`] writes anything.
3130fn parse<R: std::io::BufRead>(
3131 input: R,
3132 upcaster: &dyn crate::journal::Upcaster,
3133) -> Result<Parsed, std::io::Error> {
3134 use serde_json::Value;
3135
3136 let mut parsed = Parsed {
3137 checkpoint: Checkpoint {
3138 origin: String::new(),
3139 size: 0,
3140 root: crate::core::Digest::ZERO,
3141 },
3142 version: None,
3143 canon: None,
3144 complete: false,
3145 runs: Vec::new(),
3146 cases: Vec::new(),
3147 signed: 0,
3148 };
3149 for line in input.lines() {
3150 let line = line?;
3151 let Ok(value) = serde_json::from_str::<Value>(&line) else {
3152 continue;
3153 };
3154 match value.get("kind").and_then(Value::as_str) {
3155 Some("agentplane.export") => {
3156 parsed.version = value.get("version").and_then(Value::as_u64);
3157 parsed.canon = value.get("canon").and_then(Value::as_u64);
3158 if let Some(c) = value
3159 .get("checkpoint")
3160 .and_then(|c| serde_json::from_value::<Checkpoint>(c.clone()).ok())
3161 {
3162 parsed.checkpoint = c;
3163 }
3164 }
3165 Some("agentplane.export.run") => {
3166 if let Some(run) = value
3167 .get("run")
3168 .and_then(Value::as_str)
3169 .and_then(|s| RunId::parse(s).ok())
3170 {
3171 parsed.runs.push(RestoredRun {
3172 run,
3173 index: value.get("index").and_then(Value::as_u64),
3174 outcome: None,
3175 records: Vec::new(),
3176 hashes: Vec::new(),
3177 prev: crate::core::Digest::ZERO,
3178 });
3179 }
3180 }
3181 Some("agentplane.export.case") => {
3182 if let Some(case) = restored_case(&value)? {
3183 parsed.cases.push(case);
3184 }
3185 }
3186 Some("agentplane.export.end") => parsed.complete = true,
3187 Some(DISCLOSURE_KIND) => return Err(std::io::Error::other(PACKAGE_REFUSED)),
3188 // A line carrying a `kind` this build does not recognise is not a
3189 // record — record lines are the only unkinded lines in the format
3190 // — so it is skipped per the no-checking rule rather than held to
3191 // a record's obligations.
3192 Some(_) => {}
3193 _ => {
3194 if value.get("signature").is_some_and(|a| !a.is_null()) {
3195 parsed.signed += 1;
3196 }
3197 // The wire bytes are the source of truth, exactly as they are
3198 // for the verifier: the readable `body` is a courtesy copy,
3199 // and a restore replaying the copy would rebuild whatever the
3200 // display half said rather than what the chain covered. There
3201 // is deliberately **no fallback to that copy**: a record line
3202 // with no `raw`, or whose `raw` does not parse, is a hard
3203 // error rather than a skip or a guess — silently substituting
3204 // the one editable value two mechanisms must agree about would
3205 // rebuild a history the chain never hashed and let the
3206 // subsequent verify pass bless it.
3207 let Some(raw) = value.get("raw").and_then(Value::as_str) else {
3208 return Err(std::io::Error::other(
3209 "a record line carries no wire bytes (`raw`) — restoring its display \
3210 copy instead would rebuild what the readable half says rather than \
3211 what the chain hashed, so the file is refused instead of guessed at",
3212 ));
3213 };
3214 let Some(claimed) = value
3215 .get("hash")
3216 .and_then(|h| serde_json::from_value::<crate::core::Digest>(h.clone()).ok())
3217 else {
3218 return Err(std::io::Error::other(
3219 "a record line carries no hash — nothing ties its bytes to the chain, \
3220 so the file is refused instead of replayed",
3221 ));
3222 };
3223 let Some(current) = parsed.runs.last_mut() else {
3224 continue;
3225 };
3226 let record = replayable(raw.as_bytes(), current.prev, claimed, upcaster)?;
3227 current.prev = claimed;
3228 if let crate::journal::RecordKind::RunConcluded { outcome, .. } = &record.body.kind
3229 {
3230 current.outcome = Some(outcome.clone());
3231 }
3232 current.records.push(record);
3233 current.hashes.push(claimed);
3234 }
3235 }
3236 }
3237 Ok(parsed)
3238}