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 attestation: Option<&'r crate::core::Attestation>,
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,
attestation: r.attestation.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>,
}
#[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))?;
writeln!(
out,
"{}",
to_line(&CaseBlock {
kind: "agentplane.export.case",
case,
deadlines,
blobs,
})?
)?;
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>,
expected: Option<&Checkpoint>,
) -> 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,
};
if verifier.is_none() {
report.not_checked.push(
"signatures — no public key was supplied, so this pass cannot say who wrote \
anything"
.to_owned(),
);
}
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;
};
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, expected);
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());
}
*blob_digests += value
.get("blobs")
.and_then(Value::as_array)
.map_or(0, Vec::len);
}
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 settle(
report: &mut VerifyReport,
header_seen: bool,
mut leaves: Vec<(u64, crate::core::merkle::LeafHash)>,
expected: Option<&Checkpoint>,
) {
use crate::core::merkle;
let against = if let Some(given) = expected {
if header_seen && *given != report.checkpoint {
report.findings.push(format!(
"the export's header names log '{}' at size {} with root {}, and the \
checkpoint this pass was given 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(),
given.origin,
given.size,
given.root.to_hex(),
));
}
given.clone()
} else {
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());
}
leaves.sort_by_key(|(index, _)| *index);
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();
let Ok(body) = serde_json::from_slice::<crate::journal::RecordBody>(raw_bytes) else {
report
.findings
.push(format!("run {current}: a record line is malformed"));
pass.clean = false;
return;
};
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 attestation = value
.get("attestation")
.and_then(|a| serde_json::from_value::<Option<crate::core::Attestation>>(a.clone()).ok())
.flatten();
if let Ok(record) = crate::journal::Record::from_stored_attested(
raw_bytes.to_vec(),
pass.prev,
claimed,
attestation,
) {
pass.prev = record.hash;
pass.resealed.push(record);
} else {
report.findings.push(format!(
"run {current}: record {} does not recompute to the hash it carries \
— it was edited after it was sealed",
pass.last_seq
));
pass.clean = false;
pass.prev = crate::core::Digest::chain(pass.prev, raw_bytes);
}
}
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_attested(
&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()))?;
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?;
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>,
}
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>,
}
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(),
});
}
}
Some("agentplane.export.case") => {
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 {
continue;
};
parsed.cases.push(RestoredCase {
case,
deadlines,
blobs,
});
}
Some("agentplane.export.end") => parsed.complete = true,
Some(_) => {}
_ => {
if value.get("attestation").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 body = serde_json::from_slice::<crate::journal::RecordBody>(raw.as_bytes())
.map_err(|e| {
std::io::Error::other(format!(
"a record line's wire bytes do not parse ({e}) — the record cannot \
be replayed as written, and its display copy is not a substitute"
))
})?;
if let Some(current) = parsed.runs.last_mut() {
if let crate::journal::RecordKind::RunConcluded { outcome, .. } = &body.kind {
current.outcome = Some(outcome.clone());
}
current.bodies.push(body);
}
}
}
}
Ok(parsed)
}