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: &'r crate::journal::RecordBody,
pub prev_hash: &'r crate::core::Digest,
pub hash: &'r crate::core::Digest,
pub attestation: Option<&'r crate::core::Attestation>,
}
impl<'r> From<&'r crate::journal::Record> for ExportedRecord<'r> {
fn from(r: &'r crate::journal::Record) -> Self {
Self {
seq: r.seq(),
body: &r.body,
prev_hash: &r.prev_hash,
hash: &r.hash,
attestation: r.attestation.as_ref(),
}
}
}
#[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, Eq, serde::Serialize)]
pub struct Trailer {
pub kind: &'static str,
pub runs_requested: usize,
pub runs_exported: usize,
pub records: 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>,
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) => {
for record in &found {
writeln!(out, "{}", to_line(&ExportedRecord::from(record))?)?;
records += 1;
}
exported += 1;
}
Err(e) => unreadable.push(Unreadable {
run,
reason: e.to_string(),
}),
}
}
let trailer = Trailer {
kind: "agentplane.export.end",
runs_requested: runs.len(),
runs_exported: exported,
records,
unreadable,
};
writeln!(out, "{}", to_line(&trailer)?)?;
out.flush()?;
Ok(trailer)
}
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 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>,
) -> Result<VerifyReport, std::io::Error> {
use crate::core::{Digest, merkle};
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,
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, Digest)> = Vec::new();
let mut pass: Option<RunPass> = None;
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.run") => {
finish_run(&mut report, pass.take(), verifier);
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,
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)));
}
}
Some("agentplane.export.end") => report.complete = true,
_ => {
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);
}
}
}
finish_run(&mut report, pass, verifier);
settle(&mut report, header_seen, leaves);
Ok(report)
}
struct RunPass {
run: RunId,
declared_seal: Option<crate::core::Digest>,
prev: crate::core::Digest,
last_seq: u64,
resealed: Vec<crate::journal::Record>,
clean: bool,
}
fn settle(
report: &mut VerifyReport,
header_seen: bool,
mut leaves: Vec<(u64, crate::core::Digest)>,
) {
use crate::core::merkle;
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 == report.checkpoint.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 != report.checkpoint.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",
report.checkpoint.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",
report.checkpoint.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.findings.push(format!(
"the export was written under canonicalization rule {canon:?} and this build \
implements {} — every digest in it is unverifiable here, which is a different \
statement from wrong",
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) {
let current = pass.run;
let (Some(body), Some(claimed)) = (
value
.get("body")
.and_then(|b| serde_json::from_value::<crate::journal::RecordBody>(b.clone()).ok()),
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 is malformed"));
pass.clean = false;
return;
};
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;
match crate::journal::Record::seal(body, pass.prev) {
Ok(mut record) => {
if record.hash != claimed {
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 = record.hash;
record.attestation = value
.get("attestation")
.and_then(|a| {
serde_json::from_value::<Option<crate::core::Attestation>>(a.clone()).ok()
})
.flatten();
pass.resealed.push(record);
}
Err(e) => {
report.findings.push(format!(
"run {current}: record {} cannot be sealed: {e}",
pass.last_seq
));
pass.clean = false;
}
}
}
fn finish_run(
report: &mut VerifyReport,
pass: Option<RunPass>,
verifier: Option<&dyn crate::core::Verifier>,
) {
let Some(pass) = pass else {
return;
};
let run = pass.run;
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);
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct RestoreReport {
pub expected: Checkpoint,
pub rebuilt: Checkpoint,
pub runs: usize,
pub records: usize,
pub not_carried: Vec<String>,
}
impl RestoreReport {
#[must_use]
pub fn is_faithful(&self) -> bool {
self.expected == self.rebuilt
}
}
pub async fn from_jsonl<R: std::io::BufRead>(
store: &Arc<dyn JournalStore>,
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
)));
}
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 not_carried = Vec::new();
if parsed.canon != Some(u64::from(crate::core::canon::VERSION)) {
not_carried.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() {
not_carried.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
));
}
not_carried.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(),
);
Ok(RestoreReport {
expected: parsed.checkpoint,
rebuilt: store.checkpoint().await?,
runs: parsed.runs.len(),
records,
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>,
runs: Vec<RestoredRun>,
signed: usize,
}
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,
runs: 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.end") => {}
_ => {
if value.get("attestation").is_some_and(|a| !a.is_null()) {
parsed.signed += 1;
}
let Some(body) = value.get("body").and_then(|b| {
serde_json::from_value::<crate::journal::RecordBody>(b.clone()).ok()
}) else {
continue;
};
if let Some(current) = parsed.runs.last_mut() {
if let crate::journal::RecordKind::RunSealed { outcome, .. } = &body.kind {
current.outcome = Some(outcome.clone());
}
current.bodies.push(body);
}
}
}
}
Ok(parsed)
}