use std::collections::HashMap;
use serde::Serialize;
use crate::runtime::kernel::wire::ConfigDefaults;
use crate::runtime::kernel::wire::checkpoint::KernelCheckpoint;
use crate::runtime::kernel::wire::effect::{
EffectKind, EffectOutcome, EffectSuccess, ProviderOutcome, SpawnTasksEffect,
};
use crate::runtime::kernel::wire::record::{
KernelRecord, NormalizedPayload, RecordError, verify_record_chain,
};
use crate::runtime::kernel::wire::restore::{RestoredOperation, restore_operation};
use crate::runtime::kernel::wire::transaction::InMemoryRecordIndex;
pub const UNATTRIBUTED_SEGMENT: &str = "(unattributed)";
const DEFERRED_ALWAYS: &str = "c4.parent_chain: parent links are not journaled; an orphan spawn \
cannot resolve (no outstanding effect), which C3's re-plan enforces structurally";
const DEFERRED_WITHOUT_CHECKPOINT: &str = "c5b.launch_token_ledger: the durable LaunchToken \
ledger lives in checkpoints; provide --checkpoint to check token reuse across TaskLaunch \
payloads (the journal-direct shadow — (task_id, attempt_id) pair uniqueness — is checked \
under C4)";
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct RuleReport {
pub rule: String,
pub verdict: Verdict,
pub detail: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum Verdict {
Pass,
Fail,
Degraded,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct DegradedHop {
pub ordinal: usize,
pub step_seq: Option<u64>,
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct SegmentReport {
pub operation_id: String,
pub hops: usize,
pub degraded_hops: Vec<DegradedHop>,
pub rules: Vec<RuleReport>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct ValidationReport {
pub segments: Vec<SegmentReport>,
pub unparseable_records: usize,
#[serde(default)]
pub cross_checks: Vec<RuleReport>,
#[serde(default)]
pub unparseable_events: usize,
#[serde(default)]
pub session_events: Option<usize>,
#[serde(default)]
pub checkpoint_checks: Vec<RuleReport>,
#[serde(default)]
pub unparseable_checkpoints: usize,
#[serde(default)]
pub checkpoints: Option<usize>,
pub deferred: Vec<String>,
}
impl ValidationReport {
pub fn has_violations(&self) -> bool {
self.segments
.iter()
.flat_map(|segment| segment.rules.iter())
.chain(self.cross_checks.iter())
.chain(self.checkpoint_checks.iter())
.any(|rule| rule.verdict == Verdict::Fail)
}
pub fn has_violations_for(&self, operation_id: &str) -> bool {
self.segments
.iter()
.filter(|segment| segment.operation_id == operation_id)
.flat_map(|segment| segment.rules.iter())
.chain(self.cross_checks.iter())
.chain(self.checkpoint_checks.iter())
.any(|rule| rule.verdict == Verdict::Fail)
}
pub fn for_operation(&self, operation_id: &str) -> Self {
let mut report = self.clone();
report
.segments
.retain(|segment| segment.operation_id == operation_id);
report
}
pub fn has_insufficient_evidence(&self) -> bool {
self.unparseable_records > 0
|| self.unparseable_events > 0
|| self.unparseable_checkpoints > 0
|| matches!(self.session_events, Some(0))
|| matches!(self.checkpoints, Some(0))
}
pub fn exit_code(&self) -> i32 {
if self.has_violations() {
1
} else if self.has_insufficient_evidence() || self.segments.is_empty() {
2
} else {
0
}
}
}
enum Hop {
Complete(KernelRecord),
Degraded(DegradedRecord),
}
struct DegradedRecord {
ordinal: usize,
operation_id: Option<String>,
input_id: Option<String>,
step_seq: Option<u64>,
previous_record_digest: Option<String>,
record_digest: Option<String>,
reason: String,
integrity_failure: bool,
}
impl DegradedRecord {
fn marking(&self) -> DegradedHop {
DegradedHop {
ordinal: self.ordinal,
step_seq: self.step_seq,
reason: self.reason.clone(),
}
}
}
impl Hop {
fn operation_id(&self) -> Option<&str> {
match self {
Self::Complete(record) => Some(record.operation_id().as_str()),
Self::Degraded(degraded) => degraded.operation_id.as_deref(),
}
}
fn step_seq(&self) -> Option<u64> {
match self {
Self::Complete(record) => Some(record.step_seq().get()),
Self::Degraded(degraded) => degraded.step_seq,
}
}
}
pub fn validate_journal<B: AsRef<[u8]>>(blobs: &[B]) -> ValidationReport {
validate_with_session_log(blobs, &[] as &[Vec<Vec<u8>>])
}
pub fn validate_with_session_log<J, S>(
journal_blobs: &[J],
session_streams: &[Vec<S>],
) -> ValidationReport
where
J: AsRef<[u8]>,
S: AsRef<[u8]>,
{
validate_with_checkpoint(journal_blobs, session_streams, &[] as &[Vec<u8>], false)
}
pub fn validate_with_checkpoint<J, S, C>(
journal_blobs: &[J],
session_streams: &[Vec<S>],
checkpoint_blobs: &[C],
strict: bool,
) -> ValidationReport
where
J: AsRef<[u8]>,
S: AsRef<[u8]>,
C: AsRef<[u8]>,
{
let (outcomes, unparseable_records, segment_records) = validate_journal_plane(journal_blobs);
let mut streams: Vec<SessionStream> = Vec::with_capacity(session_streams.len());
for stream_blobs in session_streams {
let mut events = Vec::with_capacity(stream_blobs.len());
let mut unparseable_events = 0;
for blob in stream_blobs {
match classify_session_event(blob.as_ref()) {
Some(event) => events.push(event),
None => unparseable_events += 1,
}
}
streams.push(SessionStream {
events,
unparseable_events,
});
}
let session_plane_provided = !session_streams.is_empty();
let session_events =
session_plane_provided.then(|| streams.iter().map(|stream| stream.events.len()).sum());
let unparseable_events = streams.iter().map(|stream| stream.unparseable_events).sum();
let cross_checks = if session_plane_provided {
let mut checks = check_c6(&streams, &outcomes);
checks.push(check_c8(&streams, &outcomes));
checks
} else {
Vec::new()
};
let checkpoint_plane_provided = !checkpoint_blobs.is_empty();
let (checkpoint_checks, unparseable_checkpoints, checkpoints) = if checkpoint_plane_provided {
let (checks, unparseable, parsed) = check_c5(&segment_records, checkpoint_blobs, strict);
(checks, unparseable, Some(parsed))
} else {
(Vec::new(), 0, None)
};
let mut deferred = vec![DEFERRED_ALWAYS.to_string()];
if !checkpoint_plane_provided {
deferred.push(DEFERRED_WITHOUT_CHECKPOINT.to_string());
}
ValidationReport {
segments: outcomes.into_iter().map(|outcome| outcome.report).collect(),
unparseable_records,
cross_checks,
unparseable_events,
session_events,
checkpoint_checks,
unparseable_checkpoints,
checkpoints,
deferred,
}
}
fn validate_journal_plane<B: AsRef<[u8]>>(
blobs: &[B],
) -> (
Vec<SegmentOutcome>,
usize,
HashMap<String, Vec<KernelRecord>>,
) {
let mut hops: Vec<Hop> = Vec::with_capacity(blobs.len());
let mut unparseable_records = 0;
for (ordinal, blob) in blobs.iter().enumerate() {
match classify(ordinal, blob.as_ref()) {
Some(hop) => hops.push(hop),
None => unparseable_records += 1,
}
}
let mut segments: HashMap<String, Vec<Hop>> = HashMap::new();
let mut complete: HashMap<String, Vec<KernelRecord>> = HashMap::new();
for hop in hops {
let key = hop
.operation_id()
.map(str::to_string)
.unwrap_or_else(|| UNATTRIBUTED_SEGMENT.to_string());
if let Hop::Complete(record) = &hop {
complete
.entry(key.clone())
.or_default()
.push(record.clone());
}
segments.entry(key).or_default().push(hop);
}
let mut keys: Vec<String> = segments.keys().cloned().collect();
keys.sort();
let reports = keys
.iter()
.map(|key| validate_segment(key, segments.remove(key).unwrap_or_default()))
.collect();
(reports, unparseable_records, complete)
}
pub struct SessionStream {
pub events: Vec<EvidenceEvent>,
pub unparseable_events: usize,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum EvidenceEvent {
RunStarted { route_id: Option<String> },
ProviderAttempt {
effect_id: Option<String>,
request_fingerprint: Option<String>,
route_id: Option<String>,
status: Option<String>,
},
PromptMeasured {
effect_id: Option<String>,
request_fingerprint: Option<String>,
},
LlmCompleted {
effect_id: Option<String>,
invocation_id: Option<String>,
},
Other,
}
fn classify_session_event(bytes: &[u8]) -> Option<EvidenceEvent> {
let value: serde_json::Value = serde_json::from_slice(bytes).ok()?;
let object = value.as_object()?;
let string = |key: &str| {
object
.get(key)
.and_then(serde_json::Value::as_str)
.map(str::to_string)
};
let kind = string("kind");
let event = match kind.as_deref() {
Some("run_started") => EvidenceEvent::RunStarted {
route_id: route_id_of(&value),
},
Some("provider_attempt") => EvidenceEvent::ProviderAttempt {
effect_id: string("effect_id"),
request_fingerprint: string("request_fingerprint"),
route_id: route_id_of(&value),
status: string("status"),
},
Some("prompt_measured") => EvidenceEvent::PromptMeasured {
effect_id: string("effect_id"),
request_fingerprint: object
.get("measurement")
.and_then(|measurement| {
measurement
.get("requestFingerprint")
.or_else(|| measurement.get("request_fingerprint"))
})
.and_then(serde_json::Value::as_str)
.map(str::to_string),
},
Some("llm_completed") => EvidenceEvent::LlmCompleted {
effect_id: string("effect_id"),
invocation_id: string("invocation_id"),
},
_ => EvidenceEvent::Other,
};
Some(event)
}
fn route_id_of(event: &serde_json::Value) -> Option<String> {
let route = event.get("route")?;
route
.get("routeId")
.or_else(|| route.get("route_id"))
.and_then(serde_json::Value::as_str)
.map(str::to_string)
}
fn check_c6(streams: &[SessionStream], outcomes: &[SegmentOutcome]) -> Vec<RuleReport> {
let segments: HashMap<&str, &SegmentOutcome> = outcomes
.iter()
.filter(|outcome| outcome.report.operation_id != UNATTRIBUTED_SEGMENT)
.map(|outcome| (outcome.report.operation_id.as_str(), outcome))
.collect();
let mut correspondence = ClauseAccumulator::default();
let mut fingerprints = ClauseAccumulator::default();
let mut routes = ClauseAccumulator::default();
let mut total_attempts = 0usize;
for (index, stream) in streams.iter().enumerate() {
let measured: std::collections::HashSet<&str> = stream
.events
.iter()
.filter_map(|event| match event {
EvidenceEvent::PromptMeasured {
request_fingerprint,
..
} => request_fingerprint.as_deref(),
_ => None,
})
.collect();
let mut baseline: Option<&str> = None;
let mut last_pinned: Option<&str> = None;
for event in &stream.events {
match event {
EvidenceEvent::RunStarted { route_id } => {
if let Some(new_route) = route_id.as_deref() {
if let Some(previous) = last_pinned
&& previous != new_route
{
routes.degraded(format!(
"stream #{index}: run resumed on route {new_route} (was \
{previous}) — a cross-resume change, degraded per Q3"
));
}
baseline = Some(new_route);
last_pinned = Some(new_route);
} else {
baseline = None;
}
}
EvidenceEvent::ProviderAttempt {
effect_id,
request_fingerprint,
route_id,
..
} => {
total_attempts += 1;
let label = effect_id.as_deref().unwrap_or("(no effect_id)");
match effect_id.as_deref() {
None => correspondence.violation(format!(
"stream #{index}: provider_attempt without effect_id — the writers \
mint it from the kernel effect, so a keyless attempt is forged \
evidence"
)),
Some(effect) => match parse_effect_step(effect) {
None => correspondence.violation(format!(
"stream #{index}: attempt names {effect}, which is not in the \
`operation:step:N:effect:M` vocabulary"
)),
Some((operation, step)) => match segments.get(operation) {
None => correspondence.degraded(format!(
"stream #{index}: attempt names {effect}, but operation \
{operation} has no journal segment (the journal may cover \
a subset of the session)"
)),
Some(outcome) => match &outcome.effects {
None => correspondence.degraded(format!(
"stream #{index}: segment {operation} could not be \
re-planned, so {effect}'s publication is unverifiable"
)),
Some(effects) if effects.published.contains(effect) => {
correspondence.checked += 1;
}
Some(effects) if step > effects.max_step => {
correspondence.degraded(format!(
"stream #{index}: attempt names {effect} at step \
{step}, past the journal's tip (step {}) — a \
prefix cannot disprove it",
effects.max_step
));
}
Some(_) => correspondence.violation(format!(
"stream #{index}: attempt names {effect}, but the \
deterministic re-plan of {operation} never published \
it — the attempt is unmoored from the journal"
)),
},
},
},
}
match request_fingerprint.as_deref() {
None => fingerprints.violation(format!(
"stream #{index}: provider_attempt {label} without \
request_fingerprint — the writers require it (G2)"
)),
Some(fingerprint) if measured.contains(fingerprint) => {
fingerprints.checked += 1;
}
Some(fingerprint) => fingerprints.violation(format!(
"stream #{index}: attempt {label} carries fingerprint \
{fingerprint}, but no prompt_measured in this session carries it"
)),
}
match (route_id.as_deref(), baseline) {
(Some(route), Some(pinned)) if route != pinned => {
routes.violation(format!(
"stream #{index}: attempt {label} ran on route {route} inside \
a run pinned to {pinned} — an in-run route change is a \
violation (Q3)"
))
}
(Some(_), Some(_)) => routes.checked += 1,
_ => routes.unverifiable += 1,
}
}
_ => {}
}
}
}
vec![
correspondence.report(
"C6.1",
total_attempts,
"attempt↔journal effect correspondence",
),
fingerprints.report(
"C6.2",
total_attempts,
"attempt fingerprint↔prompt_measured join",
),
routes.report("C6.3", total_attempts, "in-run route stability"),
]
}
#[derive(Default)]
struct ClauseAccumulator {
checked: usize,
unverifiable: usize,
violations: Vec<String>,
degraded_notes: Vec<String>,
}
impl ClauseAccumulator {
fn violation(&mut self, detail: String) {
self.violations.push(detail);
}
fn degraded(&mut self, note: String) {
self.degraded_notes.push(note);
}
fn report(self, rule: &str, total_attempts: usize, what: &str) -> RuleReport {
let rule = rule.to_string();
if !self.violations.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Fail,
detail: self.violations.join("; "),
};
}
if total_attempts == 0 {
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: format!(
"no provider_attempt events in any stream — a pre-0.2.63 log carries none \
(C7); {what} unchecked"
),
};
}
if !self.degraded_notes.is_empty() || self.unverifiable > 0 {
let mut detail = self.degraded_notes.join("; ");
if self.unverifiable > 0 {
if !detail.is_empty() {
detail.push_str("; ");
}
detail.push_str(&format!(
"{} attempt(s) unverifiable (no pinned run route)",
self.unverifiable
));
}
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: format!("{} attempt(s) verified for {what}; {detail}", self.checked),
};
}
RuleReport {
rule,
verdict: Verdict::Pass,
detail: format!(
"{} attempt(s) verified — {what} holds across every stream",
self.checked
),
}
}
}
fn check_c8(streams: &[SessionStream], outcomes: &[SegmentOutcome]) -> RuleReport {
let rule = "C8".to_string();
let segments: HashMap<&str, &SegmentOutcome> = outcomes
.iter()
.filter(|outcome| outcome.report.operation_id != UNATTRIBUTED_SEGMENT)
.map(|outcome| (outcome.report.operation_id.as_str(), outcome))
.collect();
let mut llm_completed_events = 0usize;
let mut checked = 0usize;
let mut trivial = 0usize;
let mut unverifiable = 0usize;
let mut violations: Vec<String> = Vec::new();
let mut degraded_notes: Vec<String> = Vec::new();
for (index, stream) in streams.iter().enumerate() {
for event in &stream.events {
let EvidenceEvent::LlmCompleted {
effect_id,
invocation_id,
} = event
else {
continue;
};
llm_completed_events += 1;
let (Some(head), Some(selected)) = (invocation_id.as_deref(), effect_id.as_deref())
else {
unverifiable += 1;
continue;
};
if head == selected {
trivial += 1;
continue;
}
let (Some((head_op, head_step)), Some((selected_op, selected_step))) =
(parse_effect_step(head), parse_effect_step(selected))
else {
violations.push(format!(
"stream #{index}: llm_completed claims invocation {head} → {selected}, but \
one of the pair is not in the `operation:step:N:effect:M` vocabulary"
));
continue;
};
if head_op != selected_op {
violations.push(format!(
"stream #{index}: llm_completed claims invocation {head} → {selected} — an \
invocation chain cannot cross operations"
));
continue;
}
if selected_step <= head_step {
violations.push(format!(
"stream #{index}: llm_completed claims invocation {head} → {selected}, but \
the selected effect does not follow the chain head"
));
continue;
}
let Some(outcome) = segments.get(head_op) else {
degraded_notes.push(format!(
"stream #{index}: operation {head_op} has no journal segment (the journal \
may cover a subset of the session)"
));
continue;
};
let Some(effects) = &outcome.effects else {
degraded_notes.push(format!(
"stream #{index}: segment {head_op} could not be re-planned, so the \
invocation {head} → {selected} is unverifiable"
));
continue;
};
match effects.resolutions.get(head) {
Some(ResolutionFact::Completed) | Some(ResolutionFact::Other) => {
violations.push(format!(
"stream #{index}: llm_completed claims invocation {head} → {selected}, \
but {head} resolved to completion — a completed effect closes its \
invocation; nothing chains from it"
));
continue;
}
Some(ResolutionFact::Overflow) | Some(ResolutionFact::Failed) => {}
None if effects.published.contains(head) => degraded_notes.push(format!(
"stream #{index}: chain head {head} is published but its resolution has \
not landed in the journal (the planes are not synchronised)"
)),
None if head_step > effects.max_step => degraded_notes.push(format!(
"stream #{index}: chain head {head} claims step {head_step}, past the \
journal's tip (step {})",
effects.max_step
)),
None => {
violations.push(format!(
"stream #{index}: llm_completed claims invocation head {head}, but the \
deterministic re-plan of {head_op} never published it"
));
continue;
}
}
if effects.resolutions.contains_key(selected) {
checked += 1;
} else if effects.published.contains(selected) || selected_step > effects.max_step {
degraded_notes.push(format!(
"stream #{index}: selected effect {selected} is not resolved in the \
journal (the planes are not synchronised)"
));
} else {
violations.push(format!(
"stream #{index}: llm_completed selects {selected}, but the deterministic \
re-plan of {selected_op} never published it"
));
}
}
}
if !violations.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Fail,
detail: violations.join("; "),
};
}
if llm_completed_events == 0 {
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: "no llm_completed events in any stream — invocation adjacency unchecked"
.to_string(),
};
}
if !degraded_notes.is_empty() || unverifiable > 0 {
let mut detail = degraded_notes.join("; ");
if unverifiable > 0 {
if !detail.is_empty() {
detail.push_str("; ");
}
detail.push_str(&format!(
"{unverifiable} llm_completed event(s) without invocation_id/effect_id \
(pre-0.2.63 fields — C7)"
));
}
if !detail.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: format!(
"{checked} retried invocation(s) verified, {trivial} first-try chain(s) \
closed; {detail}"
),
};
}
}
RuleReport {
rule,
verdict: Verdict::Pass,
detail: format!(
"{checked} retried invocation(s) verified end-to-end, {trivial} first-try chain(s) \
closed"
),
}
}
fn classify(ordinal: usize, bytes: &[u8]) -> Option<Hop> {
let error = match KernelRecord::from_record_bytes(bytes) {
Ok(record) => return Some(Hop::Complete(record)),
Err(error) => error,
};
let value: serde_json::Value = serde_json::from_slice(bytes).ok()?;
let object = value.as_object()?;
let string = |key: &str| {
object
.get(key)
.and_then(serde_json::Value::as_str)
.map(str::to_string)
};
let step_seq = object.get("step_seq").and_then(|value| {
value
.as_u64()
.or_else(|| value.as_str().and_then(|text| text.parse().ok()))
});
let degraded = DegradedRecord {
ordinal,
operation_id: string("operation_id"),
input_id: string("input_id"),
step_seq,
previous_record_digest: string("previous_record_digest"),
record_digest: string("record_digest"),
reason: format!("{}: {}", error.code().as_str(), error.message()),
integrity_failure: matches!(error, RecordError::DigestMismatch(_)),
};
if degraded.operation_id.is_some() || degraded.step_seq.is_some() {
Some(Hop::Degraded(degraded))
} else {
None
}
}
fn validate_segment(operation_id: &str, mut hops: Vec<Hop>) -> SegmentOutcome {
hops.sort_by_key(|hop| {
(
hop.step_seq().unwrap_or(u64::MAX),
match hop {
Hop::Complete(_) => 0usize,
Hop::Degraded(degraded) => degraded.ordinal,
},
)
});
let degraded_hops: Vec<DegradedHop> = hops
.iter()
.filter_map(|hop| match hop {
Hop::Degraded(degraded) => Some(degraded.marking()),
Hop::Complete(_) => None,
})
.collect();
let hop_count = hops.len();
let c1 = check_c1(&hops);
let replan = replan_segment(&hops, &c1);
let c2 = check_c2(&hops);
let c3 = render_c3(&replan);
let c4 = check_c4(&hops, operation_id);
let effects = match &replan {
Replan::Restored(restored) => {
let mut published: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut resolutions: HashMap<String, ResolutionFact> = HashMap::new();
for hop in &hops {
if let Some((effect_id, fact)) = resolution_of(hop) {
published.insert(effect_id.clone());
resolutions.insert(effect_id, fact);
}
}
published.extend(
restored
.transaction
.pending_effects()
.map(|effect| effect.effect_id.as_str().to_string()),
);
let max_step = hops
.iter()
.filter_map(|hop| hop.step_seq())
.max()
.unwrap_or(0);
Some(SegmentEffects {
max_step,
published,
resolutions,
})
}
_ => None,
};
SegmentOutcome {
report: SegmentReport {
operation_id: operation_id.to_string(),
hops: hop_count,
degraded_hops,
rules: vec![c1, c2, c3, c4],
},
effects,
}
}
struct SegmentOutcome {
report: SegmentReport,
effects: Option<SegmentEffects>,
}
struct SegmentEffects {
max_step: u64,
published: std::collections::HashSet<String>,
resolutions: HashMap<String, ResolutionFact>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ResolutionFact {
Completed,
Overflow,
Failed,
Other,
}
fn resolution_of(hop: &Hop) -> Option<(String, ResolutionFact)> {
let Hop::Complete(record) = hop else {
return None;
};
let input = record.normalized_input().ok()?;
let NormalizedPayload::ResolveEffect(resolve) = &input.input else {
return None;
};
let fact = match &resolve.outcome {
EffectOutcome::Failed(_) => ResolutionFact::Failed,
EffectOutcome::Succeeded(success) => match &success.result {
EffectSuccess::Provider(provider) => match &provider.outcome {
ProviderOutcome::Completed(_) => ResolutionFact::Completed,
ProviderOutcome::ContextOverflow(_) => ResolutionFact::Overflow,
},
_ => ResolutionFact::Other,
},
};
Some((resolve.effect_id.as_str().to_string(), fact))
}
fn check_c1(hops: &[Hop]) -> RuleReport {
let rule = "C1".to_string();
if hops.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: "no records in this segment".to_string(),
};
}
let all_complete = hops.iter().all(|hop| matches!(hop, Hop::Complete(_)));
if all_complete {
let records: Vec<KernelRecord> = hops
.iter()
.filter_map(|hop| match hop {
Hop::Complete(record) => Some(record.clone()),
Hop::Degraded(_) => None,
})
.collect();
return match verify_record_chain(&records) {
Ok(genesis_digest) => RuleReport {
rule,
verdict: Verdict::Pass,
detail: format!(
"{} record(s), genesis {genesis_digest}, every link verified",
records.len()
),
},
Err(error) => RuleReport {
rule,
verdict: Verdict::Fail,
detail: format!("{}: {}", error.code().as_str(), error.message()),
},
};
}
let mut broken: Vec<String> = hops
.iter()
.filter_map(|hop| match hop {
Hop::Degraded(record) if record.integrity_failure => Some(record.reason.clone()),
_ => None,
})
.collect();
let mut unverifiable_links = 0usize;
let mut previous: Option<(&Hop, Option<&KernelRecord>)> = None;
for hop in hops {
let step = hop.step_seq();
let (prev_digest, _) = digests_of(hop);
if let Some((previous_hop, previous_complete)) = previous {
let previous_step = previous_hop.step_seq();
let previous_digest = digests_of(previous_hop).1;
match (prev_digest, previous_digest) {
(Some(expected), Some(actual)) if expected != actual => broken.push(format!(
"hop at step {} expects head {expected}, but its predecessor's digest is \
{actual}",
step.map_or("?".to_string(), |seq| seq.to_string()),
)),
(None, _) => unverifiable_links += 1,
(_, None) => unverifiable_links += 1,
_ => {}
}
match (step, previous_step) {
(Some(step), Some(previous_step)) if step != previous_step + 1 => broken.push(
format!("hop is step {step}, but its predecessor is step {previous_step}"),
),
(Some(_), Some(_)) => {}
_ => unverifiable_links += 1,
}
if let (Hop::Complete(record), Some(previous_record)) = (hop, previous_complete)
&& let Err(error) = record.verify_follows(Some(previous_record))
{
broken.push(format!("{}: {}", error.code().as_str(), error.message()));
}
} else if let Hop::Complete(record) = hop
&& let Err(error) = record.verify_follows(None)
{
broken.push(format!("{}: {}", error.code().as_str(), error.message()));
}
previous = Some((
hop,
match hop {
Hop::Complete(record) => Some(record),
Hop::Degraded(_) => None,
},
));
}
if !broken.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Fail,
detail: broken.join("; "),
};
}
RuleReport {
rule,
verdict: Verdict::Degraded,
detail: format!(
"partial chain: every surviving link verified, {unverifiable_links} link(s) \
unverifiable across degraded hop(s)"
),
}
}
fn check_c2(hops: &[Hop]) -> RuleReport {
let rule = "C2".to_string();
let mut by_input: HashMap<&str, &str> = HashMap::new();
let mut conflicts: Vec<String> = Vec::new();
let mut retries = 0usize;
let mut unverifiable = 0usize;
for hop in hops {
let (input_id, record_digest) = match hop {
Hop::Complete(record) => (
Some(record.input_id().as_str()),
Some(record.record_digest().as_str()),
),
Hop::Degraded(degraded) => (
degraded.input_id.as_deref(),
degraded.record_digest.as_deref(),
),
};
let Some(input_id) = input_id else { continue };
let Some(digest) = record_digest else {
unverifiable += 1;
continue;
};
match by_input.get(input_id) {
Some(existing) if *existing != digest => conflicts.push(format!(
"input {input_id} produced two different records ({existing} and {digest}); a \
retry must reach the same record"
)),
Some(_) => retries += 1,
None => {
by_input.insert(input_id, digest);
}
}
}
if !conflicts.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Fail,
detail: conflicts.join("; "),
};
}
if unverifiable > 0 {
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: format!(
"{} unique input(s), {retries} idempotent retry hit(s); {unverifiable} degraded \
hop(s) could not be compared",
by_input.len(),
),
};
}
RuleReport {
rule,
verdict: Verdict::Pass,
detail: format!(
"{} unique input(s), {retries} idempotent retry hit(s), no divergent duplicates",
by_input.len(),
),
}
}
enum Replan {
Unavailable(&'static str),
Failed(String),
Restored(RestoredOperation),
}
fn replan_segment(hops: &[Hop], c1: &RuleReport) -> Replan {
if hops.iter().any(|hop| matches!(hop, Hop::Degraded(_))) {
return Replan::Unavailable(
"re-plan requires complete records; this segment has degraded hops",
);
}
if c1.verdict == Verdict::Fail {
return Replan::Unavailable(
"C1 failed; a re-plan over a broken chain would only re-report that break",
);
}
let records: Vec<KernelRecord> = hops
.iter()
.filter_map(|hop| match hop {
Hop::Complete(record) => Some(record.clone()),
Hop::Degraded(_) => None,
})
.collect();
if records.is_empty() {
return Replan::Unavailable("no records in this segment");
}
match restore_operation(
None,
&records,
ConfigDefaults::default(),
InMemoryRecordIndex::from_records(&records),
) {
Ok(restored) => Replan::Restored(restored),
Err(fault) => Replan::Failed(format!("{}: {}", fault.code.as_str(), fault.message)),
}
}
fn render_c3(replan: &Replan) -> RuleReport {
let rule = "C3".to_string();
match replan {
Replan::Unavailable(reason) => RuleReport {
rule,
verdict: Verdict::Degraded,
detail: (*reason).to_string(),
},
Replan::Failed(fault) => RuleReport {
rule,
verdict: Verdict::Fail,
detail: fault.clone(),
},
Replan::Restored(restored) => RuleReport {
rule,
verdict: Verdict::Pass,
detail: format!(
"re-planned {} record(s) from genesis; every durable record digest reproduced",
restored.cost.records_before_checkpoint
),
},
}
}
fn check_c4(hops: &[Hop], operation_id: &str) -> RuleReport {
let rule = "C4".to_string();
struct LaunchFact {
task_id: String,
attempt_id: String,
step_seq: u64,
effect_id: String,
}
let mut launches: Vec<LaunchFact> = Vec::new();
let mut unreadable_inputs = 0usize;
for hop in hops {
let Hop::Complete(record) = hop else { continue };
let input = match record.normalized_input() {
Ok(input) => input,
Err(_) => {
unreadable_inputs += 1;
continue;
}
};
let NormalizedPayload::ResolveEffect(resolve) = &input.input else {
continue;
};
let EffectOutcome::Succeeded(success) = &resolve.outcome else {
continue;
};
let EffectSuccess::TasksSpawned(spawned) = &success.result else {
continue;
};
for attempt in &spawned.attempts {
launches.push(LaunchFact {
task_id: attempt.task_id.as_str().to_string(),
attempt_id: attempt.attempt_id.as_str().to_string(),
step_seq: record.step_seq().get(),
effect_id: resolve.effect_id.as_str().to_string(),
});
}
}
let mut violations: Vec<String> = Vec::new();
let mut seen: HashMap<(&str, &str), u64> = HashMap::new();
for fact in &launches {
let pair = (fact.task_id.as_str(), fact.attempt_id.as_str());
if let Some(first_step) = seen.insert(pair, fact.step_seq) {
violations.push(format!(
"task {} attempt {} launched at steps {first_step} and {}; the launch token is \
derived from that pair, so a repeated pair is a reused LaunchToken",
fact.task_id, fact.attempt_id, fact.step_seq,
));
}
match parse_effect_step(&fact.effect_id) {
Some((effect_operation, effect_step)) => {
if effect_operation != operation_id {
violations.push(format!(
"task {} launch at step {} resolves effect {} of another operation — \
causation cannot cross operations",
fact.task_id, fact.step_seq, fact.effect_id,
));
} else if effect_step >= fact.step_seq {
violations.push(format!(
"task {} launch resolved at step {} names an effect published at step \
{effect_step} — the resolution precedes the publication",
fact.task_id, fact.step_seq,
));
}
}
None => violations.push(format!(
"task {} launch at step {} names effect {}, which is not in the \
`operation:step:N:effect:M` vocabulary",
fact.task_id, fact.step_seq, fact.effect_id,
)),
}
}
if !violations.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Fail,
detail: violations.join("; "),
};
}
let degraded_hops = hops
.iter()
.filter(|hop| matches!(hop, Hop::Degraded(_)))
.count();
if degraded_hops > 0 || unreadable_inputs > 0 {
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail: format!(
"{} launch(es) checked; {degraded_hops} degraded hop(s) and \
{unreadable_inputs} unreadable input(s) could hide further launches",
launches.len(),
),
};
}
RuleReport {
rule,
verdict: Verdict::Pass,
detail: format!(
"{} launch(es), every (task_id, attempt_id) pair unique, every spawn resolution \
names an earlier step of this operation",
launches.len(),
),
}
}
fn parse_effect_step(effect_id: &str) -> Option<(&str, u64)> {
let (before_effect, _) = effect_id.rsplit_once(":effect:")?;
let (operation, step) = before_effect.rsplit_once(":step:")?;
Some((operation, step.parse().ok()?))
}
fn digests_of(hop: &Hop) -> (Option<&str>, Option<&str>) {
match hop {
Hop::Complete(record) => (
record
.previous_record_digest()
.map(|digest| digest.as_str()),
Some(record.record_digest().as_str()),
),
Hop::Degraded(degraded) => (
degraded.previous_record_digest.as_deref(),
degraded.record_digest.as_deref(),
),
}
}
struct C5Clauses {
violations: Vec<String>,
degradations: Vec<String>,
confirmations: Vec<String>,
}
impl C5Clauses {
fn new() -> Self {
Self {
violations: Vec::new(),
degradations: Vec::new(),
confirmations: Vec::new(),
}
}
fn violation(&mut self, detail: String) {
self.violations.push(detail);
}
fn degraded(&mut self, detail: String) {
self.degradations.push(detail);
}
fn confirms(&mut self, detail: String) {
self.confirmations.push(detail);
}
fn report(self, rule: &str) -> RuleReport {
let rule = rule.to_string();
if !self.violations.is_empty() {
return RuleReport {
rule,
verdict: Verdict::Fail,
detail: self.violations.join("; "),
};
}
if !self.degradations.is_empty() {
let mut detail = self.degradations.join("; ");
if !self.confirmations.is_empty() {
detail.push_str("; ");
detail.push_str(&self.confirmations.join("; "));
}
return RuleReport {
rule,
verdict: Verdict::Degraded,
detail,
};
}
RuleReport {
rule,
verdict: Verdict::Pass,
detail: if self.confirmations.is_empty() {
"nothing to anchor".to_string()
} else {
self.confirmations.join("; ")
},
}
}
}
fn check_c5<C: AsRef<[u8]>>(
segments: &HashMap<String, Vec<KernelRecord>>,
blobs: &[C],
strict: bool,
) -> (Vec<RuleReport>, usize, usize) {
let mut checks = Vec::new();
let mut unparseable = 0usize;
let mut parsed = 0usize;
for blob in blobs {
match KernelCheckpoint::from_checkpoint_bytes(blob.as_ref()) {
Ok(checkpoint) => {
parsed += 1;
let records = segments.get(checkpoint.operation_id().as_str());
let replay = strict.then(|| strict_replay(&checkpoint, records.map(Vec::as_slice)));
checks.push(check_c5a(
&checkpoint,
records.map(Vec::as_slice),
replay.as_ref(),
));
checks.push(check_c5b(&checkpoint, replay.as_ref()));
}
Err(error) => {
unparseable += 1;
checks.push(RuleReport {
rule: "C5a".to_string(),
verdict: Verdict::Degraded,
detail: format!(
"checkpoint blob did not decode — its claims are unverifiable (C7): {}",
error.message()
),
});
}
}
}
(checks, unparseable, parsed)
}
enum StrictReplay {
Skipped(String),
Faulted(String),
Done {
at_base: KernelCheckpoint,
ladder: Result<(), String>,
above_records: usize,
},
}
fn strict_replay(checkpoint: &KernelCheckpoint, records: Option<&[KernelRecord]>) -> StrictReplay {
let through = checkpoint.through_step_seq().get();
let base = checkpoint.base_step_seq().get();
let Some(records) = records else {
return StrictReplay::Skipped(
"the journal holds no segment for this operation".to_string(),
);
};
let mut steps: Vec<u64> = records
.iter()
.map(|record| record.step_seq().get())
.filter(|step| *step <= through)
.collect();
steps.sort_unstable();
steps.dedup();
if steps != (0..=through).collect::<Vec<u64>>() {
return StrictReplay::Skipped(format!(
"the journal does not hold an unbroken record run from step 0 through {through} \
(pruned prefix or partial copy); the re-plan replay cannot start at genesis"
));
}
let fold_to_base = || -> Result<KernelCheckpoint, String> {
let mut prefix: Vec<&KernelRecord> = records
.iter()
.filter(|record| record.step_seq().get() <= base)
.collect();
prefix.sort_by_key(|record| record.step_seq().get());
let prefix: Vec<KernelRecord> = prefix.into_iter().cloned().collect();
let folded = restore_operation(
None,
&prefix,
ConfigDefaults::default(),
InMemoryRecordIndex::from_records(&prefix),
)
.map_err(|fault| format!("{}: {}", fault.code.as_str(), fault.message))?;
folded
.transaction
.checkpoint_candidate(folded.driver.project_logical_state())
.map_err(|fault| format!("{}: {}", fault.code.as_str(), fault.message))?
.decode()
.map_err(|error| {
format!(
"the re-derived checkpoint does not decode: {}",
error.message()
)
})
};
let at_base = match fold_to_base() {
Ok(checkpoint) => checkpoint,
Err(fault) => return StrictReplay::Faulted(fault),
};
let mut above: Vec<&KernelRecord> = records
.iter()
.filter(|record| record.step_seq().get() > through)
.collect();
above.sort_by_key(|record| record.step_seq().get());
let above_records = above.len();
let above: Vec<KernelRecord> = above.into_iter().cloned().collect();
let ladder = restore_operation(
Some(checkpoint),
&above,
ConfigDefaults::default(),
InMemoryRecordIndex::from_records(&above),
)
.map(|_: RestoredOperation<InMemoryRecordIndex>| ())
.map_err(|fault| format!("{}: {}", fault.code.as_str(), fault.message));
StrictReplay::Done {
at_base,
ladder,
above_records,
}
}
fn check_c5a(
checkpoint: &KernelCheckpoint,
records: Option<&[KernelRecord]>,
replay: Option<&StrictReplay>,
) -> RuleReport {
let rule = "C5a";
let mut clauses = C5Clauses::new();
let Some(records) = records else {
return RuleReport {
rule: rule.to_string(),
verdict: Verdict::Degraded,
detail: format!(
"checkpoint names operation {}; the journal holds no segment for it, so no \
anchor can be checked",
checkpoint.operation_id()
),
};
};
let through = checkpoint.through_step_seq().get();
match record_at(records, 0) {
Some(genesis)
if genesis.record_digest().as_str() != checkpoint.genesis_digest().as_str() =>
{
clauses.violation(format!(
"the journal's genesis record hashes to {}, but the checkpoint binds genesis {} — \
this checkpoint was captured on another chain",
genesis.record_digest(),
checkpoint.genesis_digest()
));
}
Some(_) => clauses.confirms("genesis digest anchored".to_string()),
None => clauses.degraded(
"the genesis record is not in the journal (pruned prefix); the identity anchor is \
unverifiable"
.to_string(),
),
}
match record_at(records, through) {
Some(record)
if record.record_digest().as_str()
!= checkpoint.covered_transaction_head_digest().as_str() =>
{
clauses.violation(format!(
"the journal record at the covered step {through} hashes to {}, but the \
checkpoint's covered head is {}",
record.record_digest(),
checkpoint.covered_transaction_head_digest()
));
}
Some(_) => clauses.confirms(format!("covered head anchored at step {through}")),
None => clauses.degraded(format!(
"no journal record at the covered step {through}; the covered head is unverifiable"
)),
}
let base = checkpoint.base_step_seq().get();
if base != through {
match record_at(records, base) {
Some(record)
if record.record_digest().as_str() != checkpoint.base_record_digest().as_str() =>
{
clauses.violation(format!(
"the journal record at the tail base step {base} hashes to {}, but the \
checkpoint anchors its tail on {}",
record.record_digest(),
checkpoint.base_record_digest()
));
}
_ => {}
}
}
let mut reconciled = 0usize;
let mut pruned = 0usize;
for entry in checkpoint.tail_inputs() {
match record_at(records, entry.step_seq.get()) {
Some(record) if record.record_digest().as_str() != entry.record_digest.as_str() => {
clauses.violation(format!(
"the journal record at step {} disagrees with the checkpoint's bounded tail \
(journal {}, checkpoint {})",
entry.step_seq.get(),
record.record_digest(),
entry.record_digest
));
}
Some(_) => reconciled += 1,
None => pruned += 1,
}
}
if reconciled > 0 {
clauses.confirms(format!(
"{reconciled} bounded-tail entries reconcile with the journal"
));
}
if pruned > 0 {
clauses.degraded(format!(
"{pruned} bounded-tail entries have no journal record (pruned interval)"
));
}
match replay {
Some(StrictReplay::Skipped(reason)) => {
clauses.degraded(format!("strict replay skipped: {reason}"));
}
Some(StrictReplay::Faulted(fault)) => {
clauses.violation(format!("strict replay faulted: {fault}"));
}
Some(StrictReplay::Done {
at_base,
ladder,
above_records,
}) => {
let base = checkpoint.base_step_seq().get();
if at_base.state_digest() != checkpoint.state_digest() {
clauses.violation(format!(
"strict replay folds the journal to state digest {} at the checkpoint's \
base step {base}, but the checkpoint captured {} there",
at_base.state_digest(),
checkpoint.state_digest()
));
} else {
clauses.confirms(format!(
"strict replay reproduces the captured state digest at step {base}"
));
}
match ladder {
Ok(()) => clauses.confirms(format!(
"the checkpoint+tail restore ladder holds against the {above_records} \
journal record(s) above the covered step"
)),
Err(fault) => clauses.violation(format!(
"the checkpoint+tail restore ladder faults against this journal: {fault}"
)),
}
}
None => {}
}
clauses.report(rule)
}
fn check_c5b(checkpoint: &KernelCheckpoint, replay: Option<&StrictReplay>) -> RuleReport {
let rule = "C5b";
let mut clauses = C5Clauses::new();
let transition = &checkpoint.logical_state().transition;
let through = checkpoint.through_step_seq().get();
let mut mints: HashMap<&str, u64> = HashMap::new();
for entry in &transition.launch_tokens {
let token = entry.launch_token.as_str();
let step = entry.step_seq.get();
match mints.get(token) {
Some(previous) if *previous != step => clauses.violation(format!(
"launch token {token} is minted at step {step} and step {previous} — reuse \
across TaskLaunch payloads"
)),
Some(_) => clauses.violation(format!(
"launch token {token} is registered twice at step {step} — a duplicated ledger \
entry"
)),
None => {
mints.insert(token, step);
}
}
if step > through {
clauses.violation(format!(
"launch token {token} is minted at step {step}, beyond the covered boundary \
{through}"
));
}
}
for effect in &transition.pending_effects {
let EffectKind::SpawnTasks(spawn) = &effect.effect else {
continue;
};
let SpawnTasksEffect { tasks, .. } = spawn;
let effect_step = parse_effect_step(effect.effect_id.as_str()).map(|(_, step)| step);
for launch in tasks {
let token = launch.launch_token.as_str();
match (mints.get(token), effect_step) {
(None, _) => clauses.violation(format!(
"pending effect {} carries launch token {token} the ledger never registered",
effect.effect_id
)),
(Some(&minted), Some(step)) if minted != step => clauses.violation(format!(
"pending effect {} carries launch token {token} minted at step {minted}, not \
at the effect's own step {step}",
effect.effect_id
)),
_ => {}
}
}
}
if !transition.launch_tokens.is_empty() {
clauses.confirms(format!(
"{} launch token(s) anchored; no reuse across TaskLaunch payloads",
transition.launch_tokens.len()
));
}
match replay {
Some(StrictReplay::Skipped(reason)) => {
clauses.degraded(format!("strict replay skipped: {reason}"));
}
Some(StrictReplay::Faulted(fault)) => {
clauses.violation(format!("strict replay faulted: {fault}"));
}
Some(StrictReplay::Done { at_base, .. }) => {
let replayed_ledger = &at_base.logical_state().transition.launch_tokens;
if replayed_ledger != &transition.launch_tokens {
clauses.violation(format!(
"strict replay mints a different launch-token ledger: the journal fold \
registers {} token(s), the checkpoint carries {}",
replayed_ledger.len(),
transition.launch_tokens.len()
));
} else {
clauses.confirms(
"strict replay reproduces the launch-token ledger exactly".to_string(),
);
}
}
None => {}
}
clauses.report(rule)
}
fn record_at(records: &[KernelRecord], step_seq: u64) -> Option<&KernelRecord> {
records
.iter()
.find(|record| record.step_seq().get() == step_seq)
}
#[cfg(test)]
mod tests {
use serde_json::json;
use super::*;
use crate::runtime::kernel::wire::config::{
ConfigDefaults, ExecutionPolicy, HostEffectSupport, OperationConfig,
};
use crate::runtime::kernel::wire::driver::CanonicalOperationDriver;
use crate::runtime::kernel::wire::effect::{
EffectKindTag,
EffectSucceeded,
ProviderCompleted,
ProviderContextOverflow,
ProviderMessage,
ProviderOutcome,
ProviderSuccess,
TaskLaunchOutcome,
TaskLaunchStarted,
TaskLaunchStatus,
TasksSpawnedSuccess,
ToolCall as WireToolCall,
};
use crate::runtime::kernel::wire::envelope::{
ConfigureOperation, KernelInput, ResolveEffect, StartOperation, WireEnvelope,
};
use crate::runtime::kernel::wire::record::{KernelRecord, NormalizedInput};
use crate::runtime::kernel::wire::root::{
InitialContext, LogicalAgentSpec, LogicalMessage, LogicalTask, MessageRole, RootAgentEntry,
RootEntry, RootWorkflowEntry, WorkflowNode as WireWorkflowNode,
WorkflowSpec as WireWorkflowSpec,
};
use crate::runtime::kernel::wire::scalar::{
AttemptId, BoundedJson, CallId, EffectId, InputId, NodeId, OperationId, TaskId, WireU64,
};
use crate::runtime::kernel::wire::transaction::{InMemoryRecordIndex, KernelTransaction};
fn operation(id: &str) -> OperationId {
OperationId::new(id).unwrap()
}
fn envelope(op: &OperationId, id: &str, at: u64, input: KernelInput) -> WireEnvelope {
WireEnvelope::new(
op.clone(),
InputId::new(id).unwrap(),
WireU64::new(at),
input,
)
}
fn configure_envelope(op: &OperationId) -> WireEnvelope {
envelope(
op,
"in-configure",
1_700_000_000_000,
KernelInput::ConfigureOperation(ConfigureOperation {
config: OperationConfig {
execution_policy: Some(ExecutionPolicy {
max_turns: Some(12),
..ExecutionPolicy::default()
}),
host_effect_support: HostEffectSupport::new([
EffectKindTag::CallProvider,
EffectKindTag::SpawnTasks,
]),
..OperationConfig::default()
},
}),
)
}
fn agent_start_envelope(op: &OperationId) -> WireEnvelope {
envelope(
op,
"in-start",
1_700_000_001_000,
KernelInput::StartOperation(StartOperation {
entry: RootEntry::Agent(RootAgentEntry {
task: LogicalTask::new("write the brief"),
run_spec: Some(LogicalAgentSpec::new("write the brief")),
}),
initial_context: InitialContext::default(),
}),
)
}
fn agent_start_with_history_envelope(op: &OperationId, messages: usize) -> WireEnvelope {
envelope(
op,
"in-start",
1_700_000_001_000,
KernelInput::StartOperation(StartOperation {
entry: RootEntry::Agent(RootAgentEntry {
task: LogicalTask::new("write the brief"),
run_spec: Some(LogicalAgentSpec::new("write the brief")),
}),
initial_context: InitialContext {
messages: (0..messages)
.map(|index| LogicalMessage {
role: if index % 2 == 0 {
MessageRole::User
} else {
MessageRole::Assistant
},
content: format!(
"turn {index}: a long enough body that compaction has \
something to reclaim when the prompt stops fitting"
),
tokens: Some(64),
tool_call_id: None,
})
.collect(),
..InitialContext::default()
},
}),
)
}
fn workflow_start_envelope(op: &OperationId) -> WireEnvelope {
envelope(
op,
"in-start",
1_700_000_001_000,
KernelInput::StartOperation(StartOperation {
entry: RootEntry::Workflow(RootWorkflowEntry {
spec: WireWorkflowSpec {
name: "brief".to_string(),
nodes: vec![
WireWorkflowNode {
node_id: NodeId::new("collect").unwrap(),
task: LogicalTask::new("collect the sources"),
depends_on: vec![],
run_spec: Some(LogicalAgentSpec::new("collect the sources")),
},
WireWorkflowNode {
node_id: NodeId::new("write").unwrap(),
task: LogicalTask::new("write the brief"),
depends_on: vec![NodeId::new("collect").unwrap()],
run_spec: Some(LogicalAgentSpec::new("write the brief")),
},
],
},
}),
initial_context: InitialContext::default(),
}),
)
}
fn resolve_overflow_envelope(op: &OperationId, effect_step: u64) -> WireEnvelope {
envelope(
op,
"in-resolve",
1_700_000_002_000,
KernelInput::ResolveEffect(ResolveEffect {
effect_id: EffectId::new(format!("{op}:step:{effect_step}:effect:0")).unwrap(),
outcome: EffectOutcome::Succeeded(EffectSucceeded {
result: EffectSuccess::Provider(ProviderSuccess {
outcome: ProviderOutcome::ContextOverflow(
ProviderContextOverflow::default(),
),
}),
}),
}),
)
}
fn resolve_completed_envelope(
op: &OperationId,
id: &str,
at: u64,
effect_step: u64,
with_tool_call: bool,
) -> WireEnvelope {
envelope(
op,
id,
at,
KernelInput::ResolveEffect(ResolveEffect {
effect_id: EffectId::new(format!("{op}:step:{effect_step}:effect:0")).unwrap(),
outcome: EffectOutcome::Succeeded(EffectSucceeded {
result: EffectSuccess::Provider(ProviderSuccess {
outcome: ProviderOutcome::Completed(ProviderCompleted {
message: ProviderMessage {
role: MessageRole::Assistant,
content: "done".to_string(),
tool_calls: if with_tool_call {
vec![WireToolCall {
call_id: CallId::new("call-1").unwrap(),
name: "read_file".to_string(),
arguments: BoundedJson::new(json!({})).unwrap(),
}]
} else {
Vec::new()
},
tool_call_id: None,
tokens: None,
},
observed_input_tokens: None,
observed_output_tokens: None,
stop_reason: None,
}),
}),
}),
}),
)
}
fn resolve_spawn_envelope(
op: &OperationId,
id: &str,
at: u64,
effect_id: &str,
tasks: &[(&str, &str)],
) -> WireEnvelope {
envelope(
op,
id,
at,
KernelInput::ResolveEffect(ResolveEffect {
effect_id: EffectId::new(effect_id).unwrap(),
outcome: EffectOutcome::Succeeded(EffectSucceeded {
result: EffectSuccess::TasksSpawned(TasksSpawnedSuccess {
attempts: tasks
.iter()
.map(|(task, attempt)| TaskLaunchOutcome {
task_id: TaskId::new(*task).unwrap(),
attempt_id: AttemptId::new(*attempt).unwrap(),
outcome: TaskLaunchStatus::Started(TaskLaunchStarted {}),
})
.collect(),
}),
}),
}),
)
}
fn live_chain(envelopes: &[WireEnvelope]) -> Vec<KernelRecord> {
let mut tx = KernelTransaction::new(ConfigDefaults::default(), InMemoryRecordIndex::new());
let mut driver = CanonicalOperationDriver::new();
let mut journal = Vec::new();
for envelope in envelopes {
let preparation = tx.prepare(envelope, |context| driver.plan(context));
let token = preparation
.token()
.unwrap_or_else(|| {
panic!("expected a prepared step, got {:?}", preparation.fault())
})
.clone();
let head = preparation.record().unwrap().record_digest().clone();
let committed = tx.commit(&token, &head).expect("commit must succeed");
journal.push(committed.record.clone());
driver
.note_committed(committed.step_seq)
.expect("the driver folds the step it planned");
}
journal
}
fn hand_chain(envelopes: &[WireEnvelope]) -> Vec<KernelRecord> {
let mut records: Vec<KernelRecord> = Vec::new();
for (index, envelope) in envelopes.iter().enumerate() {
let input = NormalizedInput::normalize(envelope, &ConfigDefaults::default())
.expect("the envelope normalises");
let step = json!({ "planned": format!("step-{index}"), "effects": [] });
let record =
KernelRecord::chain(records.last(), &input, &step).expect("the record chains");
records.push(record);
}
records
}
fn blobs(records: &[KernelRecord]) -> Vec<Vec<u8>> {
records
.iter()
.map(|record| record.record_bytes().into_vec())
.collect()
}
fn rule<'a>(report: &'a ValidationReport, segment: usize, id: &str) -> &'a RuleReport {
report.segments[segment]
.rules
.iter()
.find(|rule| rule.rule == id)
.unwrap_or_else(|| panic!("segment {segment} has no {id} verdict"))
}
fn cross<'a>(report: &'a ValidationReport, id: &str) -> &'a RuleReport {
report
.cross_checks
.iter()
.find(|rule| rule.rule == id)
.unwrap_or_else(|| panic!("the report has no {id} cross-check"))
}
#[test]
fn a_green_agent_chain_passes_every_rule() {
let op = operation("op-green-agent");
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let report = validate_journal(&blobs(&chain));
assert_eq!(report.segments.len(), 1);
for id in ["C1", "C2", "C3", "C4"] {
assert_eq!(
rule(&report, 0, id).verdict,
Verdict::Pass,
"{id}: {}",
rule(&report, 0, id).detail
);
}
assert!(
rule(&report, 0, "C3")
.detail
.contains("every durable record digest reproduced"),
"C3 proves the re-plan: {}",
rule(&report, 0, "C3").detail
);
assert_eq!(report.exit_code(), 0);
assert_eq!(report.unparseable_records, 0);
}
#[test]
fn a_green_workflow_chain_passes_c4_with_real_launches() {
let op = operation("op-green-workflow");
let chain = live_chain(&[
configure_envelope(&op),
workflow_start_envelope(&op),
resolve_spawn_envelope(
&op,
"in-ack-1",
1_700_000_002_000,
"op-green-workflow:step:1:effect:0",
&[("wf-node0", "wf-node0:attempt:1")],
),
]);
let report = validate_journal(&blobs(&chain));
assert_eq!(report.segments.len(), 1);
for id in ["C1", "C2", "C3", "C4"] {
assert_eq!(
rule(&report, 0, id).verdict,
Verdict::Pass,
"{id}: {}",
rule(&report, 0, id).detail
);
}
assert!(
rule(&report, 0, "C4").detail.contains("1 launch(es)"),
"{}",
rule(&report, 0, "C4").detail
);
assert_eq!(report.deferred.len(), 2, "batch-1 scope limits are named");
assert_eq!(report.exit_code(), 0);
}
#[test]
fn input_order_is_a_storage_detail() {
let op = operation("op-shuffled");
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let mut shuffled = blobs(&chain);
shuffled.reverse();
let report = validate_journal(&shuffled);
assert_eq!(
report.exit_code(),
0,
"the chain's own links define the order"
);
}
#[test]
fn two_operations_validate_as_independent_segments() {
let op_a = operation("op-seg-a");
let op_b = operation("op-seg-b");
let chain_a = live_chain(&[configure_envelope(&op_a), agent_start_envelope(&op_a)]);
let chain_b = live_chain(&[configure_envelope(&op_b), agent_start_envelope(&op_b)]);
let mut mixed = Vec::new();
for index in 0..2 {
mixed.push(chain_a[index].record_bytes().into_vec());
mixed.push(chain_b[index].record_bytes().into_vec());
}
let report = validate_journal(&mixed);
assert_eq!(report.segments.len(), 2);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn a_gap_in_the_chain_fails_c1_and_degrades_c3() {
let op = operation("op-gapped");
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let gapped = blobs(&[chain[0].clone(), chain[2].clone()]);
let report = validate_journal(&gapped);
assert_eq!(rule(&report, 0, "C1").verdict, Verdict::Fail);
assert_eq!(
rule(&report, 0, "C3").verdict,
Verdict::Degraded,
"a re-plan over a broken chain would only re-report the C1 break"
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn two_different_records_for_one_input_fail_c2() {
let op = operation("op-dup-input");
let chain_a = hand_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let mut later_start = agent_start_envelope(&op);
later_start.observed_at_ms = WireU64::new(1_700_000_001_500);
let chain_b = hand_chain(&[configure_envelope(&op), later_start]);
assert_ne!(
chain_a[1].record_digest(),
chain_b[1].record_digest(),
"the fixture must produce two different records for one input id"
);
let report = validate_journal(&blobs(&[
chain_a[0].clone(),
chain_a[1].clone(),
chain_b[1].clone(),
]));
assert_eq!(rule(&report, 0, "C2").verdict, Verdict::Fail);
assert!(
rule(&report, 0, "C2").detail.contains("in-start"),
"{}",
rule(&report, 0, "C2").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_repeated_attempt_pair_is_a_reused_launch_token() {
let op = operation("op-dup-launch");
let chain = hand_chain(&[
configure_envelope(&op),
resolve_spawn_envelope(
&op,
"in-ack-1",
1_700_000_001_000,
"op-dup-launch:step:0:effect:0",
&[("writer", "writer:attempt:1")],
),
resolve_spawn_envelope(
&op,
"in-ack-2",
1_700_000_002_000,
"op-dup-launch:step:0:effect:0",
&[("writer", "writer:attempt:1")],
),
]);
let report = validate_journal(&blobs(&chain));
assert_eq!(rule(&report, 0, "C1").verdict, Verdict::Pass);
assert_eq!(rule(&report, 0, "C4").verdict, Verdict::Fail);
assert!(
rule(&report, 0, "C4").detail.contains("LaunchToken"),
"the verdict names the token reuse: {}",
rule(&report, 0, "C4").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_resolution_naming_a_future_step_fails_c4() {
let op = operation("op-future-effect");
let chain = hand_chain(&[
configure_envelope(&op),
resolve_spawn_envelope(
&op,
"in-ack-1",
1_700_000_001_000,
"op-future-effect:step:5:effect:0",
&[("writer", "writer:attempt:1")],
),
]);
let report = validate_journal(&blobs(&chain));
assert_eq!(rule(&report, 0, "C4").verdict, Verdict::Fail);
assert!(
rule(&report, 0, "C4")
.detail
.contains("precedes the publication"),
"{}",
rule(&report, 0, "C4").detail
);
}
#[test]
fn a_resolution_naming_another_operation_fails_c4() {
let op = operation("op-foreign-effect");
let chain = hand_chain(&[
configure_envelope(&op),
resolve_spawn_envelope(
&op,
"in-ack-1",
1_700_000_001_000,
"op-somewhere-else:step:0:effect:0",
&[("writer", "writer:attempt:1")],
),
]);
let report = validate_journal(&blobs(&chain));
assert_eq!(rule(&report, 0, "C4").verdict, Verdict::Fail);
assert!(
rule(&report, 0, "C4").detail.contains("another operation"),
"{}",
rule(&report, 0, "C4").detail
);
}
#[test]
fn a_tampered_hop_fails_integrity_validation() {
let op = operation("op-tampered");
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let mut input = blobs(&chain);
let mut forged: serde_json::Value = serde_json::from_slice(&input[1]).unwrap();
forged["step_digest"] = serde_json::Value::String(chain[0].record_digest().to_string());
input[1] = serde_json::to_vec(&forged).unwrap();
let report = validate_journal(&input);
assert_eq!(report.segments.len(), 1);
assert_eq!(report.segments[0].degraded_hops.len(), 1);
assert_eq!(
rule(&report, 0, "C1").verdict,
Verdict::Fail,
"a digest mismatch must fail C1: {}",
rule(&report, 0, "C1").detail
);
assert_eq!(rule(&report, 0, "C3").verdict, Verdict::Degraded);
assert_eq!(rule(&report, 0, "C4").verdict, Verdict::Degraded);
assert_eq!(
report.exit_code(),
1,
"proven digest corruption must fail the validator"
);
}
#[test]
fn missing_legacy_digest_degrades_without_claiming_corruption() {
let op = operation("op-legacy");
let chain = live_chain(&[configure_envelope(&op)]);
let mut legacy: serde_json::Value = serde_json::from_slice(&blobs(&chain)[0]).unwrap();
legacy.as_object_mut().unwrap().remove("step_digest");
let report = validate_journal(&[serde_json::to_vec(&legacy).unwrap()]);
assert_eq!(rule(&report, 0, "C1").verdict, Verdict::Degraded);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn unparseable_input_is_evidence_insufficient_not_guilty() {
let report = validate_journal(&[b"this is not a record".to_vec()]);
assert!(report.segments.is_empty());
assert_eq!(report.unparseable_records, 1);
assert_eq!(report.exit_code(), 2);
}
#[test]
fn garbage_beside_a_green_chain_stays_exit_2_without_a_violation() {
let op = operation("op-plus-garbage");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let mut input = blobs(&chain);
input.push(b"this is not a record".to_vec());
let report = validate_journal(&input);
assert_eq!(report.segments.len(), 1);
assert_eq!(rule(&report, 0, "C1").verdict, Verdict::Pass);
assert_eq!(report.unparseable_records, 1);
assert_eq!(
report.exit_code(),
2,
"no violation was proven, but the evidence was partially unreadable"
);
}
#[test]
fn an_empty_journal_is_evidence_insufficient() {
let report = validate_journal::<Vec<u8>>(&[]);
assert_eq!(report.exit_code(), 2);
}
fn session_event(value: serde_json::Value) -> Vec<u8> {
serde_json::to_vec(&value).unwrap()
}
#[test]
fn session_events_classify_leniently_across_host_spellings() {
let node_attempt = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op:step:1:effect:0",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a", "provider": "p" },
"status": "success"
}));
let py_attempt = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op:step:2:effect:0",
"route": { "route_id": "route-b" }
}));
let run_started = session_event(json!({
"kind": "run_started",
"run_id": "run-1",
"route": { "routeId": "route-a" }
}));
let node_measured = session_event(json!({
"kind": "prompt_measured",
"turn": 1,
"effect_id": "op:step:1:effect:0",
"measurement": { "requestFingerprint": "fp-1", "inputTokens": 10 }
}));
let py_measured = session_event(json!({
"kind": "prompt_measured",
"measurement": { "request_fingerprint": "fp-2" }
}));
let llm_completed = session_event(json!({
"kind": "llm_completed",
"effect_id": "op:step:2:effect:0",
"invocation_id": "op:step:1:effect:0"
}));
let unknown_kind = session_event(json!({ "kind": "compressed", "turn": 3 }));
let kindless = session_event(json!({ "turn": 3 }));
assert_eq!(
classify_session_event(&node_attempt),
Some(EvidenceEvent::ProviderAttempt {
effect_id: Some("op:step:1:effect:0".to_string()),
request_fingerprint: Some("fp-1".to_string()),
route_id: Some("route-a".to_string()),
status: Some("success".to_string()),
})
);
assert_eq!(
classify_session_event(&py_attempt),
Some(EvidenceEvent::ProviderAttempt {
effect_id: Some("op:step:2:effect:0".to_string()),
request_fingerprint: None,
route_id: Some("route-b".to_string()),
status: None,
})
);
assert_eq!(
classify_session_event(&run_started),
Some(EvidenceEvent::RunStarted {
route_id: Some("route-a".to_string())
})
);
assert_eq!(
classify_session_event(&node_measured),
Some(EvidenceEvent::PromptMeasured {
effect_id: Some("op:step:1:effect:0".to_string()),
request_fingerprint: Some("fp-1".to_string()),
})
);
assert_eq!(
classify_session_event(&py_measured),
Some(EvidenceEvent::PromptMeasured {
effect_id: None,
request_fingerprint: Some("fp-2".to_string()),
})
);
assert_eq!(
classify_session_event(&llm_completed),
Some(EvidenceEvent::LlmCompleted {
effect_id: Some("op:step:2:effect:0".to_string()),
invocation_id: Some("op:step:1:effect:0".to_string()),
})
);
assert_eq!(
classify_session_event(&unknown_kind),
Some(EvidenceEvent::Other),
"unknown kinds are parseable but ignored — the vocabulary evolves"
);
assert_eq!(
classify_session_event(&kindless),
Some(EvidenceEvent::Other)
);
assert_eq!(
classify_session_event(b"not json"),
None,
"a non-object event blob is unparseable input, never a violation"
);
}
#[test]
fn dual_input_with_a_green_journal_and_real_events_stays_green() {
let op = operation("op-dual-green");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let stream = vec![
session_event(json!({
"kind": "run_started",
"run_id": "r1",
"route": { "routeId": "route-a" }
})),
session_event(json!({
"kind": "prompt_measured",
"turn": 1,
"effect_id": "op-dual-green:step:1:effect:0",
"measurement": { "requestFingerprint": "fp-1", "inputTokens": 10 }
})),
session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-dual-green:step:1:effect:0",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a" },
"status": "success"
})),
session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-dual-green:step:1:effect:0",
"invocation_id": "op-dual-green:step:1:effect:0"
})),
];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(report.session_events, Some(4));
assert_eq!(report.unparseable_events, 0);
for id in ["C6.1", "C6.2", "C6.3", "C8"] {
assert_eq!(
cross(&report, id).verdict,
Verdict::Pass,
"{id}: {}",
cross(&report, id).detail
);
}
assert!(
!report
.deferred
.iter()
.any(|line| line.starts_with("c6.") || line.starts_with("c8.")),
"C6/C8 are implemented — the interim scope notes are gone"
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn journal_only_validation_carries_no_session_plane() {
let op = operation("op-journal-only");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let report = validate_journal(&blobs(&chain));
assert_eq!(report.session_events, None);
assert_eq!(report.unparseable_events, 0);
assert_eq!(
report.deferred.len(),
2,
"batch-1 deferred scope is unchanged"
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn an_empty_session_plane_is_evidence_insufficient() {
let op = operation("op-empty-session");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let report = validate_with_session_log(&blobs(&chain), &[Vec::<Vec<u8>>::new()]);
assert_eq!(report.session_events, Some(0));
assert!(
!report.has_violations(),
"an empty log proves nothing either way"
);
assert_eq!(report.exit_code(), 2);
}
#[test]
fn garbage_session_events_are_evidence_insufficient_not_guilty() {
let op = operation("op-garbage-session");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let stream = vec![
session_event(json!({ "kind": "run_started", "run_id": "r1" })),
b"this is not an event".to_vec(),
];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(report.session_events, Some(1));
assert_eq!(report.unparseable_events, 1);
assert!(!report.has_violations());
assert_eq!(report.exit_code(), 2);
}
fn honest_dual_input(op_name: &str) -> (Vec<KernelRecord>, Vec<Vec<u8>>) {
let op = operation(op_name);
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let effect = format!("{op_name}:step:1:effect:0");
let stream = vec![
session_event(json!({
"kind": "run_started",
"run_id": "r1",
"route": { "routeId": "route-a" }
})),
session_event(json!({
"kind": "prompt_measured",
"turn": 1,
"effect_id": effect,
"measurement": { "requestFingerprint": "fp-1", "inputTokens": 10 }
})),
session_event(json!({
"kind": "provider_attempt",
"effect_id": effect,
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a" },
"status": "success"
})),
];
(chain, stream)
}
#[test]
fn an_attempt_naming_an_effect_the_replan_never_published_fails_c6() {
let (chain, mut stream) = honest_dual_input("op-forged-effect");
stream[2] = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-forged-effect:step:1:effect:7",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a" },
"status": "success"
}));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C6.1").verdict, Verdict::Fail);
assert!(
cross(&report, "C6.1").detail.contains("never published"),
"{}",
cross(&report, "C6.1").detail
);
assert_eq!(
report.exit_code(),
1,
"the forged attempt turns the run red"
);
}
#[test]
fn an_attempt_without_an_effect_id_fails_c6_as_forged_evidence() {
let (chain, mut stream) = honest_dual_input("op-keyless-attempt");
stream[2] = session_event(json!({
"kind": "provider_attempt",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a" },
"status": "success"
}));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C6.1").verdict, Verdict::Fail);
assert!(
cross(&report, "C6.1").detail.contains("without effect_id"),
"{}",
cross(&report, "C6.1").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn an_attempt_past_the_journal_tip_degrades_c6_instead_of_failing() {
let (chain, mut stream) = honest_dual_input("op-prefix-attempt");
stream[2] = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-prefix-attempt:step:9:effect:0",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a" },
"status": "success"
}));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(
cross(&report, "C6.1").verdict,
Verdict::Degraded,
"a journal prefix cannot disprove an effect past its tip: {}",
cross(&report, "C6.1").detail
);
assert_eq!(report.exit_code(), 0, "degradation never turns the run red");
}
#[test]
fn an_attempt_on_an_operation_without_a_segment_degrades_c6() {
let (chain, mut stream) = honest_dual_input("op-subset-journal");
stream[2] = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-elsewhere:step:1:effect:0",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-a" },
"status": "success"
}));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C6.1").verdict, Verdict::Degraded);
assert!(
cross(&report, "C6.1").detail.contains("no journal segment"),
"{}",
cross(&report, "C6.1").detail
);
}
#[test]
fn an_orphan_fingerprint_fails_c6() {
let (chain, mut stream) = honest_dual_input("op-orphan-fp");
stream.remove(1); let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C6.2").verdict, Verdict::Fail);
assert!(
cross(&report, "C6.2").detail.contains("fp-1"),
"{}",
cross(&report, "C6.2").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn an_attempt_without_a_fingerprint_fails_c6() {
let (chain, mut stream) = honest_dual_input("op-fpless-attempt");
stream[2] = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-fpless-attempt:step:1:effect:0",
"route": { "routeId": "route-a" },
"status": "success"
}));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C6.2").verdict, Verdict::Fail);
assert!(
cross(&report, "C6.2")
.detail
.contains("without request_fingerprint"),
"{}",
cross(&report, "C6.2").detail
);
}
#[test]
fn an_in_run_route_change_fails_c6() {
let (chain, mut stream) = honest_dual_input("op-route-flip");
stream[2] = session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-route-flip:step:1:effect:0",
"request_fingerprint": "fp-1",
"route": { "routeId": "route-b" },
"status": "success"
}));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C6.3").verdict, Verdict::Fail);
assert!(
cross(&report, "C6.3")
.detail
.contains("in-run route change"),
"{}",
cross(&report, "C6.3").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_cross_resume_route_change_degrades_c6_per_q3() {
let (chain, mut stream) = honest_dual_input("op-route-resume");
stream.push(session_event(json!({
"kind": "run_started",
"run_id": "r1",
"route": { "routeId": "route-b" }
})));
stream.push(session_event(json!({
"kind": "prompt_measured",
"turn": 2,
"effect_id": "op-route-resume:step:1:effect:0",
"measurement": { "requestFingerprint": "fp-2", "inputTokens": 11 }
})));
stream.push(session_event(json!({
"kind": "provider_attempt",
"effect_id": "op-route-resume:step:1:effect:0",
"request_fingerprint": "fp-2",
"route": { "routeId": "route-b" },
"status": "success"
})));
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(
cross(&report, "C6.3").verdict,
Verdict::Degraded,
"cross-resume route changes mark, they do not fail: {}",
cross(&report, "C6.3").detail
);
assert!(
cross(&report, "C6.3").detail.contains("cross-resume"),
"{}",
cross(&report, "C6.3").detail
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn a_pre_0_2_63_log_without_attempts_degrades_every_c6_clause() {
let op = operation("op-old-log");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let stream = vec![
session_event(json!({ "kind": "run_started", "run_id": "r1" })),
session_event(json!({ "kind": "llm_completed", "turn": 1, "content": "done" })),
];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
for id in ["C6.1", "C6.2", "C6.3", "C8"] {
assert_eq!(
cross(&report, id).verdict,
Verdict::Degraded,
"{id}: old logs degrade (C7), never fail — {}",
cross(&report, id).detail
);
}
assert_eq!(report.exit_code(), 0);
}
fn overflow_retry_chain(op_name: &str) -> Vec<KernelRecord> {
let op = operation(op_name);
live_chain(&[
configure_envelope(&op),
agent_start_with_history_envelope(&op, 14),
resolve_overflow_envelope(&op, 1),
resolve_completed_envelope(&op, "in-resolve-2", 1_700_000_003_000, 2, false),
])
}
#[test]
fn an_honest_overflow_retry_invocation_passes_c8() {
let chain = overflow_retry_chain("op-c8-green");
let stream = vec![
session_event(json!({ "kind": "run_started", "run_id": "r1" })),
session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-c8-green:step:2:effect:0",
"invocation_id": "op-c8-green:step:1:effect:0"
})),
];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(
cross(&report, "C8").verdict,
Verdict::Pass,
"{}",
cross(&report, "C8").detail
);
assert!(
cross(&report, "C8")
.detail
.contains("1 retried invocation(s)"),
"{}",
cross(&report, "C8").detail
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn a_completed_chain_head_is_the_merge_forgery() {
let op = operation("op-c8-merged");
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_completed_envelope(&op, "in-resolve-1", 1_700_000_002_000, 1, true),
]);
let stream = vec![session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-c8-merged:step:2:effect:0",
"invocation_id": "op-c8-merged:step:1:effect:0"
}))];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C8").verdict, Verdict::Fail);
assert!(
cross(&report, "C8")
.detail
.contains("closes its invocation"),
"{}",
cross(&report, "C8").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_selected_effect_preceding_the_chain_head_fails_c8() {
let chain = overflow_retry_chain("op-c8-backwards");
let stream = vec![session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-c8-backwards:step:1:effect:0",
"invocation_id": "op-c8-backwards:step:2:effect:0"
}))];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C8").verdict, Verdict::Fail);
assert!(
cross(&report, "C8")
.detail
.contains("does not follow the chain head"),
"{}",
cross(&report, "C8").detail
);
}
#[test]
fn a_selected_effect_the_replan_never_published_fails_c8() {
let chain = overflow_retry_chain("op-c8-phantom");
let stream = vec![session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-c8-phantom:step:2:effect:9",
"invocation_id": "op-c8-phantom:step:1:effect:0"
}))];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(cross(&report, "C8").verdict, Verdict::Fail);
assert!(
cross(&report, "C8").detail.contains("never published"),
"{}",
cross(&report, "C8").detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_first_try_invocation_has_no_adjacency_to_prove() {
let op = operation("op-c8-first-try");
let chain = live_chain(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_completed_envelope(&op, "in-resolve-1", 1_700_000_002_000, 1, false),
]);
let stream = vec![session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-c8-first-try:step:1:effect:0",
"invocation_id": "op-c8-first-try:step:1:effect:0"
}))];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(
cross(&report, "C8").verdict,
Verdict::Pass,
"{}",
cross(&report, "C8").detail
);
assert!(
cross(&report, "C8").detail.contains("1 first-try"),
"{}",
cross(&report, "C8").detail
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn an_invocation_past_the_journal_tip_degrades_c8() {
let op = operation("op-c8-lag");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let stream = vec![session_event(json!({
"kind": "llm_completed",
"turn": 1,
"effect_id": "op-c8-lag:step:2:effect:0",
"invocation_id": "op-c8-lag:step:1:effect:0"
}))];
let report = validate_with_session_log(&blobs(&chain), &[stream]);
assert_eq!(
cross(&report, "C8").verdict,
Verdict::Degraded,
"{}",
cross(&report, "C8").detail
);
assert_eq!(report.exit_code(), 0);
}
use crate::runtime::kernel::wire::checkpoint::{
CheckpointDraft, KernelCheckpoint, LaunchTokenState,
};
use crate::runtime::kernel::wire::driver::PlannedStep;
use crate::runtime::kernel::wire::transaction::CheckpointBoundary;
use std::path::PathBuf;
fn live_runtime(
envelopes: &[WireEnvelope],
) -> (
Vec<KernelRecord>,
KernelTransaction<PlannedStep, InMemoryRecordIndex>,
CanonicalOperationDriver,
) {
let mut tx = KernelTransaction::new(ConfigDefaults::default(), InMemoryRecordIndex::new());
let mut driver = CanonicalOperationDriver::new();
let mut journal = Vec::new();
for envelope in envelopes {
let preparation = tx.prepare(envelope, |context| driver.plan(context));
let token = preparation
.token()
.unwrap_or_else(|| {
panic!("expected a prepared step, got {:?}", preparation.fault())
})
.clone();
let head = preparation.record().unwrap().record_digest().clone();
let committed = tx.commit(&token, &head).expect("commit must succeed");
journal.push(committed.record.clone());
driver
.note_committed(committed.step_seq)
.expect("the driver folds the step it planned");
}
(journal, tx, driver)
}
fn checkpoint_at_head(
tx: &KernelTransaction<PlannedStep, InMemoryRecordIndex>,
driver: &CanonicalOperationDriver,
) -> KernelCheckpoint {
tx.checkpoint_candidate(driver.project_logical_state())
.expect("the head checkpoints")
.decode()
.expect("the candidate decodes")
}
fn checkpoint_blob(checkpoint: &KernelCheckpoint) -> Vec<u8> {
checkpoint.checkpoint_bytes().into_vec()
}
fn checkpoint_check<'a>(report: &'a ValidationReport, id: &str) -> &'a RuleReport {
report
.checkpoint_checks
.iter()
.find(|rule| rule.rule == id)
.unwrap_or_else(|| panic!("the report has no {id} checkpoint-check"))
}
fn no_streams() -> Vec<Vec<Vec<u8>>> {
Vec::new()
}
#[test]
fn without_a_checkpoint_plane_c5_is_deferred_not_red() {
let op = operation("op-c5-deferred");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let report = validate_journal(&blobs(&chain));
assert!(report.checkpoint_checks.is_empty());
assert_eq!(report.checkpoints, None);
assert_eq!(report.unparseable_checkpoints, 0);
assert_eq!(
report.deferred.len(),
2,
"the checkpoint plane was never offered"
);
assert!(
report.deferred[1].contains("c5b.launch_token_ledger"),
"{}",
report.deferred[1]
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn a_real_checkpoint_anchors_its_journal() {
let op = operation("op-c5-green");
let (chain, tx, driver) = live_runtime(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let checkpoint = checkpoint_at_head(&tx, &driver);
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&checkpoint)],
false,
);
assert_eq!(report.checkpoints, Some(1));
assert_eq!(report.unparseable_checkpoints, 0);
for id in ["C5a", "C5b"] {
assert_eq!(
checkpoint_check(&report, id).verdict,
Verdict::Pass,
"{id}: {}",
checkpoint_check(&report, id).detail
);
}
assert!(
checkpoint_check(&report, "C5a")
.detail
.contains("covered head anchored at step"),
"{}",
checkpoint_check(&report, "C5a").detail
);
assert_eq!(
report.deferred.len(),
1,
"the checkpoint plane retires the c5b deferral"
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn strict_replay_reproduces_the_covered_state_and_the_ladder_holds() {
let op = operation("op-c5-strict");
let (chain, tx, driver) = live_runtime(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let checkpoint = checkpoint_at_head(&tx, &driver);
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&checkpoint)],
true,
);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Pass, "{}", c5a.detail);
assert!(
c5a.detail
.contains("strict replay reproduces the captured state digest"),
"{}",
c5a.detail
);
assert!(
c5a.detail.contains("restore ladder holds"),
"{}",
c5a.detail
);
let c5b = checkpoint_check(&report, "C5b");
assert_eq!(c5b.verdict, Verdict::Pass, "{}", c5b.detail);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn strict_replay_with_tail_records_drives_the_whole_ladder() {
let op = operation("op-c5-strict-tail");
let (mut chain, mut tx, mut driver) =
live_runtime(&[configure_envelope(&op), agent_start_envelope(&op)]);
let checkpoint = checkpoint_at_head(&tx, &driver);
let preparation = tx.prepare(&resolve_overflow_envelope(&op, 1), |context| {
driver.plan(context)
});
let token = preparation.token().expect("the tail step prepares").clone();
let head = preparation.record().unwrap().record_digest().clone();
let committed = tx.commit(&token, &head).expect("the tail step commits");
chain.push(committed.record.clone());
driver
.note_committed(committed.step_seq)
.expect("the driver folds the tail step");
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&checkpoint)],
true,
);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Pass, "{}", c5a.detail);
assert!(
c5a.detail.contains("against the 1 journal record(s) above"),
"{}",
c5a.detail
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn a_bounded_tail_checkpoint_reconciles_with_the_journal() {
let op = operation("op-c5-bounded");
let (mut chain, mut tx, mut driver) =
live_runtime(&[configure_envelope(&op), agent_start_envelope(&op)]);
let candidate = tx
.checkpoint_candidate(driver.project_logical_state())
.expect("the base checkpoints");
let boundary: CheckpointBoundary = candidate.boundary();
let base_state = candidate
.decode()
.expect("the base decodes")
.logical_state()
.clone();
let preparation = tx.prepare(&resolve_overflow_envelope(&op, 1), |context| {
driver.plan(context)
});
let token = preparation.token().expect("the tail step prepares").clone();
let head = preparation.record().unwrap().record_digest().clone();
let committed = tx.commit(&token, &head).expect("the tail step commits");
chain.push(committed.record.clone());
driver
.note_committed(committed.step_seq)
.expect("the driver folds the tail step");
let rebased = tx
.checkpoint_rebase(&boundary, base_state)
.expect("the window rebases")
.decode()
.expect("the rebase decodes");
assert_eq!(rebased.base_step_seq().get(), 1);
assert_eq!(rebased.through_step_seq().get(), 2);
assert_eq!(rebased.tail_inputs().len(), 1);
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&rebased)],
true,
);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Pass, "{}", c5a.detail);
assert!(
c5a.detail.contains("1 bounded-tail entries reconcile"),
"{}",
c5a.detail
);
assert!(
c5a.detail.contains("the captured state digest at step 1"),
"{}",
c5a.detail
);
assert_eq!(report.exit_code(), 0);
}
#[test]
fn a_checkpoint_from_another_chain_fails_c5a() {
let op = operation("op-c5-foreign");
let (chain, _tx, _driver) =
live_runtime(&[configure_envelope(&op), agent_start_envelope(&op)]);
let foreign_config = |max_turns: u32| {
envelope(
&op,
"in-configure",
1_700_000_000_000,
KernelInput::ConfigureOperation(ConfigureOperation {
config: OperationConfig {
execution_policy: Some(ExecutionPolicy {
max_turns: Some(max_turns),
..ExecutionPolicy::default()
}),
host_effect_support: HostEffectSupport::new([
EffectKindTag::CallProvider,
EffectKindTag::SpawnTasks,
]),
..OperationConfig::default()
},
}),
)
};
let (_chain_other, other_tx, other_driver) =
live_runtime(&[foreign_config(24), agent_start_envelope(&op)]);
let foreign = checkpoint_at_head(&other_tx, &other_driver);
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&foreign)],
false,
);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Fail, "{}", c5a.detail);
assert!(
c5a.detail.contains("captured on another chain"),
"{}",
c5a.detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_reused_launch_token_fails_c5b() {
let op = operation("op-c5-tokens");
let (chain, tx, driver) = live_runtime(&[
configure_envelope(&op),
workflow_start_envelope(&op),
resolve_spawn_envelope(
&op,
"in-ack-1",
1_700_000_002_000,
"op-c5-tokens:step:1:effect:0",
&[("wf-node0", "wf-node0:attempt:1")],
),
]);
let checkpoint = checkpoint_at_head(&tx, &driver);
assert_eq!(
checkpoint.logical_state().transition.launch_tokens.len(),
1,
"the workflow start minted one launch token"
);
let mut state = checkpoint.logical_state().clone();
let minted = state.transition.launch_tokens[0].clone();
state.transition.launch_tokens.push(LaunchTokenState {
launch_token: minted.launch_token,
step_seq: WireU64::new(2),
});
let forged = KernelCheckpoint::assemble(CheckpointDraft {
operation_id: checkpoint.operation_id().clone(),
genesis_digest: checkpoint.genesis_digest().clone(),
base_step_seq: checkpoint.base_step_seq(),
base_record_digest: checkpoint.base_record_digest().clone(),
through_step_seq: checkpoint.through_step_seq(),
covered_transaction_head_digest: checkpoint.covered_transaction_head_digest().clone(),
logical_state: state,
tail_inputs: checkpoint.tail_inputs().to_vec(),
})
.expect("the forged draft assembles");
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&forged)],
false,
);
let c5b = checkpoint_check(&report, "C5b");
assert_eq!(c5b.verdict, Verdict::Fail, "{}", c5b.detail);
assert!(
c5b.detail.contains("reuse across TaskLaunch payloads"),
"{}",
c5b.detail
);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn strict_replay_catches_a_ledger_the_journal_never_minted() {
let op = operation("op-c5-moved-mint");
let (chain, tx, driver) = live_runtime(&[
configure_envelope(&op),
workflow_start_envelope(&op),
resolve_spawn_envelope(
&op,
"in-ack-1",
1_700_000_002_000,
"op-c5-moved-mint:step:1:effect:0",
&[("wf-node0", "wf-node0:attempt:1")],
),
]);
let checkpoint = checkpoint_at_head(&tx, &driver);
let mut state = checkpoint.logical_state().clone();
state.transition.launch_tokens[0].step_seq = WireU64::new(2);
let forged = KernelCheckpoint::assemble(CheckpointDraft {
operation_id: checkpoint.operation_id().clone(),
genesis_digest: checkpoint.genesis_digest().clone(),
base_step_seq: checkpoint.base_step_seq(),
base_record_digest: checkpoint.base_record_digest().clone(),
through_step_seq: checkpoint.through_step_seq(),
covered_transaction_head_digest: checkpoint.covered_transaction_head_digest().clone(),
logical_state: state,
tail_inputs: checkpoint.tail_inputs().to_vec(),
})
.expect("the forged draft assembles");
let default_report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&forged)],
false,
);
assert_eq!(
checkpoint_check(&default_report, "C5b").verdict,
Verdict::Pass,
"the default plane cannot see a moved mint: {}",
checkpoint_check(&default_report, "C5b").detail
);
let strict_report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[checkpoint_blob(&forged)],
true,
);
let c5b = checkpoint_check(&strict_report, "C5b");
assert_eq!(c5b.verdict, Verdict::Fail, "{}", c5b.detail);
assert!(
c5b.detail.contains("different launch-token ledger"),
"{}",
c5b.detail
);
assert_eq!(
checkpoint_check(&strict_report, "C5a").verdict,
Verdict::Fail,
"the ledger is part of the state, so the state digest moves too"
);
assert_eq!(strict_report.exit_code(), 1);
}
#[test]
fn a_tail_disconnected_from_the_journal_fails_c5a() {
let op = operation("op-c5-spliced");
let (chain_a, tx, driver) = live_runtime(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let checkpoint = checkpoint_at_head(&tx, &driver);
let other_resolution = envelope(
&op,
"in-resolve-other",
1_700_000_002_000,
KernelInput::ResolveEffect(ResolveEffect {
effect_id: EffectId::new("op-c5-spliced:step:1:effect:0").unwrap(),
outcome: EffectOutcome::Succeeded(EffectSucceeded {
result: EffectSuccess::Provider(ProviderSuccess {
outcome: ProviderOutcome::ContextOverflow(
ProviderContextOverflow::default(),
),
}),
}),
}),
);
let (chain_b, _tx_b, _driver_b) = live_runtime(&[
configure_envelope(&op),
agent_start_envelope(&op),
other_resolution,
]);
let mut spliced = chain_a[..2].to_vec();
spliced.push(chain_b[2].clone());
let report = validate_with_checkpoint(
&blobs(&spliced),
&no_streams(),
&[checkpoint_blob(&checkpoint)],
false,
);
assert_eq!(
rule(&report, 0, "C1").verdict,
Verdict::Pass,
"the splice chains cleanly — C1 is not the witness here"
);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Fail, "{}", c5a.detail);
assert!(c5a.detail.contains("covered step 2"), "{}", c5a.detail);
assert_eq!(report.exit_code(), 1);
}
#[test]
fn a_pruned_journal_degrades_the_c5_anchors_and_c1_keeps_its_verdict() {
let op = operation("op-c5-pruned");
let (chain, tx, driver) = live_runtime(&[
configure_envelope(&op),
agent_start_envelope(&op),
resolve_overflow_envelope(&op, 1),
]);
let checkpoint = checkpoint_at_head(&tx, &driver);
let pruned = chain[2..].to_vec();
let report = validate_with_checkpoint(
&blobs(&pruned),
&no_streams(),
&[checkpoint_blob(&checkpoint)],
true,
);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Degraded, "{}", c5a.detail);
assert!(
c5a.detail.contains("identity anchor is unverifiable"),
"{}",
c5a.detail
);
assert!(
c5a.detail.contains("covered head anchored at step 2"),
"{}",
c5a.detail
);
assert!(
c5a.detail.contains("strict replay skipped"),
"{}",
c5a.detail
);
assert_eq!(rule(&report, 0, "C1").verdict, Verdict::Fail);
assert!(report.has_violations());
assert_eq!(report.exit_code(), 1);
}
#[test]
fn an_unparseable_checkpoint_leaves_evidence_insufficient() {
let op = operation("op-c5-junk");
let chain = live_chain(&[configure_envelope(&op), agent_start_envelope(&op)]);
let report = validate_with_checkpoint(
&blobs(&chain),
&no_streams(),
&[b"{not a checkpoint".to_vec()],
false,
);
assert_eq!(report.unparseable_checkpoints, 1);
assert_eq!(report.checkpoints, Some(0));
assert_eq!(
checkpoint_check(&report, "C5a").verdict,
Verdict::Degraded,
"{}",
checkpoint_check(&report, "C5a").detail
);
assert_eq!(report.exit_code(), 2);
}
#[test]
fn the_published_golden_checkpoints_decode_through_the_checkpoint_plane() {
let fixture_dir =
PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../tests/fixtures/kernel-wire");
let mut blobs = Vec::new();
for name in [
"golden_checkpoint_agent_turn",
"golden_checkpoint_bounded_tail",
] {
let wrapper: serde_json::Value = serde_json::from_str(
&std::fs::read_to_string(fixture_dir.join(format!("{name}.json")))
.expect("fixture readable"),
)
.expect("fixture json");
blobs.push(
serde_json::to_vec(&wrapper["checkpoint"]).expect("the nested checkpoint writes"),
);
}
let report = validate_with_checkpoint(&[] as &[Vec<u8>], &no_streams(), &blobs, false);
assert_eq!(report.checkpoints, Some(2));
assert_eq!(report.unparseable_checkpoints, 0);
let c5a = checkpoint_check(&report, "C5a");
assert_eq!(c5a.verdict, Verdict::Degraded, "{}", c5a.detail);
assert!(c5a.detail.contains("holds no segment"), "{}", c5a.detail);
let c5b = checkpoint_check(&report, "C5b");
assert_ne!(c5b.verdict, Verdict::Fail, "{}", c5b.detail);
assert!(!report.has_violations());
}
}