use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use crate::event::Envelope;
use crate::journal;
use crate::ledger::{self, RunPaths};
use crate::projection::{self, RunState};
pub(crate) const CHECKPOINT_SCHEMA_VERSION: u32 = 4;
fn this_version<'de, D: serde::Deserializer<'de>>(reader: D) -> Result<u32, D::Error> {
let found = u32::deserialize(reader)?;
if found != CHECKPOINT_SCHEMA_VERSION {
return Err(serde::de::Error::custom(format!(
"checkpoint schema_version {found}, and this build reads {CHECKPOINT_SCHEMA_VERSION}"
)));
}
Ok(found)
}
#[derive(Debug, Clone, Default, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Placed {
pub(crate) ts: String,
pub(crate) stream: String,
}
fn placed(event: &Envelope) -> Placed {
Placed {
ts: event.ts.clone(),
stream: event.stream.clone(),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Coverage {
pub(crate) bytes: u64,
pub(crate) records: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(crate) at: Option<Placed>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub(crate) streams: BTreeMap<String, u64>,
#[serde(serialize_with = "as_hex", deserialize_with = "of_hex")]
pub(crate) digest: u128,
}
pub(crate) const NOTHING_DIGESTED: u128 = 0x6c62_272e_07bb_0142_62b8_2175_6295_c58d;
const FNV_PRIME: u128 = (1 << 88) | 0x13b;
pub(crate) fn digested(from: u128, bytes: &[u8]) -> u128 {
bytes.iter().fold(from, |digest, byte| {
(digest ^ u128::from(*byte)).wrapping_mul(FNV_PRIME)
})
}
fn digest_of_prefix(journal: &std::path::Path, bytes: u64) -> Option<u128> {
ledger::read_range(journal, 0, bytes).map(|prefix| digested(NOTHING_DIGESTED, &prefix))
}
pub(crate) fn as_hex<S: serde::Serializer>(digest: &u128, writer: S) -> Result<S::Ok, S::Error> {
writer.serialize_str(&format!("{digest:032x}"))
}
pub(crate) fn of_hex<'de, D: serde::Deserializer<'de>>(reader: D) -> Result<u128, D::Error> {
let written = String::deserialize(reader)?;
let refuse = |why: &str| serde::de::Error::custom(format!("digest '{written}': {why}"));
let a_digit = |digit: &u8| matches!(digit, b'0'..=b'9' | b'a'..=b'f');
if written.len() != 32 || !written.as_bytes().iter().all(a_digit) {
return Err(refuse("not 32 lower-case hex digits"));
}
u128::from_str_radix(&written, 16).map_err(|e| refuse(&e.to_string()))
}
impl Default for Coverage {
fn default() -> Self {
Self {
bytes: 0,
records: 0,
at: None,
streams: BTreeMap::new(),
digest: NOTHING_DIGESTED,
}
}
}
impl Coverage {
fn is_in_front_of(&self, event: &Envelope) -> bool {
let placed = placed(event);
self.at.as_ref().is_none_or(|at| *at <= placed)
&& self
.streams
.get(&event.stream)
.is_none_or(|reached| *reached <= event.seq)
}
fn sorts_in_front_of(&self, grown: &[(Option<Envelope>, u64)]) -> bool {
grown
.iter()
.filter_map(|(event, _)| event.as_ref())
.all(|event| self.is_in_front_of(event))
}
fn sealed_with(&self, digested_bytes: u128, state: &RunState) -> u128 {
let mut sealed = digested(digested_bytes, &self.bytes.to_le_bytes());
sealed = digested(sealed, &self.records.to_le_bytes());
if let Some(at) = &self.at {
sealed = digested(sealed, at.ts.as_bytes());
sealed = digested(sealed, at.stream.as_bytes());
}
for (stream, reached) in &self.streams {
sealed = digested(sealed, stream.as_bytes());
sealed = digested(sealed, &reached.to_le_bytes());
}
match serde_json::to_vec(state) {
Ok(written) => digested(sealed, &written),
Err(_) => sealed.wrapping_add(1),
}
}
fn corroborated_by(&self, journal: &std::path::Path, state: &RunState) -> Option<u128> {
if length_of(journal) < self.bytes {
return None;
}
if self.bytes > 0 && !ends_a_record(journal, self.bytes) {
return None;
}
digest_of_prefix(journal, self.bytes)
.filter(|digested_bytes| self.sealed_with(*digested_bytes, state) == self.digest)
}
fn absorb(&mut self, event: Option<&Envelope>, bytes: u64) {
self.bytes += bytes;
self.records += 1;
let Some(event) = event else {
return;
};
let placed = placed(event);
if self.at.as_ref().is_none_or(|at| *at < placed) {
self.at = Some(placed);
}
let reached = self.streams.entry(event.stream.clone()).or_default();
*reached = (*reached).max(event.seq);
}
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Checkpoint {
#[serde(deserialize_with = "this_version")]
schema_version: u32,
run_id: String,
coverage: Coverage,
state: RunState,
}
#[derive(Debug)]
pub(crate) struct Projected {
covered: RunState,
coverage: Coverage,
state: RunState,
took: u64,
digested_bytes: u128,
}
impl std::ops::Deref for Projected {
type Target = RunState;
fn deref(&self) -> &RunState {
&self.state
}
}
impl std::ops::DerefMut for Projected {
fn deref_mut(&mut self) -> &mut RunState {
&mut self.state
}
}
impl Projected {
pub(crate) fn open(paths: &RunPaths) -> Self {
let resumed = readable(paths).and_then(|checkpoint| {
let digested_bytes = checkpoint
.coverage
.corroborated_by(&paths.journal(), &checkpoint.state)?;
Some((checkpoint, digested_bytes))
});
let mut projected = match resumed {
Some((checkpoint, digested_bytes)) => Self {
state: checkpoint.state.clone(),
covered: checkpoint.state,
coverage: checkpoint.coverage,
took: 0,
digested_bytes,
},
None => Self::empty(),
};
projected.refresh(paths);
projected
}
pub(crate) fn refresh(&mut self, paths: &RunPaths) {
let journal = paths.journal();
let mut grown = journal::finished_records_after(&journal, self.coverage.bytes);
if length_of(&journal) < self.coverage.bytes || !self.coverage.sorts_in_front_of(&grown) {
*self = Self::empty();
grown = journal::finished_records_after(&journal, 0);
}
self.took = grown.len() as u64;
self.take(paths, &grown);
self.state.cross_dag = BTreeMap::new();
}
pub(crate) fn into_state(self) -> RunState {
self.state
}
#[cfg(test)]
pub(crate) fn took(&self) -> u64 {
self.took
}
#[cfg(test)]
pub(crate) fn coverage(&self) -> &Coverage {
&self.coverage
}
fn empty() -> Self {
let empty = RunState {
strict: true,
..RunState::default()
};
Self {
state: empty.clone(),
covered: empty,
coverage: Coverage::default(),
took: 0,
digested_bytes: NOTHING_DIGESTED,
}
}
fn take(&mut self, paths: &RunPaths, grown: &[(Option<Envelope>, u64)]) {
crate::loopstats::records_folded(grown.len() as u64);
let accountable = extent(&self.coverage, grown);
let taking: u64 = grown[..accountable].iter().map(|(_, bytes)| bytes).sum();
let taken = match taking {
0 => Some(Vec::new()),
_ => ledger::read_range(&paths.journal(), self.coverage.bytes, taking),
};
let (covered, ahead) = grown.split_at(taken.as_ref().map_or(0, |_| accountable));
fold_in_merge_order(&mut self.covered, covered);
for (event, bytes) in covered {
self.coverage.absorb(event.as_ref(), *bytes);
}
if let Some(taken) = taken.filter(|taken| !taken.is_empty()) {
self.digested_bytes = digested(self.digested_bytes, &taken);
self.coverage.digest = self
.coverage
.sealed_with(self.digested_bytes, &self.covered);
self.write(paths);
}
self.state = self.covered.clone();
fold_in_merge_order(&mut self.state, ahead);
}
fn write(&self, paths: &RunPaths) {
let _ = ledger::write_json(
&paths.checkpoint(),
&Checkpoint {
schema_version: CHECKPOINT_SCHEMA_VERSION,
run_id: paths.run.clone(),
coverage: self.coverage.clone(),
state: self.covered.clone(),
},
);
}
}
fn fold_in_merge_order(state: &mut RunState, records: &[(Option<Envelope>, u64)]) {
let mut ordered: Vec<Envelope> = records
.iter()
.filter_map(|(event, _)| event.clone())
.collect();
journal::merge_order(&mut ordered);
for event in &ordered {
projection::fold_one(state, event);
}
}
pub(crate) fn fold_and_checkpoint(paths: &RunPaths) -> RunState {
Projected::open(paths).into_state()
}
fn readable(paths: &RunPaths) -> Option<Checkpoint> {
ledger::read_json_opt::<Checkpoint>(&paths.checkpoint())
.filter(|checkpoint| checkpoint.run_id == paths.run)
}
fn ends_a_record(journal: &std::path::Path, at: u64) -> bool {
use std::io::{Read, Seek, SeekFrom};
let Ok(mut file) = std::fs::File::open(journal) else {
return false;
};
if file.seek(SeekFrom::Start(at - 1)).is_err() {
return false;
}
let mut byte = [0u8; 1];
file.read_exact(&mut byte).is_ok() && byte[0] == b'\n'
}
fn length_of(journal: &std::path::Path) -> u64 {
std::fs::metadata(journal).map_or(0, |about| about.len())
}
fn extent(coverage: &Coverage, grown: &[(Option<Envelope>, u64)]) -> usize {
let cap = before_the_open_instant(grown);
let mut ruled_out = vec![0i64; cap + 2];
let mut rule_out = |from: usize, to: usize| {
if from <= cap {
ruled_out[from] += 1;
ruled_out[to.min(cap) + 1] -= 1;
}
};
let mut reached: Vec<Option<Placed>> = vec![coverage.at.clone()];
let mut per_stream: BTreeMap<&str, Vec<(usize, u64)>> = BTreeMap::new();
for (at, (event, _)) in grown.iter().enumerate() {
let front = reached[at].clone();
let Some(event) = event else {
reached.push(front);
continue;
};
let placed = placed(event);
if front.as_ref().is_some_and(|front| *front > placed) {
let first =
reached.partition_point(|reached| reached.as_ref().is_none_or(|at| *at <= placed));
rule_out(first, at);
}
let stream = per_stream.entry(event.stream.as_str()).or_default();
if stream
.last()
.is_some_and(|(_, reached)| *reached > event.seq)
{
let first = stream.partition_point(|(_, reached)| *reached <= event.seq);
rule_out(stream[first].0 + 1, at);
}
let carried = stream.last().map_or(event.seq, |(_, reached)| *reached);
stream.push((at, carried.max(event.seq)));
reached.push(Some(match front {
Some(front) if front > placed => front,
_ => placed,
}));
}
let mut covered = 0;
let mut spanning = 0;
for (boundary, opened) in ruled_out.iter().enumerate().take(cap + 1) {
spanning += opened;
if spanning == 0 {
covered = boundary;
}
}
covered
}
fn before_the_open_instant(grown: &[(Option<Envelope>, u64)]) -> usize {
let Some(last) = grown
.iter()
.rev()
.find_map(|(event, _)| event.as_ref().map(|event| event.ts.clone()))
else {
return grown.len();
};
grown
.iter()
.rposition(|(event, _)| event.as_ref().is_some_and(|event| event.ts != last))
.map_or(0, |before| before + 1)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::{EventKind, Labels, Source, ENVELOPE_VERSION};
use crate::journal::{Journal, PipelineKind};
use crate::ledger::LaunchRecord;
use crate::plan::{Goal, Node, Plan, PLAN_SCHEMA_VERSION};
use serde_json::{json, Value};
use std::path::{Path, PathBuf};
fn scratch(name: &str) -> PathBuf {
let opened = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |since| since.as_nanos());
let root = std::env::temp_dir().join(format!("onepipeline-checkpoint-{name}-{opened}"));
let _ = std::fs::remove_dir_all(&root);
std::fs::create_dir_all(&root).expect("a scratch root");
root
}
fn plan(nodes: &[&str]) -> Plan {
Plan {
schema_version: PLAN_SCHEMA_VERSION,
goal: Some(Goal {
text: "fold a run without replaying it".into(),
}),
name: Some("demo".into()),
concurrency: 4,
tasks: nodes
.iter()
.map(|id| Node {
id: (*id).to_string(),
persona: Some("engineer".into()),
task: Some("## What\ndo it".into()),
..Node::default()
})
.collect(),
}
}
fn a_run(root: &Path, run: &str) -> RunPaths {
let paths = RunPaths::under(root, run);
paths.create().expect("the run directory");
let record = LaunchRecord {
run_id: run.to_string(),
project: "plans:demo".into(),
dir: PathBuf::from("/tmp/launch"),
graph: String::new(),
graph_run: String::new(),
observer_runs: Vec::new(),
observer_ending: String::new(),
node_graph: "graph".into(),
pr_author_graph: String::new(),
node_validator: String::new(),
envelope_reviewer: String::new(),
launcher: "checkpoint-journey".into(),
session: "a-session".into(),
pid: 0,
host: String::new(),
started: String::new(),
started_at: crate::sys::now_rfc3339(),
heartbeat_interval: 1_800,
writeback_item_budget: 0,
success_hook: String::new(),
failure_hook: String::new(),
hook_timeout: 0,
dispatch_env_hook: String::new(),
dispatch_env_hook_timeout: 0,
dag_sets: Vec::new(),
node_sets: Vec::new(),
adoptions: 0,
filters: crate::filter::Filters::default(),
bus_config: Default::default(),
envelope_reviewer_bar: Default::default(),
};
crate::ledger::write_json(&paths.launch(), &record).expect("a launch record");
paths
}
fn an_instant_later() {
std::thread::sleep(std::time::Duration::from_millis(2));
}
fn a_recorded_run(root: &Path, run: &str) -> RunPaths {
let paths = a_run(root, run);
let mut journal = Journal::open(&paths);
journal
.emit(
PipelineKind::RunStarted,
crate::journal::labels(run, None),
crate::journal::payload(&[("plan", json!(plan(&["build", "ship"])))]),
)
.expect("appended");
an_instant_later();
journal
.emit(
PipelineKind::NodeDispatched,
crate::journal::labels(run, Some("build")),
crate::journal::payload(&[("persona", json!("engineer")), ("attempt", json!(1))]),
)
.expect("appended");
paths
}
fn settle(paths: &RunPaths, node: &str, status: &str) {
an_instant_later();
Journal::open(paths)
.emit(
PipelineKind::NodeSettled,
crate::journal::labels(&paths.run, Some(node)),
crate::journal::payload(&[("status", json!(status))]),
)
.expect("appended");
}
fn relayed(paths: &RunPaths, stream: &str, seq: u64, ts: &str, node: &str) {
let envelope = Envelope {
v: ENVELOPE_VERSION,
ts: ts.to_string(),
stream: stream.to_string(),
seq,
source: Source::Agentgraph,
kind: EventKind("turn-activity".into()),
dimensions: Default::default(),
labels: Labels {
run_id: Some(paths.run.clone()),
node: Some(node.to_string()),
..Labels::default()
},
payload: crate::journal::payload(&[("tool", json!("Edit"))]),
artifacts: Vec::new(),
};
crate::ledger::append_line(
&paths.journal(),
&serde_json::to_string(&envelope).expect("an envelope serializes"),
)
.expect("appended");
}
fn folded_as(state: &RunState) -> Value {
serde_json::to_value(state).expect("a folded state serializes")
}
fn without_a_checkpoint(paths: &RunPaths) -> Value {
let held = std::fs::read(paths.checkpoint()).ok();
let _ = std::fs::remove_file(paths.checkpoint());
let whole = folded_as(&fold_and_checkpoint(paths));
match held {
Some(bytes) => std::fs::write(paths.checkpoint(), bytes).expect("put back"),
None => {
let _ = std::fs::remove_file(paths.checkpoint());
}
}
whole
}
fn seal_and_write(paths: &RunPaths, document: &mut Value) {
let parsed: Checkpoint =
serde_json::from_value(document.clone()).expect("a document this build reads");
let bytes = digest_of_prefix(&paths.journal(), parsed.coverage.bytes)
.expect("the journal holds the bytes the marker claims");
let sealed = parsed.coverage.sealed_with(bytes, &parsed.state);
document["coverage"]["digest"] = json!(format!("{sealed:032x}"));
crate::ledger::write_json(&paths.checkpoint(), document).expect("written");
}
fn document(paths: &RunPaths) -> Value {
crate::ledger::read_json_opt(&paths.checkpoint()).expect("the checkpoint this run carries")
}
fn folding(paths: &RunPaths) -> (Value, u64) {
let projected = Projected::open(paths);
(folded_as(&projected), projected.took())
}
#[test]
fn a_resumed_fold_lands_on_the_state_the_whole_store_folds_to() {
let root = scratch("resumed-is-whole");
let paths = a_recorded_run(&root, "r-resumed");
let _ = fold_and_checkpoint(&paths);
settle(&paths, "build", "done");
relayed(&paths, "graph-1", 0, "2099-01-01T00:00:01.000Z", "ship");
settle(&paths, "ship", "done");
assert!(paths.checkpoint().is_file(), "no checkpoint was written");
let resumed = folded_as(&fold_and_checkpoint(&paths));
assert_eq!(resumed, without_a_checkpoint(&paths));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_resumed_fold_takes_only_the_records_the_checkpoint_does_not_account_for() {
let root = scratch("takes-only-the-tail");
let paths = a_recorded_run(&root, "r-tail");
settle(&paths, "build", "done");
let _ = fold_and_checkpoint(&paths);
let held = crate::ledger::read_records(&paths.journal()).len() as u64;
let stored = super::readable(&paths).expect("the checkpoint this read wrote");
let covered = stored.coverage.records;
assert!(
covered > 0 && covered < held,
"a marker over a {held}-record store accounts for {covered}"
);
let mut planted = document(&paths);
planted["state"]["outcomes"]["build"] = json!("carried-from-the-checkpoint");
seal_and_write(&paths, &mut planted);
settle(&paths, "ship", "done");
let grew_to = crate::ledger::read_records(&paths.journal()).len() as u64;
let (resumed, took) = folding(&paths);
assert_eq!(
resumed["outcomes"]["build"], "carried-from-the-checkpoint",
"the records the marker accounts for were folded again: {resumed}"
);
assert_eq!(
resumed["recorded"]["ship"],
json!({"at": "done"}),
"the records past the marker were not folded: {resumed}"
);
assert_eq!(
took,
grew_to - covered,
"a resumed fold took more than the store grew by"
);
let poisoned = std::fs::read(paths.checkpoint()).expect("the poisoned checkpoint");
std::fs::remove_file(paths.checkpoint()).expect("the checkpoint goes away");
let (whole, took_whole) = folding(&paths);
std::fs::write(paths.checkpoint(), poisoned).expect("put back");
assert_eq!(
whole["outcomes"].get("build"),
None,
"the control fold kept an account only the checkpoint carried"
);
assert_eq!(
took_whole, grew_to,
"the control fold did not take the whole store"
);
assert!(
took < took_whole,
"a resumed fold took {took} records and a full one took {took_whole}"
);
let _ = std::fs::remove_dir_all(&root);
}
fn a_view_of(paths: &RunPaths) -> Value {
folded_as(
&crate::views::RunView::open(paths)
.expect("the run reads")
.state,
)
}
fn a_run_with_a_checkpoint(name: &str) -> (PathBuf, RunPaths, Value) {
let root = scratch(name);
let paths = a_recorded_run(&root, "r-fallback");
settle(&paths, "build", "done");
let _ = fold_and_checkpoint(&paths);
let mut planted = document(&paths);
planted["state"]["outcomes"]["build"] = json!("carried-from-the-checkpoint");
seal_and_write(&paths, &mut planted);
settle(&paths, "ship", "done");
let whole = without_a_checkpoint(&paths);
assert_eq!(
whole["outcomes"].get("build"),
None,
"the control fold carried an account only the checkpoint holds"
);
(root, paths, whole)
}
#[test]
fn an_absent_checkpoint_folds_the_whole_store() {
let (root, paths, whole) = a_run_with_a_checkpoint("absent");
std::fs::remove_file(paths.checkpoint()).expect("the checkpoint goes away");
assert_eq!(a_view_of(&paths), whole);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_checkpoint_that_cannot_be_read_folds_the_whole_store() {
let (root, paths, whole) = a_run_with_a_checkpoint("unreadable");
std::fs::write(paths.checkpoint(), b"{ this is not a checkpoint")
.expect("the checkpoint is mangled");
assert_eq!(a_view_of(&paths), whole);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_checkpoint_this_build_did_not_write_folds_the_whole_store() {
let (root, paths, whole) = a_run_with_a_checkpoint("another-version");
let mut later = document(&paths);
later["schema_version"] = json!(CHECKPOINT_SCHEMA_VERSION + 1);
crate::ledger::write_json(&paths.checkpoint(), &later).expect("written");
assert_eq!(a_view_of(&paths), whole);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_coverage_the_journal_does_not_corroborate_folds_the_whole_store() {
let (root, paths, whole) = a_run_with_a_checkpoint("uncorroborated");
let mut ahead = document(&paths);
ahead["coverage"]["at"] = json!({"ts": "2199-01-01T00:00:00.000Z", "stream": "zzzz"});
seal_and_write(&paths, &mut ahead);
assert_eq!(a_view_of(&paths), whole);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_store_the_merge_order_rearranges_folds_the_same_state_either_way() {
let root = scratch("rearranged");
let paths = a_recorded_run(&root, "r-rearranged");
settle(&paths, "build", "done");
let _ = fold_and_checkpoint(&paths);
let covered = super::readable(&paths).expect("a checkpoint").coverage;
assert!(covered.records > 0, "nothing was accounted for");
relayed(&paths, "graph-0", 0, "1999-01-01T00:00:00.000Z", "build");
let resumed = folded_as(&fold_and_checkpoint(&paths));
assert_eq!(resumed, without_a_checkpoint(&paths));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_marker_is_never_extended_over_a_record_a_reordering_would_move() {
let root = scratch("never-extended");
let paths = a_run(&root, "r-unplaceable");
relayed(&paths, "graph-b", 0, "2099-01-01T00:00:03.000Z", "build");
relayed(&paths, "graph-a", 0, "2099-01-01T00:00:01.000Z", "build");
let state = folded_as(&fold_and_checkpoint(&paths));
let covered = super::readable(&paths).map(|stored| stored.coverage);
assert!(
covered.is_none_or(|covered| covered.records == 0),
"a record the merge order moves was accounted for"
);
assert_eq!(state, without_a_checkpoint(&paths));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_marker_passes_a_local_inversion_rather_than_stopping_at_it() {
let root = scratch("past-an-inversion");
let paths = a_run(&root, "r-inverted");
for (stream, ts) in [
("graph-a", "2099-01-01T00:00:01.000Z"),
("graph-c", "2099-01-01T00:00:05.000Z"),
("graph-b", "2099-01-01T00:00:03.000Z"),
("graph-d", "2099-01-01T00:00:09.000Z"),
("graph-e", "2099-01-01T00:00:11.000Z"),
] {
relayed(&paths, stream, 0, ts, "build");
}
let state = folded_as(&fold_and_checkpoint(&paths));
let covered = super::readable(&paths)
.expect("a checkpoint")
.coverage
.records;
assert_eq!(
covered, 4,
"a marker stopped at an inversion instead of passing it"
);
assert_eq!(state, without_a_checkpoint(&paths));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_line_this_build_cannot_read_is_accounted_for_and_placed_by_nothing() {
let root = scratch("unreadable-line");
let paths = a_recorded_run(&root, "r-unreadable");
an_instant_later();
crate::ledger::append_line(&paths.journal(), "this is not a record").expect("appended");
settle(&paths, "build", "done");
settle(&paths, "ship", "done");
let held = crate::ledger::read_records(&paths.journal()).len() as u64;
let state = folded_as(&fold_and_checkpoint(&paths));
let covered = super::readable(&paths).expect("a checkpoint").coverage;
assert_eq!(
covered.bytes,
crate::ledger::read_records(&paths.journal())
.iter()
.take(covered.records as usize)
.map(|record| record.bytes + 1)
.sum::<u64>(),
"the marker's bytes and the records it counts describe different stores"
);
assert!(
covered.records > 2 && covered.records < held,
"a {held}-record store with a line this build cannot read accounts for \
{}",
covered.records
);
assert_eq!(state, without_a_checkpoint(&paths));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_state_read_back_goes_through_the_checks_that_made_it() {
let root = scratch("checked-on-the-way-back");
let paths = a_recorded_run(&root, "r-checked");
settle(&paths, "build", "done");
let _ = fold_and_checkpoint(&paths);
let whole = without_a_checkpoint(&paths);
for (planting, named) in [
(
json!({"parks": {"build": {"by": "planner", "reason": " "}}}),
"reason",
),
(
json!({"sessions": {"build": {"token": "../somewhere-else", "branch": "work"}}}),
"session handle",
),
] {
let mut planted = document(&paths);
for (field, value) in planting.as_object().expect("one field to plant") {
planted["state"][field] = value.clone();
}
let refusal = serde_json::from_value::<Checkpoint>(planted.clone())
.expect_err("a value this crate refuses off a stream is refused off a document");
assert!(
refusal.to_string().contains(named),
"the refusal of {planting} does not name what it refused: {refusal}"
);
crate::ledger::write_json(&paths.checkpoint(), &planted).expect("written");
assert_eq!(
folded_as(&fold_and_checkpoint(&paths)),
whole,
"a document carrying {planting} was folded from"
);
}
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_store_with_nothing_readable_in_it_holds_no_instant_open() {
let root = scratch("nothing-readable");
let paths = a_run(&root, "r-illegible");
for line in ["not a record", "nor is this"] {
crate::ledger::append_line(&paths.journal(), line).expect("appended");
}
let state = folded_as(&fold_and_checkpoint(&paths));
let covered = super::readable(&paths).expect("a checkpoint").coverage;
assert_eq!(
covered.records, 2,
"a line the fold skipped was left unaccounted"
);
assert_eq!(covered.at, None, "a line this build cannot read was placed");
assert_eq!(state, without_a_checkpoint(&paths));
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_covered_record_changed_at_identical_length_folds_the_whole_store() {
let root = scratch("prefix-changed");
let paths = a_recorded_run(&root, "r-changed");
settle(&paths, "build", "done");
let _ = fold_and_checkpoint(&paths);
let marker = super::readable(&paths)
.expect("the checkpoint that read wrote")
.coverage;
let mut planted = document(&paths);
planted["state"]["outcomes"]["build"] = json!("carried-from-the-checkpoint");
seal_and_write(&paths, &mut planted);
settle(&paths, "ship", "done");
let store = paths.journal();
let held = std::fs::read(&store).expect("the run's journal");
let mut lines: Vec<Vec<u8>> = held
.split_inclusive(|byte| *byte == b'\n')
.map(<[u8]>::to_vec)
.collect();
let changed = lines[0].len() - 1;
lines[0] = b"#".repeat(changed).into_iter().chain(*b"\n").collect();
let mangled = lines.concat();
assert_eq!(
mangled.len(),
held.len(),
"this journey did not hold the store's length"
);
std::fs::write(&store, &mangled).expect("the journal is rewritten");
let after = super::readable(&paths)
.expect("the checkpoint is still there")
.coverage;
assert_eq!(
(after.bytes, after.records, &after.at, &after.streams),
(marker.bytes, marker.records, &marker.at, &marker.streams),
"the marker moved, so this journey is not about a store that did not"
);
let whole = without_a_checkpoint(&paths);
let read = folded_as(&fold_and_checkpoint(&paths));
assert_eq!(
read.get("outcomes")
.and_then(|outcomes| outcomes.get("build")),
None,
"a checkpoint over a rewritten prefix was folded from anyway: {read}"
);
assert_eq!(
read, whole,
"the state after a rewritten prefix is not the state the whole store folds to"
);
let _ = std::fs::remove_dir_all(&root);
}
const GOLDEN: &str = include_str!("../tests/golden/checkpoint-v4.json");
const GOLDEN_BYTES_DIGESTED: u128 = 0x1234_5678_9abc_def0_1234_5678_9abc_def0;
fn a_checkpoint() -> Checkpoint {
let mut graph = crate::graph::Graph::with_concurrency(4);
graph.insert(Node {
id: "build".into(),
persona: Some("engineer".into()),
task: Some("## What\ndo it".into()),
..Node::default()
});
let state = RunState {
graph,
recorded: BTreeMap::from([(
"build".to_string(),
crate::projection::Recorded::At(crate::graph::NodeStatus::Done),
)]),
last_write_at: Some(1_786_000_000_000),
strict: true,
..RunState::default()
};
let mut coverage = Coverage {
bytes: 8_192,
records: 42,
at: Some(Placed {
ts: "2026-09-08T12:00:00.000Z".into(),
stream: "golden-host-1".into(),
}),
streams: BTreeMap::from([("golden-host-1".to_string(), 41)]),
digest: 0,
};
coverage.digest = coverage.sealed_with(GOLDEN_BYTES_DIGESTED, &state);
Checkpoint {
schema_version: CHECKPOINT_SCHEMA_VERSION,
run_id: "golden".into(),
coverage,
state,
}
}
const GOLDEN_EARLIER: [(u32, &str); 3] = [
(1, include_str!("../tests/golden/checkpoint-v1.json")),
(2, include_str!("../tests/golden/checkpoint-v2.json")),
(3, include_str!("../tests/golden/checkpoint-v3.json")),
];
#[test]
fn the_documents_earlier_builds_wrote_are_refused_rather_than_read() {
for (version, earlier) in GOLDEN_EARLIER {
let refused = serde_json::from_str::<Checkpoint>(earlier).expect_err("it is refused");
assert!(
refused
.to_string()
.contains(&format!("schema_version {version}")),
"the refusal does not name the version it met: {refused}"
);
}
}
#[test]
fn a_schema_4_document_is_the_shape_the_golden_pins() {
let rendered = serde_json::to_string_pretty(&a_checkpoint()).expect("it serialises");
assert_eq!(
rendered.trim(),
GOLDEN.trim(),
"the checkpoint document changed shape. If that was deliberate, bump \
CHECKPOINT_SCHEMA_VERSION and update tests/golden/checkpoint-v4.json together"
);
}
#[test]
fn a_schema_4_document_round_trips_and_a_version_this_build_does_not_read_is_refused() {
let read: Checkpoint =
serde_json::from_str(GOLDEN).expect("the golden reads back into the types");
assert_eq!(
serde_json::to_string_pretty(&read).expect("it serialises"),
GOLDEN.trim(),
"a document this build wrote does not read back as itself"
);
assert_eq!(read.coverage, a_checkpoint().coverage);
assert_eq!(
read.coverage.digest,
read.coverage
.sealed_with(GOLDEN_BYTES_DIGESTED, &read.state),
"the document read back does not seal to what it carries"
);
let mut later: serde_json::Value = serde_json::from_str(GOLDEN).expect("it parses");
later["schema_version"] = json!(CHECKPOINT_SCHEMA_VERSION + 1);
let refused =
serde_json::from_value::<Checkpoint>(later).expect_err("a later version is refused");
assert!(
refused.to_string().contains("schema_version"),
"the refusal does not name what it refused: {refused}"
);
for refused in [
"not a digest",
"ABCDEF01234567890123456789ABCDEF",
"+bcdef01234567890123456789abcdef",
"bcdef01234567890123456789abcdef",
] {
let mut mangled: serde_json::Value = serde_json::from_str(GOLDEN).expect("it parses");
mangled["coverage"]["digest"] = json!(refused);
let refusal = serde_json::from_value::<Checkpoint>(mangled)
.expect_err("a digest no writer here produced is refused");
assert!(
refusal.to_string().contains("digest"),
"the refusal of '{refused}' does not name what it refused: {refusal}"
);
}
}
#[test]
fn a_session_the_writer_produced_is_one_the_reader_accepts() {
let written = json!({"token": "s-abc", "branch": "work/build"});
let read: crate::vcs::DispatchSession =
serde_json::from_value(written.clone()).expect("a session this crate accepts");
assert_eq!(
serde_json::to_value(&read).expect("it serialises"),
written,
"a session does not read back as the document it was written as"
);
for refused in [
json!({"token": "../somewhere-else", "branch": "work"}),
json!({"token": "s-abc", "branch": " "}),
json!({"token": "s-abc", "branch": "work", "extra": 1}),
] {
assert!(
serde_json::from_value::<crate::vcs::DispatchSession>(refused.clone()).is_err(),
"{refused} was accepted as a session"
);
}
}
#[test]
fn a_park_the_writer_produced_is_one_the_reader_accepts() {
let written = crate::edits::Park::of(crate::channel::Author::planner(), Some("waiting"));
let document = serde_json::to_value(&written).expect("it serialises");
assert_eq!(
serde_json::from_value::<crate::edits::Park>(document.clone())
.expect("a park this crate accepts"),
written,
"a park does not read back as the document it was written as"
);
let none = serde_json::to_value(crate::edits::Park::of(
crate::channel::Author::planner(),
None,
))
.expect("it serialises");
assert_eq!(none, json!({"by": "planner"}));
for refused in [
json!({"by": "planner", "reason": " "}),
json!({"by": "planner", "reason": ""}),
json!({"by": "planner", "extra": 1}),
] {
assert!(
serde_json::from_value::<crate::edits::Park>(refused.clone()).is_err(),
"{refused} was accepted as a park"
);
}
}
#[test]
fn a_state_edited_without_the_seal_moving_folds_the_whole_store() {
let root = scratch("state-edited");
let paths = a_recorded_run(&root, "r-edited");
settle(&paths, "build", "done");
let _ = fold_and_checkpoint(&paths);
settle(&paths, "ship", "done");
let whole = without_a_checkpoint(&paths);
let mut edited = document(&paths);
edited["state"]["outcomes"]["build"] = json!("carried-from-the-checkpoint");
crate::ledger::write_json(&paths.checkpoint(), &edited).expect("written");
let read = folded_as(&fold_and_checkpoint(&paths));
assert_eq!(
read.get("outcomes")
.and_then(|outcomes| outcomes.get("build")),
None,
"a fold edited in the document was served: {read}"
);
assert_eq!(read, whole);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_run_that_has_recorded_nothing_leaves_no_checkpoint() {
let root = scratch("nothing-recorded");
let paths = a_run(&root, "r-empty");
let state = fold_and_checkpoint(&paths);
assert!(state.strict, "the empty fold is not the fold's own start");
assert!(state.graph.is_empty());
assert!(
!paths.checkpoint().exists(),
"a run with no store was given a checkpoint of it"
);
let _ = std::fs::remove_dir_all(&root);
}
#[test]
fn a_refold_takes_what_the_store_grew_by_rather_than_the_whole_journal() {
let root = scratch("refold");
let paths = a_recorded_run(&root, "r-refold");
for nth in 0..20 {
settle(&paths, if nth % 2 == 0 { "build" } else { "ship" }, "done");
}
let mut projected = Projected::open(&paths);
assert_eq!(
projected.took(),
22,
"the loop did not open on the whole store"
);
let covered = projected.coverage().records;
assert!(
covered > 0,
"the loop accounted for none of the 22-record store it opened"
);
settle(&paths, "build", "failed");
projected.refresh(&paths);
let took = projected.took();
assert_eq!(
took,
23 - covered,
"a re-fold took {took} records of a 23-record store"
);
assert!(took <= 4, "the open instant is not a bounded run: {took}");
assert_eq!(
folded_as(&projected),
without_a_checkpoint(&paths),
"the loop's state and a full fold's disagree"
);
let _ = std::fs::remove_dir_all(&root);
}
}