use serde_json::{Value, json};
use super::checkpoint::KernelCheckpoint;
use super::config::{
ConfigDefaults, ExecutionPolicy, HostEffectSupport, MemoryPolicy, OperationConfig,
ResourceQuota, SkillMetadata,
};
use super::driver::{CanonicalOperationDriver, PlannedStep, SYSCALL_TOOL_NAMES};
use super::effect::{
EffectKindTag, EffectOutcome, EffectSucceeded, EffectSuccess, InlineToolResult,
MemoryAccessBinding, MemoryCapabilities, ProviderCompleted, ProviderMessage, ProviderOutcome,
ProviderStopReason, ProviderSuccess, ToolCall as WireToolCall, ToolResult as WireToolResult,
ToolResultDisposition, ToolResultPayload as WireToolResultPayload, ToolsSuccess,
};
use super::envelope::{
ConfigureOperation, KernelInput, ResolveEffect, StartOperation, WireEnvelope,
};
use super::fault::{KernelFaultCode, PrepareToken};
use super::record::{KernelRecord, RecordPreparation};
use super::restore::{RestoreCost, restore_operation};
use super::root::{InitialContext, LogicalAgentSpec, LogicalTask, RootAgentEntry, RootEntry};
use super::scalar::{BoundedJson, CallId, EffectId, InputId, MemoryBindingId, WireU64};
use super::transaction::{CommittedTransition, InMemoryRecordIndex, KernelTransaction};
const OPERATION: &str = "op-crash-1";
struct Runtime {
tx: KernelTransaction<PlannedStep, InMemoryRecordIndex>,
driver: CanonicalOperationDriver,
journal: Vec<KernelRecord>,
restore_cost: Option<RestoreCost>,
}
impl Runtime {
fn new() -> Self {
Self {
tx: KernelTransaction::new(ConfigDefaults::default(), InMemoryRecordIndex::new()),
driver: CanonicalOperationDriver::new(),
journal: Vec::new(),
restore_cost: None,
}
}
fn prepare(&mut self, envelope: &WireEnvelope) -> RecordPreparation<PlannedStep> {
let Self { tx, driver, .. } = self;
tx.prepare(envelope, |context| driver.plan(context))
}
fn append_and_commit(
&mut self,
preparation: RecordPreparation<PlannedStep>,
) -> CommittedTransition<PlannedStep> {
let token: PrepareToken = preparation
.token()
.unwrap_or_else(|| panic!("expected a prepared step, got {:?}", preparation.fault()))
.clone();
let head = preparation.record().unwrap().record_digest().clone();
let committed = self.tx.commit(&token, &head).expect("commit must succeed");
self.journal.push(committed.record.clone());
committed
}
fn submit(&mut self, envelope: &WireEnvelope) -> CommittedTransition<PlannedStep> {
let preparation = self.prepare(envelope);
let committed = self.append_and_commit(preparation);
self.driver
.note_committed(committed.step_seq)
.expect("the driver folds the step it planned");
committed
}
fn rebuild(records: &[KernelRecord]) -> Runtime {
let Restored {
transaction,
driver,
} = {
let restored = restore_operation(
None,
records,
ConfigDefaults::default(),
InMemoryRecordIndex::from_records(records),
)
.expect("the journal-only ladder runs");
Restored {
transaction: restored.transaction,
driver: restored.driver,
}
};
Runtime {
tx: transaction,
driver,
journal: records.to_vec(),
restore_cost: None,
}
}
fn restore(checkpoint: &KernelCheckpoint, records: &[KernelRecord]) -> Runtime {
let restored = restore_operation(
Some(checkpoint),
records,
ConfigDefaults::default(),
InMemoryRecordIndex::from_records(records),
)
.expect("the checkpoint ladder runs");
Runtime {
tx: restored.transaction,
driver: restored.driver,
journal: records.to_vec(),
restore_cost: Some(restored.cost),
}
}
fn checkpoint(&self) -> super::checkpoint::CheckpointCandidate {
self.tx
.checkpoint_candidate(self.driver.project_logical_state())
.expect("a configured operation has a logical state to checkpoint")
}
}
struct Restored {
transaction: KernelTransaction<PlannedStep, InMemoryRecordIndex>,
driver: CanonicalOperationDriver,
}
fn drive(runtime: &mut Runtime, envelopes: &[WireEnvelope]) {
for envelope in envelopes {
runtime.submit(envelope);
}
}
fn fingerprint(runtime: &Runtime) -> Value {
json!({
"head": runtime.tx.head().map(|head| json!({
"digest": head.digest.as_str(),
"step_seq": head.step_seq.to_string(),
})),
"lifecycle": format!("{:?}", runtime.tx.lifecycle()),
"pending_effects": runtime
.tx
.pending_effects_in_order()
.iter()
.map(|effect| effect.effect_id.as_str())
.collect::<Vec<_>>(),
"terminal": runtime.tx.terminal().map(|t| serde_json::to_value(t).unwrap()),
})
}
fn operation() -> super::scalar::OperationId {
super::scalar::OperationId::new(OPERATION).unwrap()
}
fn envelope(id: &str, observed_at_ms: u64, input: KernelInput) -> WireEnvelope {
WireEnvelope::new(
operation(),
InputId::new(id).unwrap(),
WireU64::new(observed_at_ms),
input,
)
}
fn syscall_tool_catalog() -> Vec<super::effect::ToolSchema> {
SYSCALL_TOOL_NAMES
.iter()
.chain(std::iter::once(&"search"))
.map(|name| super::effect::ToolSchema {
name: (*name).to_string(),
description: String::new(),
parameters: Default::default(),
})
.collect()
}
fn test_agent_spec(goal: &str) -> LogicalAgentSpec {
LogicalAgentSpec {
exposure_baseline: Some(
syscall_tool_catalog()
.into_iter()
.map(|tool| tool.name)
.collect(),
),
..LogicalAgentSpec::new(goal)
}
}
fn syscall_config() -> WireEnvelope {
let config = OperationConfig {
execution_policy: Some(ExecutionPolicy {
max_turns: Some(12),
..ExecutionPolicy::default()
}),
host_effect_support: HostEffectSupport::new([
EffectKindTag::CallProvider,
EffectKindTag::ExecuteTools,
EffectKindTag::LoadPayload,
EffectKindTag::SpawnTasks,
EffectKindTag::PreemptTasks,
EffectKindTag::PersistMemory,
EffectKindTag::QueryMemory,
]),
tool_catalog: syscall_tool_catalog(),
skill_catalog: vec![SkillMetadata {
name: "debug".to_string(),
description: "debug helper".to_string(),
when_to_use: None,
allowed_tools: Vec::new(),
capability_grants: Vec::new(),
effort: None,
estimated_tokens: None,
}],
memory_access: Some(MemoryAccessBinding {
binding_id: MemoryBindingId::new("mem-binding-1").unwrap(),
capabilities: MemoryCapabilities {
read: true,
write: true,
},
}),
memory_policy: Some(MemoryPolicy {
retrieval_top_k: Some(4),
..MemoryPolicy::default()
}),
resource_quota: Some(ResourceQuota {
max_workflow_nodes: Some(3),
..ResourceQuota::default()
}),
..OperationConfig::default()
};
envelope(
"in-configure",
1_700_000_000_000,
KernelInput::ConfigureOperation(ConfigureOperation { config }),
)
}
fn agent_start(id: &str, observed_at_ms: u64) -> WireEnvelope {
envelope(
id,
observed_at_ms,
KernelInput::StartOperation(StartOperation {
entry: RootEntry::Agent(RootAgentEntry {
task: LogicalTask::new("write the research brief"),
run_spec: Some(test_agent_spec("write the research brief")),
}),
initial_context: InitialContext::default(),
}),
)
}
fn effect_id(step_seq: WireU64) -> EffectId {
EffectId::new(format!("{OPERATION}:step:{step_seq}:effect:0")).unwrap()
}
fn tool_call(call_id: &str, name: &str, arguments: Value) -> WireToolCall {
WireToolCall {
call_id: CallId::new(call_id).unwrap(),
name: name.to_string(),
arguments: BoundedJson::new(arguments).unwrap(),
}
}
fn resolved(id: &str, at: u64, effect: &EffectId, result: EffectSuccess) -> WireEnvelope {
envelope(
id,
at,
KernelInput::ResolveEffect(ResolveEffect {
effect_id: effect.clone(),
outcome: EffectOutcome::Succeeded(EffectSucceeded { result }),
}),
)
}
fn provider_result(
id: &str,
observed_at_ms: u64,
effect: &EffectId,
calls: Vec<WireToolCall>,
) -> WireEnvelope {
resolved(
id,
observed_at_ms,
effect,
EffectSuccess::Provider(ProviderSuccess {
outcome: ProviderOutcome::Completed(ProviderCompleted {
message: ProviderMessage {
role: super::root::MessageRole::Assistant,
content: String::new(),
tool_calls: calls,
tool_call_id: None,
tokens: None,
},
observed_input_tokens: None,
observed_output_tokens: None,
stop_reason: None,
}),
}),
)
}
fn tools_resolved(
id: &str,
at: u64,
effect: &EffectId,
results: &[(&str, &str, bool)],
) -> WireEnvelope {
resolved(
id,
at,
effect,
EffectSuccess::Tools(ToolsSuccess {
results: results
.iter()
.map(|(call_id, output, is_error)| {
WireToolResultPayload::Inline(InlineToolResult {
call_id: CallId::new(*call_id).unwrap(),
result: WireToolResult {
output: (*output).to_string(),
durable_content: None,
is_error: *is_error,
disposition: ToolResultDisposition::Recoverable,
},
})
})
.collect(),
measurements: Vec::new(),
}),
)
}
fn provider_answer(id: &str, at: u64, effect: &EffectId, text: &str) -> WireEnvelope {
resolved(
id,
at,
effect,
EffectSuccess::Provider(ProviderSuccess {
outcome: ProviderOutcome::Completed(ProviderCompleted {
message: ProviderMessage {
role: super::root::MessageRole::Assistant,
content: text.to_string(),
tool_calls: Vec::new(),
tool_call_id: None,
tokens: None,
},
observed_input_tokens: None,
observed_output_tokens: None,
stop_reason: Some(ProviderStopReason::EndTurn),
}),
}),
)
}
fn turn_envelopes() -> Vec<WireEnvelope> {
vec![
syscall_config(),
agent_start("in-start", 1_700_000_001_000),
provider_result(
"in-acted",
1_700_000_002_000,
&effect_id(WireU64::new(1)),
vec![tool_call("call-1", "search", json!({"q": "sources"}))],
),
tools_resolved(
"in-results",
1_700_000_003_000,
&effect_id(WireU64::new(2)),
&[("call-1", "three sources found", false)],
),
provider_result(
"in-acted-2",
1_700_000_004_000,
&effect_id(WireU64::new(3)),
vec![tool_call("call-2", "search", json!({"q": "more"}))],
),
tools_resolved(
"in-results-2",
1_700_000_005_000,
&effect_id(WireU64::new(4)),
&[("call-2", "two more sources", false)],
),
provider_answer(
"in-answer",
1_700_000_006_000,
&effect_id(WireU64::new(5)),
"the brief cites five sources",
),
]
}
fn digests(runtime: &Runtime) -> Vec<String> {
runtime
.journal
.iter()
.map(|record| record.record_digest().to_string())
.collect()
}
#[test]
fn crash_before_append_leaves_no_residue_and_redelivers_the_same_step() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
runtime.submit(&envelopes[0]);
let preparation = runtime.prepare(&envelopes[1]);
let crashed_digest = preparation.record().unwrap().record_digest().clone();
drop(preparation);
assert!(
runtime.journal.len() == 1,
"no durable residue: the journal holds only the configure record"
);
let mut recovered = Runtime::rebuild(&runtime.journal);
let prepared = match recovered.prepare(&envelopes[1]) {
RecordPreparation::Prepared(prepared) => prepared,
other => panic!(
"an input that never reached the journal must be accepted again, got {:?}",
other.fault().map(|fault| fault.code)
),
};
assert_eq!(
prepared.record.record_digest(),
&crashed_digest,
"the same input plans the same step — the crash never happened, durably speaking",
);
let committed = recovered.append_and_commit(RecordPreparation::Prepared(prepared));
recovered
.driver
.note_committed(committed.step_seq)
.expect("the driver folds the step it planned");
drive(&mut recovered, &envelopes[2..]);
let mut uninterrupted = Runtime::new();
drive(&mut uninterrupted, &envelopes);
assert_eq!(digests(&recovered), digests(&uninterrupted));
}
#[test]
fn crash_after_append_before_commit_reaches_one_head_both_ways() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
runtime.submit(&envelopes[0]);
let preparation = runtime.prepare(&envelopes[1]);
let token: PrepareToken = preparation.token().unwrap().clone();
let appended_head = preparation.record().unwrap().record_digest().clone();
runtime.journal.push(preparation.record().unwrap().clone());
runtime
.tx
.commit(&token, &appended_head)
.expect("the append happened, so the commit lands");
runtime
.driver
.note_committed(runtime.tx.head().unwrap().step_seq)
.expect("the fold catches up");
let journal_snapshot = runtime.journal.clone();
let rebuilt = Runtime::rebuild(&journal_snapshot);
assert_eq!(
runtime.tx.head().map(|head| head.digest),
rebuilt.tx.head().map(|head| head.digest),
"commit-after-the-fact and rebuild agree on the head",
);
assert_eq!(
runtime.tx.head().unwrap().digest,
appended_head,
"and that head is exactly the appended record",
);
drive(&mut runtime, &envelopes[2..]);
let mut rebuilt = rebuilt;
drive(&mut rebuilt, &envelopes[2..]);
let mut uninterrupted = Runtime::new();
drive(&mut uninterrupted, &envelopes);
assert_eq!(digests(&runtime), digests(&uninterrupted));
assert_eq!(digests(&rebuilt), digests(&uninterrupted));
}
#[test]
fn crash_after_commit_before_publish_republishes_the_same_effect_ids() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
runtime.submit(&envelopes[0]);
let started = runtime.submit(&envelopes[1]);
let published: Vec<String> = started
.published_effects()
.iter()
.map(|effect| effect.effect_id.as_str().to_string())
.collect();
assert_eq!(
published.len(),
1,
"the fixture publishes exactly one effect"
);
assert_eq!(
started.published_effects()[0].tag(),
EffectKindTag::CallProvider
);
let rebuilt = Runtime::rebuild(&runtime.journal);
let republished: Vec<String> = rebuilt
.tx
.pending_effects_in_order()
.iter()
.map(|effect| effect.effect_id.as_str().to_string())
.collect();
assert_eq!(
republished, published,
"§5g-1 · the recovery hands back the same effect identities to (re-)execute",
);
assert_eq!(fingerprint(&rebuilt), fingerprint(&runtime));
assert_eq!(
rebuilt.tx.head().unwrap().digest,
started.record.record_digest().clone(),
"the rebuild stands on the committed record, not before it",
);
}
#[test]
fn crash_after_publish_before_resolve_redelivers_to_the_same_record() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
runtime.submit(&envelopes[0]);
runtime.submit(&envelopes[1]);
let rebuilt = Runtime::rebuild(&runtime.journal);
let mut rebuilt = rebuilt;
rebuilt.submit(&envelopes[2]);
let mut uninterrupted = Runtime::new();
drive(&mut uninterrupted, &envelopes[..3]);
assert_eq!(
digests(&rebuilt),
digests(&uninterrupted),
"the post-crash resolution writes the same record the uninterrupted run wrote",
);
assert_eq!(fingerprint(&rebuilt), fingerprint(&uninterrupted));
}
#[test]
fn crash_before_install_a_candidate_is_pure_and_regenerable() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
drive(&mut runtime, &envelopes[..4]);
let head_before = runtime.tx.head();
let tail_before = runtime.tx.tail_usage();
let boundary_before = runtime.tx.checkpoint_boundary();
let first = runtime
.checkpoint()
.decode()
.expect("the host decoded the blob");
let first_digest = first.checkpoint_digest().clone();
drop(first);
let regenerated = runtime
.checkpoint()
.decode()
.expect("regenerates just as well");
assert_eq!(
runtime.tx.head(),
head_before,
"generation is a read: the durable head did not move",
);
assert_eq!(runtime.tx.tail_usage(), tail_before);
assert_eq!(runtime.tx.checkpoint_boundary(), boundary_before);
assert_eq!(
regenerated.checkpoint_digest(),
&first_digest,
"the regenerated candidate is byte-identical — nothing was in flight to lose",
);
}
#[test]
fn crash_after_install_before_ack_still_installs_acks_and_restores() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
drive(&mut runtime, &envelopes[..4]);
let candidate = runtime.checkpoint();
let installed = candidate.decode().expect("the host installed this blob");
let installed_digest = installed.checkpoint_digest().clone();
let reread = candidate.decode().expect("re-read from the store");
assert_eq!(reread.checkpoint_digest(), &installed_digest);
let usage_before = runtime.tx.tail_usage();
assert!(usage_before.records > 0, "the fixture carries a live tail");
let usage = runtime
.tx
.note_checkpoint_acked(&installed.boundary())
.expect("the ack is a retention signal and lands whenever it is delivered");
assert_eq!(
usage.records, 0,
"a full-state candidate at the head covers the whole tail",
);
assert_eq!(
runtime.tx.head().map(|head| head.digest),
Some(installed.covered_transaction_head_digest().clone()),
"the ack reclaims accounting, never history",
);
let above: Vec<KernelRecord> = runtime
.journal
.iter()
.filter(|record| record.step_seq().get() > installed.through_step_seq().get())
.cloned()
.collect();
let restored = Runtime::restore(&installed, &above);
assert_eq!(restored.restore_cost.unwrap().records_before_checkpoint, 0);
assert_eq!(
restored
.tx
.pending_effects_in_order()
.iter()
.map(|effect| effect.effect_id.as_str().to_string())
.collect::<Vec<_>>(),
vec![effect_id(WireU64::new(3)).to_string()],
"§5g-1 · the effect the operation is waiting on is exposed again",
);
}
#[test]
fn crash_after_ack_before_prune_reack_converges() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
drive(&mut runtime, &envelopes[..4]);
let checkpoint = runtime.checkpoint().decode().expect("verifies");
let boundary = checkpoint.boundary();
runtime
.tx
.note_checkpoint_acked(&boundary)
.expect("the boundary names a prefix of this journal");
let head_after_ack = runtime.tx.head();
let mut rebuilt = Runtime::rebuild(&runtime.journal);
assert_eq!(
rebuilt.tx.head(),
head_after_ack,
"the un-pruned prefix replays to the same head",
);
assert!(
rebuilt.tx.tail_usage().records > 0,
"the rebuild re-materialises the tail the ack had reclaimed — physical state is honest again",
);
let usage = rebuilt
.tx
.note_checkpoint_acked(&boundary)
.expect("the replayed prune hits the same boundary");
assert_eq!(
usage.records, 0,
"the reclaim lands exactly as the first one did"
);
let redelivered = rebuilt.prepare(&envelopes[1]);
assert!(
redelivered.token().is_none(),
"a replay offers nothing to commit",
);
let RecordPreparation::Replayed(replay) = redelivered else {
panic!("an acked input must not be accepted a second time");
};
assert_eq!(
replay.record_digest,
*runtime.journal[1].record_digest(),
"the answer is still the original record digest",
);
assert!(
replay.record.is_some() && replay.committed_step.is_some(),
"the un-pruned journal reproduces the step from disk rather than re-executing it",
);
}
#[test]
fn append_conflict_discards_the_candidate_and_rebuild_replays() {
let envelopes = turn_envelopes();
let mut base = Runtime::new();
drive(&mut base, &envelopes[..2]);
let shared = base.journal.clone();
let mut winner = Runtime::rebuild(&shared);
let mut loser = Runtime::rebuild(&shared);
let preparation = loser.prepare(&envelopes[2]);
let token: PrepareToken = preparation.token().unwrap().clone();
let loser_head = preparation.record().unwrap().record_digest().clone();
winner.submit(&envelopes[2]);
let actual_head = winner.tx.head().unwrap().digest.clone();
let conflict = loser.tx.note_append_conflict(&token, Some(&actual_head));
assert_eq!(conflict.code, KernelFaultCode::TransactionConflict);
assert!(
conflict.message.contains(&loser_head.to_string()),
"the fault names the head the loser expected",
);
assert!(
conflict.message.contains(&actual_head.to_string()),
"and the head the journal actually holds",
);
let refused = loser.prepare(&envelopes[3]);
assert_eq!(
refused.fault().map(|fault| fault.code),
Some(KernelFaultCode::TransactionConflict),
"the poison propagates until the host rebuilds",
);
let mut recovered = Runtime::rebuild(&winner.journal);
let redelivered = recovered.prepare(&envelopes[2]);
assert!(
redelivered.token().is_none(),
"the replay offers nothing to commit",
);
let RecordPreparation::Replayed(replay) = redelivered else {
panic!("the winner's record answers the redelivery");
};
assert_eq!(
replay.record_digest,
*winner.journal[2].record_digest(),
"the disputed input resolves to the winner's record, not a competing one",
);
recovered.submit(&envelopes[3]);
winner.submit(&envelopes[3]);
let mut control = Runtime::new();
drive(&mut control, &envelopes[..4]);
assert_eq!(digests(&winner), digests(&control));
assert_eq!(digests(&recovered), digests(&control));
assert_eq!(fingerprint(&recovered), fingerprint(&control));
}
#[test]
fn checkpoint_plus_incomplete_tail_restores_to_head() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
drive(&mut runtime, &envelopes[..2]);
let checkpoint = runtime.checkpoint().decode().expect("verifies");
runtime
.tx
.note_checkpoint_acked(&checkpoint.boundary())
.expect("the host acked the install");
drive(&mut runtime, &envelopes[2..4]);
let pruned: Vec<KernelRecord> = runtime
.journal
.iter()
.filter(|record| record.step_seq().get() > checkpoint.through_step_seq().get())
.cloned()
.collect();
assert_eq!(pruned.len(), 2, "the host pruned the covered prefix");
let restored = Runtime::restore(&checkpoint, &pruned);
let cost = restored.restore_cost.unwrap();
assert_eq!(
cost.records_before_checkpoint, 0,
"nothing below the base is read"
);
assert_eq!(
cost.tail_inputs_replayed, 0,
"a full-state candidate carries no tail of its own"
);
assert_eq!(
cost.records_after_checkpoint, 2,
"the incomplete tail is what is replayed"
);
let mut uninterrupted = Runtime::new();
drive(&mut uninterrupted, &envelopes[..4]);
assert_eq!(fingerprint(&restored), fingerprint(&uninterrupted));
let mut restored = restored;
drive(&mut restored, &envelopes[4..]);
drive(&mut uninterrupted, &envelopes[4..]);
let offset = checkpoint.through_step_seq().get() as usize + 1;
assert_eq!(
digests(&restored),
digests(&uninterrupted)[offset..],
"the records above the checkpoint are written identically",
);
}
#[test]
fn redelivery_below_the_base_is_an_ack_not_a_step() {
let envelopes = turn_envelopes();
let mut runtime = Runtime::new();
drive(&mut runtime, &envelopes[..4]);
let checkpoint = runtime.checkpoint().decode().expect("verifies");
runtime
.tx
.note_checkpoint_acked(&checkpoint.boundary())
.expect("the host acked the install");
let mut restored = Runtime::restore(&checkpoint, &[]);
assert_eq!(
restored.restore_cost.unwrap().records_before_checkpoint,
0,
"the covered prefix is gone from every surface",
);
let head_before = restored.tx.head();
for (index, envelope) in envelopes.iter().take(2).enumerate() {
let redelivered = restored.prepare(envelope);
let RecordPreparation::Replayed(replay) = redelivered else {
panic!("a redelivery below the checkpoint base must not be accepted again");
};
assert_eq!(
replay.record_digest,
*runtime.journal[index].record_digest(),
"§12.3 rule 10 · the answer is the original step and record digest",
);
assert!(
replay.committed_step.is_none() && replay.record.is_none(),
"idempotent acknowledgement, not step reproduction",
);
assert!(replay.step_seq == runtime.journal[index].step_seq());
}
assert_eq!(
restored.tx.head(),
head_before,
"an acknowledgement writes nothing",
);
let offset = checkpoint.through_step_seq().get() as usize + 1;
drive(&mut restored, &envelopes[offset..]);
let mut uninterrupted = Runtime::new();
drive(&mut uninterrupted, &envelopes);
assert_eq!(
digests(&restored),
digests(&uninterrupted)[offset..],
"the post-checkpoint history is written identically",
);
assert_eq!(fingerprint(&restored), fingerprint(&uninterrupted));
}