use std::sync::Arc;
use crate::core::{RunId, StoreError};
use crate::journal::{Append, Checkpoint, JournalStore};
pub const FORMAT_VERSION: u32 = 1;
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct Header {
pub kind: &'static str,
pub version: u32,
pub checkpoint: Checkpoint,
pub canon: u16,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize)]
pub struct ExportedRecord<'r> {
pub seq: crate::core::Seq,
pub body: crate::journal::RecordBody,
pub prev_hash: &'r crate::core::Digest,
pub hash: &'r crate::core::Digest,
pub signature: Option<&'r crate::core::KeySignature>,
pub raw: std::borrow::Cow<'r, str>,
}
impl<'r> ExportedRecord<'r> {
fn from_stored(r: &'r crate::journal::Record) -> Result<Self, String> {
let body = serde_json::from_slice::<crate::journal::RecordBody>(r.raw())
.map_err(|e| format!("record {}'s wire bytes do not parse: {e}", r.seq()))?;
Ok(Self {
seq: r.seq(),
body,
prev_hash: &r.prev_hash,
hash: &r.hash,
signature: r.signature.as_ref(),
raw: String::from_utf8_lossy(r.raw()),
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct RunBlock {
pub kind: &'static str,
pub run: RunId,
#[serde(skip_serializing_if = "Option::is_none")]
pub index: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub seal: Option<crate::core::Digest>,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize)]
pub struct CaseBlock {
pub kind: &'static str,
pub case: crate::core::Case,
pub deadlines: Vec<crate::core::Deadline>,
pub blobs: Vec<crate::core::Digest>,
pub hold: Option<crate::core::LegalHold>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct Trailer {
pub kind: &'static str,
pub runs_requested: usize,
pub runs_exported: usize,
pub records: usize,
pub cases: usize,
pub unreadable: Vec<Unreadable>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct Unreadable {
pub run: RunId,
pub reason: String,
}
pub async fn to_jsonl<W: std::io::Write>(
store: &Arc<dyn JournalStore>,
cases: Option<&Arc<dyn crate::case::CaseStore>>,
runs: &[RunId],
mut out: W,
) -> Result<Trailer, std::io::Error> {
let checkpoint = store.checkpoint().await.map_err(|e| as_io(&e))?;
let log_size = checkpoint.size;
let header = Header {
kind: "agentplane.export",
version: FORMAT_VERSION,
checkpoint,
canon: crate::core::canon::VERSION,
};
writeln!(out, "{}", to_line(&header)?)?;
let mut records = 0usize;
let mut exported = 0usize;
let mut unreadable = Vec::new();
for &run in runs {
let placed = store.inclusion_proof(run).await.ok().flatten();
let placed = placed.filter(|i| i.index < log_size);
writeln!(
out,
"{}",
to_line(&RunBlock {
kind: "agentplane.export.run",
run,
index: placed.as_ref().map(|i| i.index),
seal: placed.as_ref().map(|i| i.seal),
})?
)?;
match store.read(run, 1).await {
Ok(found) if found.is_empty() => unreadable.push(Unreadable {
run,
reason: "the store holds no records for this run".to_owned(),
}),
Ok(found) => {
match found
.iter()
.map(ExportedRecord::from_stored)
.collect::<Result<Vec<_>, _>>()
{
Ok(lines) => {
for line in &lines {
writeln!(out, "{}", to_line(line)?)?;
records += 1;
}
exported += 1;
}
Err(reason) => unreadable.push(Unreadable { run, reason }),
}
}
Err(e) => unreadable.push(Unreadable {
run,
reason: e.to_string(),
}),
}
}
let mut case_count = 0usize;
if let Some(case_store) = cases {
let mut after: Option<crate::core::CaseId> = None;
loop {
let page = case_store
.cases(after, CASE_PAGE)
.await
.map_err(|e| as_io(&e))?;
let Some(last) = page.last() else { break };
after = Some(last.id);
let full = page.len() >= CASE_PAGE;
for case in page {
let deadlines = case_store.deadlines(case.id).await.map_err(|e| as_io(&e))?;
let blobs = case_store.blobs_of(case.id).await.map_err(|e| as_io(&e))?;
let hold = case_store.hold(case.id).await.map_err(|e| as_io(&e))?;
writeln!(
out,
"{}",
to_line(&CaseBlock {
kind: "agentplane.export.case",
case,
deadlines,
blobs,
hold,
})?
)?;
case_count += 1;
}
if !full {
break;
}
}
}
let trailer = Trailer {
kind: "agentplane.export.end",
runs_requested: runs.len(),
runs_exported: exported,
records,
cases: case_count,
unreadable,
};
writeln!(out, "{}", to_line(&trailer)?)?;
out.flush()?;
Ok(trailer)
}
pub(crate) const CASE_PAGE: usize = 256;
fn to_line<T: serde::Serialize>(value: &T) -> Result<String, std::io::Error> {
serde_json::to_string(value).map_err(|e| std::io::Error::other(e.to_string()))
}
fn as_io(e: &StoreError) -> std::io::Error {
std::io::Error::other(e.to_string())
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct VerifyReport {
pub checkpoint: Checkpoint,
pub sound: Vec<RunId>,
pub findings: Vec<String>,
pub not_checked: Vec<String>,
pub records: usize,
pub cases: usize,
pub complete: bool,
}
impl VerifyReport {
#[must_use]
pub fn is_sound(&self) -> bool {
self.findings.is_empty() && self.complete
}
}
pub fn verify<R: std::io::BufRead>(
input: R,
verifier: Option<&dyn crate::core::Verifier>,
anchors: &[crate::journal::Anchor],
) -> Result<VerifyReport, std::io::Error> {
use crate::core::Digest;
use serde_json::Value;
let mut report = VerifyReport {
checkpoint: Checkpoint {
origin: String::new(),
size: 0,
root: Digest::ZERO,
},
sound: Vec::new(),
findings: Vec::new(),
not_checked: Vec::new(),
records: 0,
cases: 0,
complete: false,
};
unanswerable(&mut report, verifier.is_some());
let mut header_seen = false;
let mut leaves: Vec<(u64, crate::core::merkle::LeafHash)> = Vec::new();
let mut pass: Option<RunPass> = None;
let mut run_blocks = 0usize;
let mut read_runs = 0usize;
let mut empty_blocks: Vec<RunId> = Vec::new();
let mut claims = TrailerClaims::default();
let mut stamped: std::collections::BTreeSet<crate::core::CaseId> =
std::collections::BTreeSet::new();
let mut carried: std::collections::BTreeSet<crate::core::CaseId> =
std::collections::BTreeSet::new();
let mut blob_digests = 0usize;
for line in input.lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let Ok(value) = serde_json::from_str::<Value>(&line) else {
report
.findings
.push("a line is not valid JSON, so the export is unreadable from there on".into());
break;
};
if let Some(kind) = value.get("kind").and_then(Value::as_str) {
note_unknown_members(kind, &value, &mut report);
}
match value.get("kind").and_then(Value::as_str) {
Some("agentplane.export") => {
header_seen = true;
read_header(&value, &mut report);
}
Some("agentplane.export.case") => {
read_case_block(&value, &mut report, &mut carried, &mut blob_digests);
}
Some("agentplane.export.run") => {
run_blocks += 1;
finish_run(
&mut report,
pass.take(),
verifier,
&mut read_runs,
&mut empty_blocks,
);
pass = open_run_block(&value, &mut leaves);
}
Some("agentplane.export.end") => read_trailer(&value, &mut report, &mut claims),
_ => {
report.records += 1;
let Some(pass) = pass.as_mut() else {
report.findings.push(
"a record appears before any run block, so nothing says which run it \
belongs to"
.into(),
);
continue;
};
read_record(&value, pass, &mut report, &mut stamped);
}
}
}
finish_run(
&mut report,
pass,
verifier,
&mut read_runs,
&mut empty_blocks,
);
let open_runs = run_blocks.saturating_sub(leaves.len());
settle(&mut report, header_seen, leaves, anchors);
settle_trailer(
&mut report,
&claims,
run_blocks,
read_runs,
open_runs,
&empty_blocks,
);
settle_cases(&mut report, &stamped, &carried, blob_digests);
Ok(report)
}
#[derive(Default)]
struct TrailerClaims {
runs_requested: Option<u64>,
runs_exported: Option<u64>,
records: Option<u64>,
unreadable: Vec<(RunId, String)>,
}
fn read_trailer(value: &serde_json::Value, report: &mut VerifyReport, claims: &mut TrailerClaims) {
report.complete = true;
if let Some(declared) = value.get("cases").and_then(serde_json::Value::as_u64)
&& declared != report.cases as u64
{
report.findings.push(format!(
"the trailer says {declared} case(s) were exported and this file \
carries {} — the case layer was cut after the export was taken",
report.cases
));
}
claims.runs_requested = value
.get("runs_requested")
.and_then(serde_json::Value::as_u64);
claims.runs_exported = value
.get("runs_exported")
.and_then(serde_json::Value::as_u64);
claims.records = value.get("records").and_then(serde_json::Value::as_u64);
if let Some(list) = value
.get("unreadable")
.and_then(serde_json::Value::as_array)
{
for entry in list {
let Some(run) = entry
.get("run")
.and_then(serde_json::Value::as_str)
.and_then(|s| RunId::parse(s).ok())
else {
continue;
};
let reason = entry
.get("reason")
.and_then(serde_json::Value::as_str)
.unwrap_or("no reason recorded")
.to_owned();
claims.unreadable.push((run, reason));
}
}
}
fn settle_trailer(
report: &mut VerifyReport,
claims: &TrailerClaims,
run_blocks: usize,
read_runs: usize,
open_runs: usize,
empty_blocks: &[RunId],
) {
if open_runs > 0 {
report.not_checked.push(format!(
"{open_runs} open run(s): a run that has not concluded has no position in the \
Merkle log, so the root proves nothing about it — its chain and signatures \
were verified, and records cut from its tail before the export was taken are \
undetectable from this file"
));
}
for (run, reason) in &claims.unreadable {
report.not_checked.push(format!(
"run {run}: the export declares it unreadable ({reason}), so its records are \
not in this file and nothing about it was verified"
));
}
for run in empty_blocks {
if claims.unreadable.iter().any(|(u, _)| u == run) {
continue;
}
report.findings.push(format!(
"run {run}: its block carries no records and the export does not declare it \
unreadable — either the records were removed after the export was taken, or \
the file was cut short before them"
));
}
if !report.complete {
return;
}
match (claims.runs_requested, claims.runs_exported, claims.records) {
(Some(requested), Some(exported), Some(records)) => {
if requested != run_blocks as u64 {
report.findings.push(format!(
"the trailer says {requested} run(s) were requested and this file carries \
{run_blocks} run block(s) — whole runs were removed or added after the \
export was taken"
));
}
if exported != read_runs as u64 {
report.findings.push(format!(
"the trailer says {exported} run(s) were exported in full and this file \
carries records for {read_runs} — a run's records were removed after the \
export was taken"
));
}
if records != report.records as u64 {
report.findings.push(format!(
"the trailer says {records} record(s) were written and this file carries \
{} — record lines were removed or added after the export was taken",
report.records
));
}
}
_ => report.findings.push(
"the trailer is missing counts this format always writes (runs_requested, \
runs_exported, records) — a reader cannot hold the file to its own accounting"
.to_owned(),
),
}
}
fn read_case_block(
value: &serde_json::Value,
report: &mut VerifyReport,
carried: &mut std::collections::BTreeSet<crate::core::CaseId>,
blob_digests: &mut usize,
) {
use serde_json::Value;
report.cases += 1;
match serde_json::from_value::<crate::core::Case>(
value.get("case").cloned().unwrap_or(Value::Null),
) {
Ok(case) => {
carried.insert(case.id);
}
Err(e) => report
.findings
.push(format!("a case block is malformed: {e}")),
}
if value
.get("deadlines")
.is_none_or(|d| serde_json::from_value::<Vec<crate::core::Deadline>>(d.clone()).is_err())
{
report
.findings
.push("a case block's deadlines are malformed".to_owned());
}
if let Err(e) = case_hold(value) {
report.findings.push(e);
}
*blob_digests += value
.get("blobs")
.and_then(Value::as_array)
.map_or(0, Vec::len);
}
fn unanswerable(report: &mut VerifyReport, verifier_supplied: bool) {
if !verifier_supplied {
report.not_checked.push(
"signatures — no public key was supplied, so this pass cannot say who wrote \
anything"
.to_owned(),
);
}
report
.not_checked
.extend(UNCARRIED.iter().map(|limit| (*limit).to_owned()));
}
const UNCARRIED: [&str; 2] = [
"webhook delivery cursors — this file carries no push registrations, so a restored plane \
re-delivers from the start of each subscriber's history rather than from where it got \
to. Receivers deduplicate on the event's own identity, so the cost is repetition rather \
than loss",
"worklist decisions no run has consumed — a decision recorded against a task and not yet \
read back by the run it answers is a store row, not a record, so it does not survive \
here. The task re-opens and is decided again under the same four-eyes and expiry",
];
fn settle_cases(
report: &mut VerifyReport,
stamped: &std::collections::BTreeSet<crate::core::CaseId>,
carried: &std::collections::BTreeSet<crate::core::CaseId>,
blob_digests: usize,
) {
if carried.is_empty() {
if !stamped.is_empty() {
report.not_checked.push(format!(
"the case layer — {} case(s) are stamped on records and this file carries no \
case blocks, so either the plane's case store was not supplied to the export \
or the layer was dropped; the two cannot be told apart from the file alone",
stamped.len()
));
}
return;
}
for case in stamped.difference(carried) {
report.findings.push(format!(
"case {case} is stamped on exported records and missing from the case layer — \
the journal names a matter this file does not carry"
));
}
if blob_digests > 0 {
report.not_checked.push(format!(
"blob bytes — the case layer references {blob_digests} blob digest(s) and this \
file carries digests, not bytes; presence and integrity are a question about a \
live blob store"
));
}
report.not_checked.push(
"sealed-state keys — whether sealed case state can still be opened is a question \
about a live key ring, which an offline file cannot answer"
.to_owned(),
);
}
fn open_run_block(
value: &serde_json::Value,
leaves: &mut Vec<(u64, crate::core::merkle::LeafHash)>,
) -> Option<RunPass> {
use crate::core::{Digest, merkle};
use serde_json::Value;
let pass = value
.get("run")
.and_then(Value::as_str)
.and_then(|s| RunId::parse(s).ok())
.map(|run| RunPass {
run,
declared_seal: value
.get("seal")
.and_then(|s| serde_json::from_value::<Digest>(s.clone()).ok()),
prev: Digest::ZERO,
last_seq: 0,
records: 0,
resealed: Vec::new(),
clean: true,
});
if let Some(pass) = &pass
&& let (Some(index), Some(seal)) = (
value.get("index").and_then(Value::as_u64),
pass.declared_seal,
)
{
leaves.push((index, merkle::leaf_hash(&seal)));
}
pass
}
struct RunPass {
run: RunId,
declared_seal: Option<crate::core::Digest>,
prev: crate::core::Digest,
last_seq: u64,
records: usize,
resealed: Vec<crate::journal::Record>,
clean: bool,
}
fn compare_anchors(
report: &mut VerifyReport,
header_seen: bool,
anchors: &[crate::journal::Anchor],
leaves: &[(u64, crate::core::merkle::LeafHash)],
) -> Option<Checkpoint> {
let mut matched = None;
if !header_seen {
return matched;
}
for anchor in anchors {
let given = &anchor.checkpoint;
if given.origin != report.checkpoint.origin || given.size > report.checkpoint.size {
report.findings.push(format!(
"the export's header names log '{}' at size {} with root {}, and the \
checkpoint held by {} names '{}' at size {} with root {} — the file \
describes a different history than the one it is being checked against",
report.checkpoint.origin,
report.checkpoint.size,
report.checkpoint.root.to_hex(),
anchor.obtained_from,
given.origin,
given.size,
given.root.to_hex(),
));
} else if given.size == report.checkpoint.size {
if given.root == report.checkpoint.root {
if matched.is_none() {
matched = Some(given.clone());
}
} else {
report.findings.push(format!(
"the export's header names log '{}' at size {} with root {}, and the \
checkpoint held by {} holds that same size with root {} — one tree of \
a given size has one root, so these are two histories",
report.checkpoint.origin,
report.checkpoint.size,
report.checkpoint.root.to_hex(),
anchor.obtained_from,
given.root.to_hex(),
));
}
} else if prefix_root(leaves, given.size) != Some(given.root) {
report.findings.push(format!(
"the checkpoint held by {} commits to the first {} run(s) of log '{}' with \
root {}, and this export's first {} run(s) do not rebuild to it — a run \
inside that prefix was removed, replaced or moved",
anchor.obtained_from,
given.size,
given.origin,
given.root.to_hex(),
given.size,
));
} else {
report.not_checked.push(format!(
"the checkpoint held by {} is at size {} and matches this export's first {} \
run(s); the {} after it are held only to the file's own header",
anchor.obtained_from,
given.size,
given.size,
report.checkpoint.size - given.size
));
}
}
matched
}
fn prefix_root(
leaves: &[(u64, crate::core::merkle::LeafHash)],
size: u64,
) -> Option<crate::core::Digest> {
let size = usize::try_from(size).ok()?;
let prefix = leaves.get(..size)?;
prefix
.iter()
.enumerate()
.all(|(at, (index, _))| u64::try_from(at) == Ok(*index))
.then(|| {
crate::core::merkle::root(&prefix.iter().map(|(_, leaf)| *leaf).collect::<Vec<_>>())
})
}
fn settle(
report: &mut VerifyReport,
header_seen: bool,
mut leaves: Vec<(u64, crate::core::merkle::LeafHash)>,
anchors: &[crate::journal::Anchor],
) {
use crate::core::merkle;
leaves.sort_by_key(|(index, _)| *index);
let matched = compare_anchors(report, header_seen, anchors, &leaves);
let against = if let Some(checkpoint) = matched {
checkpoint
} else {
if anchors.is_empty() {
report.not_checked.push(
"deletion — no checkpoint was supplied, so the Merkle root could only be \
rebuilt and compared against this file's own header. That proves the \
file is internally consistent, which is also what an editor who dropped \
a run and rewrote the header achieves. Pass the checkpoint an earlier \
audit printed, or one a witness cosigned"
.to_owned(),
);
}
report.checkpoint.clone()
};
if !header_seen {
report
.findings
.push("the export has no header, so nothing says which log it came from".into());
}
let size = u64::try_from(leaves.len()).unwrap_or(u64::MAX);
if size == against.size {
let contiguous = leaves
.iter()
.enumerate()
.all(|(at, (index, _))| u64::try_from(at) == Ok(*index));
if contiguous {
let rebuilt =
merkle::root(&leaves.into_iter().map(|(_, leaf)| leaf).collect::<Vec<_>>());
if rebuilt != against.root {
report.findings.push(
"the Merkle root rebuilt from this export does not match the checkpoint it \
claims to be a copy of"
.to_owned(),
);
}
} else {
report.findings.push(format!(
"the run blocks' log positions are not the contiguous 0..{} the checkpoint \
commits to — a position is duplicated or missing, so this file describes a \
different log than the one it names",
against.size
));
}
} else {
report.findings.push(format!(
"the export carries {size} sealed run(s) and its checkpoint commits to {} — the \
difference is runs that were in the log and are not in this file",
against.size
));
}
if !report.complete {
report.findings.push(
"the export has no trailer, so it was cut short — every line in it is still valid, \
which is why the frame is the signal"
.to_owned(),
);
}
}
fn read_header(value: &serde_json::Value, report: &mut VerifyReport) {
let version = value.get("version").and_then(serde_json::Value::as_u64);
if version != Some(u64::from(FORMAT_VERSION)) {
report.findings.push(format!(
"the export claims format version {version:?} and this build reads {FORMAT_VERSION} \
— the findings below describe the lines this build could interpret, which may not \
be all of them"
));
}
let canon = value.get("canon").and_then(serde_json::Value::as_u64);
if canon != Some(u64::from(crate::core::canon::VERSION)) {
report.not_checked.push(format!(
"derived digests — the export was written under canonicalization rule {canon:?} \
and this build implements {}. The chain, leaf, root and signature checks still \
ran and still hold (they hash the bytes as written, never a re-serialization); \
what this build cannot do is re-derive the digests inside the bodies, such as \
effect keys, under the rule that produced them",
crate::core::canon::VERSION
));
}
match value
.get("checkpoint")
.and_then(|c| serde_json::from_value::<Checkpoint>(c.clone()).ok())
{
Some(c) => report.checkpoint = c,
None => report
.findings
.push("the header carries no readable checkpoint".to_owned()),
}
}
fn read_record(
value: &serde_json::Value,
pass: &mut RunPass,
report: &mut VerifyReport,
stamped: &mut std::collections::BTreeSet<crate::core::CaseId>,
) {
let current = pass.run;
pass.records += 1;
let (Some(raw), Some(claimed)) = (
value.get("raw").and_then(serde_json::Value::as_str),
value
.get("hash")
.and_then(|h| serde_json::from_value::<crate::core::Digest>(h.clone()).ok()),
) else {
report.findings.push(format!(
"run {current}: a record line carries no wire bytes or no hash — nothing ties \
it to the chain"
));
pass.clean = false;
return;
};
let raw_bytes = raw.as_bytes();
if crate::core::Digest::chain(pass.prev, raw_bytes) != claimed {
return edited_record(raw_bytes, pass, report);
}
let body = match serde_json::from_slice::<crate::journal::RecordBody>(raw_bytes) {
Ok(body) => body,
Err(parse) => return unparsed_record(raw_bytes, parse, pass, report),
};
let wire: serde_json::Value = serde_json::from_slice(raw_bytes).unwrap_or_default();
if value.get("body") != Some(&wire) {
report.findings.push(format!(
"run {current}: record {}'s readable body does not match its wire bytes — the \
display copy was edited, and every hash still verifies over the real one",
body.seq
));
pass.clean = false;
}
if let Some(case) = body.case {
stamped.insert(case);
}
if body.run != current {
report.findings.push(format!(
"run {current}: a record in this block belongs to run {} — the block was relabelled, \
or spliced from another history",
body.run
));
pass.clean = false;
}
if body.seq != pass.last_seq + 1 {
report.findings.push(format!(
"run {current}: seq {} follows {}, so a record is missing from the middle — \
every chain link either side of the gap still joins",
body.seq, pass.last_seq
));
pass.clean = false;
}
pass.last_seq = body.seq;
if let crate::journal::RecordKind::RunConcluded { chain_head, .. } = &body.kind
&& *chain_head != pass.prev
{
report.findings.push(format!(
"run {current}: the sealing record claims a chain head that is not the head it \
sits on — the conclusion was drawn over a different history"
));
pass.clean = false;
}
let signature = value
.get("signature")
.and_then(|a| serde_json::from_value::<Option<crate::core::KeySignature>>(a.clone()).ok())
.flatten();
match crate::journal::Record::from_stored_signed(
raw_bytes.to_vec(),
pass.prev,
claimed,
signature,
) {
Ok(record) => {
pass.prev = record.hash;
pass.resealed.push(record);
}
Err(skew) => {
report.findings.push(format!(
"run {current}: record {} is at a version this build does not read, and \
its bytes hash as written — this is a build skew rather than an edit: \
{skew}",
pass.last_seq
));
pass.clean = false;
pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
}
}
}
fn note_unknown_members(kind: &str, value: &serde_json::Value, report: &mut VerifyReport) {
let known: &[&str] = match kind {
"agentplane.export" => &["kind", "version", "checkpoint", "canon"],
"agentplane.export.run" => &["kind", "run", "index", "seal"],
"agentplane.export.case" => &["kind", "case", "deadlines", "blobs", "hold"],
"agentplane.export.end" => &[
"kind",
"runs_requested",
"runs_exported",
"records",
"cases",
"unreadable",
],
_ => return,
};
let Some(object) = value.as_object() else {
return;
};
let unknown: Vec<&str> = object
.keys()
.map(String::as_str)
.filter(|member| !known.contains(member))
.collect();
if !unknown.is_empty() {
report.not_checked.push(format!(
"a {kind} line carries {} this build does not know — whatever they claim was \
not checked, and a later build wrote this file",
unknown.join(", ")
));
}
}
fn edited_record(raw_bytes: &[u8], pass: &mut RunPass, report: &mut VerifyReport) {
let seq = serde_json::from_slice::<serde_json::Value>(raw_bytes)
.ok()
.and_then(|v| v.get("seq").and_then(serde_json::Value::as_u64))
.unwrap_or(pass.last_seq + 1);
report.findings.push(format!(
"run {}: record {seq} does not recompute to the hash it carries — it was edited \
after it was sealed",
pass.run
));
pass.clean = false;
pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
pass.last_seq = seq;
}
fn unparsed_record(
raw_bytes: &[u8],
parse: serde_json::Error,
pass: &mut RunPass,
report: &mut VerifyReport,
) {
let current = pass.run;
let classified = crate::journal::unreadable(raw_bytes, parse);
match &classified {
crate::core::StoreError::UnreadableRecordShape { .. } => {
report.findings.push(format!(
"run {current}: a record is at a shape this build does not read — a build \
skew rather than a damaged file: {classified}"
));
}
_ => report
.findings
.push(format!("run {current}: a record line is malformed")),
}
pass.clean = false;
pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
if let Some(seq) = serde_json::from_slice::<serde_json::Value>(raw_bytes)
.ok()
.and_then(|v| v.get("seq").and_then(serde_json::Value::as_u64))
{
pass.last_seq = seq;
}
}
fn finish_run(
report: &mut VerifyReport,
pass: Option<RunPass>,
verifier: Option<&dyn crate::core::Verifier>,
read_runs: &mut usize,
empty_blocks: &mut Vec<RunId>,
) {
let Some(pass) = pass else {
return;
};
let run = pass.run;
if pass.records == 0 {
empty_blocks.push(run);
return;
}
*read_runs += 1;
let mut ok = pass.clean;
if let Some(seal) = pass.declared_seal
&& seal != pass.prev
{
report.findings.push(format!(
"run {run}: the log's leaf is not this run's terminal hash, so the chain in this \
file is not the chain the checkpoint committed to"
));
ok = false;
}
if let Some(v) = verifier
&& let Err(e) = crate::journal::Record::verify_signed(
&pass.resealed,
crate::core::Digest::ZERO,
v,
true,
)
{
report.findings.push(format!("run {run}: {e}"));
ok = false;
}
if ok {
report.sound.push(run);
}
}
fn losses(parsed: &Parsed) -> Vec<String> {
let mut out = Vec::new();
if parsed.canon != Some(u64::from(crate::core::canon::VERSION)) {
out.push(format!(
"the digests — the export was written under canonicalization rule {:?} and this \
build implements {}, so the rebuilt store re-derives every digest under the new \
rule and its checkpoint cannot match the export's. The data is restored; \
`is_faithful` is unprovable, not false",
parsed.canon,
crate::core::canon::VERSION
));
}
if parsed.signed > 0 && !parsed.runs.is_empty() {
out.push(format!(
"{} record(s) carried a signature that this store did not reproduce — `append` \
attests as the restoring store's own signer, so authorship is lost unless it \
holds the original key. Hashes and the Merkle root are unaffected",
parsed.signed
));
}
out.push(
"activity timestamps — `recent_runs` now orders by restore time rather than by when \
history happened. It is a discovery index for listing, and nothing derives a decision \
from it"
.to_owned(),
);
out.push(
"every wait's registration — a timer, a subscription and any worklist row a wait \
opened live in stores this export does not carry, so a restored run that was \
waiting has nothing to wake it: no timer fires, no subscription matches, and its \
lease was released cleanly when it suspended, so recovery does not see it either. \
Resuming each run in `awaiting` re-arms the wait from the journal"
.to_owned(),
);
out.push(
"the worklist, unclaimed inbound events, webhook registrations and their delivery \
cursors, governed memory, and the batch, quota and standing-authority ledgers — \
none of these layers is in the export. A decision a run already consumed survives \
because that run journaled it; one nobody had consumed does not"
.to_owned(),
);
out
}
async fn awaiting_runs(
store: &Arc<dyn JournalStore>,
runs: &[RestoredRun],
) -> Result<Vec<RunId>, StoreError> {
let mut awaiting = Vec::new();
for run in runs {
let records = store.read(run.run, 1).await?;
if matches!(
crate::runtime::observed_status(&records),
Some(crate::runtime::RunStatus::Suspended(_))
) {
awaiting.push(run.run);
}
}
Ok(awaiting)
}
pub async fn runs_in_flight(
store: &Arc<dyn JournalStore>,
limit: usize,
) -> Result<InFlight, StoreError> {
let mut found = InFlight::default();
let mut after: Option<(u64, RunId)> = None;
while found.runs.len() < limit {
let page = store.recent_runs(after, CASE_PAGE).await?;
if page.is_empty() {
break;
}
after = page.last().map(|(run, at)| (*at, *run));
for (run, _) in page {
if found.runs.len() == limit {
found.truncated = true;
return Ok(found);
}
let head = store.head(run).await?;
if head.seq == 0 {
continue;
}
match store.read(run, head.seq).await {
Ok(last) => {
let in_flight = match crate::runtime::observed_status(&last) {
None | Some(crate::runtime::RunStatus::Suspended(_)) => true,
Some(_) => false,
};
if in_flight {
found.runs.push(run);
}
}
Err(e) => found.unreadable.push((run, e.to_string())),
}
}
}
Ok(found)
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct InFlight {
pub runs: Vec<RunId>,
pub truncated: bool,
pub unreadable: Vec<(RunId, String)>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct RestoreReport {
pub expected: Checkpoint,
pub rebuilt: Checkpoint,
pub runs: usize,
pub records: usize,
pub cases: usize,
pub awaiting: Vec<RunId>,
pub not_carried: Vec<String>,
}
impl RestoreReport {
#[must_use]
pub fn is_faithful(&self) -> bool {
self.expected.size == self.rebuilt.size && self.expected.root == self.rebuilt.root
}
}
pub async fn from_jsonl<R: std::io::BufRead>(
store: &Arc<dyn JournalStore>,
cases: Option<&Arc<dyn crate::case::CaseStore>>,
input: R,
) -> Result<RestoreReport, StoreError> {
let parsed = parse(input).map_err(|e| StoreError::Backend(e.to_string()))?;
restore_parsed(store, cases, parsed).await
}
#[cfg(feature = "redb")]
#[derive(Debug)]
pub struct ReplaySource {
pub store: Arc<crate::store::RedbStore>,
pub runs: Vec<RunId>,
pub report: RestoreReport,
}
#[cfg(feature = "redb")]
pub async fn open_for_replay<R: std::io::BufRead>(input: R) -> Result<ReplaySource, StoreError> {
let parsed = parse(input).map_err(|e| StoreError::Backend(e.to_string()))?;
let runs = parsed.runs.iter().map(|r| r.run).collect();
let store = Arc::new(crate::store::RedbStore::open_in_memory()?);
let journal = Arc::clone(&store) as Arc<dyn JournalStore>;
let cases = Arc::clone(&store) as Arc<dyn crate::case::CaseStore>;
let report = restore_parsed(&journal, Some(&cases), parsed).await?;
if !report.is_faithful() {
return Err(StoreError::Backend(format!(
"the export does not rebuild to its own checkpoint: it claims {} records under \
root {}, and restoring it produced {} under {}",
report.expected.size, report.expected.root, report.rebuilt.size, report.rebuilt.root
)));
}
Ok(ReplaySource {
store,
runs,
report,
})
}
async fn restore_parsed(
store: &Arc<dyn JournalStore>,
cases: Option<&Arc<dyn crate::case::CaseStore>>,
parsed: Parsed,
) -> Result<RestoreReport, StoreError> {
if parsed.version != Some(u64::from(FORMAT_VERSION)) {
return Err(StoreError::Backend(format!(
"the export claims format version {:?} and this build reads {FORMAT_VERSION} — \
restoring a format this build cannot fully parse would rebuild an unknowable \
subset and call it a history",
parsed.version
)));
}
if !parsed.complete {
return Err(StoreError::Backend(
"the export has no trailer, so it was cut short — every line in it is a valid \
prefix, and restoring a prefix would rebuild a partial history shaped exactly \
like a whole one. Re-take the export"
.to_owned(),
));
}
let mut records = 0usize;
for run in &parsed.runs {
for batch in run.bodies.chunk_by(|a, b| a.epoch == b.epoch) {
let Some(epoch) = batch.first().map(|b| b.epoch) else {
continue;
};
let appends: Vec<Append> = batch.iter().cloned().map(Append::from_body).collect();
records += appends.len();
store.append(epoch, appends).await?;
}
}
let mut sealed: Vec<&RestoredRun> = parsed
.runs
.iter()
.filter(|r| r.index.is_some())
.collect::<Vec<_>>();
sealed.sort_by_key(|r| r.index);
for run in sealed {
let (Some(outcome), Some(epoch)) =
(run.outcome.as_deref(), run.bodies.last().map(|b| b.epoch))
else {
continue;
};
store.seal(run.run, epoch, outcome).await?;
}
let mut imported = 0usize;
let mut not_carried = Vec::new();
match (cases, parsed.cases.is_empty()) {
(Some(case_store), false) => {
for block in &parsed.cases {
case_store
.import_case(&block.case, &block.deadlines, &block.blobs)
.await?;
if let Some(hold) = &block.hold {
case_store.place_hold(block.case.id, hold).await?;
}
imported += 1;
}
not_carried.push(
"blob link timestamps — the export carries a case's blob digests without \
the instant each link was written, so erasure reachability survives and \
the original ordering does not"
.to_owned(),
);
}
(None, false) => not_carried.push(format!(
"the case layer — the export carries {} case(s) and no case store was supplied, \
so the journal is rebuilt and the matters it names are not",
parsed.cases.len()
)),
(_, true) => {}
}
not_carried.extend(losses(&parsed));
let awaiting = awaiting_runs(store, &parsed.runs).await?;
let rebuilt = store.checkpoint().await?;
if rebuilt.origin != parsed.checkpoint.origin {
not_carried.push(format!(
"the log identity — this history was written by '{}' and is now held by '{}'. \
A legitimate recovery: a restore is pointed at a store, and one tenant's \
history put back under another tenant's name is a different log with the \
same contents. It is named because a checkpoint an auditor holds from \
before the disaster will report the new log as the wrong one",
parsed.checkpoint.origin, rebuilt.origin
));
}
Ok(RestoreReport {
expected: parsed.checkpoint,
rebuilt,
runs: parsed.runs.len(),
records,
cases: imported,
awaiting,
not_carried,
})
}
struct RestoredRun {
run: RunId,
index: Option<u64>,
outcome: Option<String>,
bodies: Vec<crate::journal::RecordBody>,
prev: crate::core::Digest,
}
fn replayable(
raw: &[u8],
prev: crate::core::Digest,
claimed: crate::core::Digest,
) -> Result<crate::journal::RecordBody, std::io::Error> {
crate::journal::Record::from_stored_signed(raw.to_vec(), prev, claimed, None)
.map(|record| record.body)
.map_err(|e| {
std::io::Error::other(match e {
StoreError::Corrupt { .. } => format!(
"a record's claimed hash does not cover its wire bytes and the chain \
before it ({e}) — replaying it would rebuild a history the export never \
committed to, so the file is refused before anything is written"
),
StoreError::UnknownRecordVersion { .. } => format!(
"a record is at a version this build does not restore ({e}) — `append` \
would rewrite it at this build's version, so the file is refused before \
anything is written"
),
other => format!(
"a record line's wire bytes do not parse ({other}) — the record cannot be \
replayed as written, and its display copy is not a substitute"
),
})
})
}
struct Parsed {
checkpoint: Checkpoint,
version: Option<u64>,
canon: Option<u64>,
complete: bool,
runs: Vec<RestoredRun>,
cases: Vec<RestoredCase>,
signed: usize,
}
struct RestoredCase {
case: crate::core::Case,
deadlines: Vec<crate::core::Deadline>,
blobs: Vec<crate::core::Digest>,
hold: Option<crate::core::LegalHold>,
}
fn case_hold(value: &serde_json::Value) -> Result<Option<crate::core::LegalHold>, String> {
let Some(hold) = value.get("hold") else {
return Err(
"a case block carries no `hold` member, so whether the matter is under \
a legal hold is unknown"
.to_owned(),
);
};
serde_json::from_value::<Option<crate::core::LegalHold>>(hold.clone())
.map_err(|e| format!("a case block's legal hold is malformed: {e}"))
}
fn restored_case(value: &serde_json::Value) -> Result<Option<RestoredCase>, std::io::Error> {
use serde_json::Value;
let (Ok(case), Some(deadlines), Some(blobs)) = (
serde_json::from_value::<crate::core::Case>(
value.get("case").cloned().unwrap_or(Value::Null),
),
value
.get("deadlines")
.and_then(|d| serde_json::from_value::<Vec<crate::core::Deadline>>(d.clone()).ok()),
value
.get("blobs")
.and_then(|b| serde_json::from_value::<Vec<crate::core::Digest>>(b.clone()).ok()),
) else {
return Ok(None);
};
let hold = case_hold(value).map_err(|e| {
std::io::Error::other(format!(
"case {}: {e} — restoring the matter without it would let retention \
erase it, so the file is refused",
case.id
))
})?;
Ok(Some(RestoredCase {
case,
deadlines,
blobs,
hold,
}))
}
fn parse<R: std::io::BufRead>(input: R) -> Result<Parsed, std::io::Error> {
use serde_json::Value;
let mut parsed = Parsed {
checkpoint: Checkpoint {
origin: String::new(),
size: 0,
root: crate::core::Digest::ZERO,
},
version: None,
canon: None,
complete: false,
runs: Vec::new(),
cases: Vec::new(),
signed: 0,
};
for line in input.lines() {
let line = line?;
let Ok(value) = serde_json::from_str::<Value>(&line) else {
continue;
};
match value.get("kind").and_then(Value::as_str) {
Some("agentplane.export") => {
parsed.version = value.get("version").and_then(Value::as_u64);
parsed.canon = value.get("canon").and_then(Value::as_u64);
if let Some(c) = value
.get("checkpoint")
.and_then(|c| serde_json::from_value::<Checkpoint>(c.clone()).ok())
{
parsed.checkpoint = c;
}
}
Some("agentplane.export.run") => {
if let Some(run) = value
.get("run")
.and_then(Value::as_str)
.and_then(|s| RunId::parse(s).ok())
{
parsed.runs.push(RestoredRun {
run,
index: value.get("index").and_then(Value::as_u64),
outcome: None,
bodies: Vec::new(),
prev: crate::core::Digest::ZERO,
});
}
}
Some("agentplane.export.case") => {
if let Some(case) = restored_case(&value)? {
parsed.cases.push(case);
}
}
Some("agentplane.export.end") => parsed.complete = true,
Some(_) => {}
_ => {
if value.get("signature").is_some_and(|a| !a.is_null()) {
parsed.signed += 1;
}
let Some(raw) = value.get("raw").and_then(Value::as_str) else {
return Err(std::io::Error::other(
"a record line carries no wire bytes (`raw`) — restoring its display \
copy instead would rebuild what the readable half says rather than \
what the chain hashed, so the file is refused instead of guessed at",
));
};
let Some(claimed) = value
.get("hash")
.and_then(|h| serde_json::from_value::<crate::core::Digest>(h.clone()).ok())
else {
return Err(std::io::Error::other(
"a record line carries no hash — nothing ties its bytes to the chain, \
so the file is refused instead of replayed",
));
};
let Some(current) = parsed.runs.last_mut() else {
continue;
};
let body = replayable(raw.as_bytes(), current.prev, claimed)?;
current.prev = claimed;
if let crate::journal::RecordKind::RunConcluded { outcome, .. } = &body.kind {
current.outcome = Some(outcome.clone());
}
current.bodies.push(body);
}
}
}
Ok(parsed)
}