1use super::state::{ArchivedCommandRequestOutcome, ArchivedIngressRequest};
2use super::{
3 ADMISSION_CURSOR_FORMAT_VERSION, BoundaryReceipt, BoundaryRecord, BoundaryRequest, CanwuError,
4 CauseRef, CommandAttemptOutcome, CommandAttemptRecord, CommandEnvelope, CommandOutcome,
5 CommandReceipt, CommandRecord, CommandRequest, CommitmentRoots, DecisionState,
6 DeterministicRng, DomainRecord, DomainRecordRef, DomainRecordSchema, DomainRecordType,
7 DomainRecordVersionRef, DomainRecordVersionSource, ENGINE_VERSION, ErrorCode, EvidenceRef,
8 IngressPayload, IngressReceipt, IngressRecord, KeyedDrawReservation, KnowledgeSnapshot,
9 PayloadProperty, PayloadSchema, PayloadValueType, PluginComponentRecord, PluginDescriptor,
10 PluginIngressRequest, RandomDrawAddress, RandomDrawRecord, RandomStreamState,
11 RunConfigurationSnapshot, RunManifest, RuntimeEvidence, SNAPSHOT_FORMAT_VERSION,
12 STATE_REVISION_FORMAT_VERSION, Scenario, ScheduledAction, ScheduledRecord, SchemaRegistry,
13 SimDuration, SimEvent, SimTime, Simulation, SimulationPlugin, SystemCadence,
14 TypedDomainRecordRef, WorldSnapshot, has_unqueued_command_history, invalid_snapshot_error,
15 is_one_u64, is_zero_u32, is_zero_u64, legacy_v4, one_u64,
16};
17use serde::{Deserialize, Serialize};
18use serde_json::Value;
19use std::collections::{BTreeMap, BTreeSet};
20
21#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
22pub struct SimulationSnapshot {
23 pub engine_version: String,
24 pub snapshot_format_version: u32,
25 #[serde(default)]
26 pub run_manifest: Option<RunManifest>,
27 #[serde(default)]
28 pub run_manifest_hash: String,
29 #[serde(default, skip_serializing_if = "Option::is_none")]
30 pub run_configuration: Option<RunConfigurationSnapshot>,
31 #[serde(default)]
32 pub checkpoint_hash: String,
33 #[serde(default, skip_serializing_if = "is_zero_u32")]
34 pub commitment_format_version: u32,
36 #[serde(default, skip_serializing_if = "Option::is_none")]
37 pub commitment_roots: Option<CommitmentRoots>,
39 #[serde(default)]
40 pub revision_format_version: u32,
42 #[serde(default, skip_serializing_if = "is_zero_u64")]
43 pub state_revision: u64,
45 #[serde(default, skip_serializing_if = "is_zero_u32")]
46 pub replay_revision_format_version: u32,
48 #[serde(default, skip_serializing_if = "is_zero_u32")]
49 pub admission_cursor_format_version: u32,
51 #[serde(default, skip_serializing_if = "is_zero_u64")]
52 pub admitted_attempt_count: u64,
54 #[serde(default, skip_serializing_if = "is_zero_u64")]
55 pub admitted_command_count: u64,
57 #[serde(default, skip_serializing_if = "is_zero_u64")]
58 pub admitted_event_count: u64,
60 pub initial_time: SimTime,
61 #[serde(default, skip_serializing_if = "Option::is_none")]
62 pub initial_scenario: Option<Scenario>,
63 pub now: SimTime,
64 pub plugin_registration_closed: bool,
65 pub world: WorldSnapshot,
66 pub knowledge: KnowledgeSnapshot,
67 pub events: Vec<SimEvent>,
68 pub commands: Vec<CommandRecord>,
69 #[serde(default, skip_serializing_if = "Vec::is_empty")]
70 pub command_attempts: Vec<CommandAttemptRecord>,
71 #[serde(default, skip_serializing_if = "Vec::is_empty")]
72 pub ingress: Vec<IngressRecord>,
73 #[serde(default)]
74 pub boundaries: Vec<BoundaryRecord>,
75 pub plugin_components: Vec<PluginComponentRecord>,
76 #[serde(default, skip_serializing_if = "Vec::is_empty")]
77 pub domain_records: Vec<DomainRecord>,
78 #[serde(default, skip_serializing_if = "DecisionState::is_empty")]
79 pub decisions: DecisionState,
80 pub plugin_descriptors: Vec<PluginDescriptor>,
81 pub schema: SchemaRegistry,
82 #[serde(default)]
83 pub root_seed: u64,
84 #[serde(default)]
85 pub random_streams: Vec<RandomStreamState>,
86 #[serde(default)]
87 pub random_draws: Vec<RandomDrawRecord>,
88 pub(super) scheduled: Vec<ScheduledRecord>,
89 #[serde(default, rename = "rng", skip_serializing_if = "Option::is_none")]
90 pub(super) legacy_rng: Option<DeterministicRng>,
91 pub(super) next_event_id: u64,
92 pub(super) next_command_id: u64,
93 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
94 pub(super) next_command_attempt_id: u64,
95 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
96 pub(super) next_ingress_id: u64,
97 #[serde(default)]
98 pub(super) next_boundary_id: u64,
99 #[serde(default)]
100 pub(super) next_random_draw_id: u64,
101 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
102 pub(super) next_knowledge_record_id: u64,
103 pub(super) next_schedule_sequence: u64,
104 pub(super) next_correlation_id: u64,
105 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
106 pub(super) next_decision_trace_id: u64,
107}
108
109#[derive(Clone, Debug, PartialEq, Serialize)]
110pub struct ReplayJournal {
112 pub engine_version: String,
113 pub snapshot_format_version: u32,
114 pub root_seed: u64,
115 pub run_manifest: RunManifest,
116 pub run_manifest_hash: String,
117 pub run_configuration: RunConfigurationSnapshot,
118 pub plugin_descriptors: Vec<PluginDescriptor>,
119 pub plugin_registration_closed: bool,
120 pub commands: Vec<CommandRecord>,
121 pub command_attempts: Vec<CommandAttemptRecord>,
122 #[serde(default, skip_serializing_if = "Vec::is_empty")]
123 pub ingress: Vec<IngressRecord>,
124 pub boundaries: Vec<BoundaryRecord>,
125 pub final_time: SimTime,
126 pub checkpoint_hash: String,
127 pub commitment_format_version: u32,
129 pub revision_format_version: u32,
131 pub final_revision: u64,
133}
134
135#[derive(Deserialize, Serialize)]
136#[serde(deny_unknown_fields)]
137pub(super) struct ReplayJournalWire {
138 pub(super) engine_version: String,
139 pub(super) snapshot_format_version: u32,
140 pub(super) root_seed: u64,
141 pub(super) run_manifest: RunManifest,
142 pub(super) run_manifest_hash: String,
143 #[serde(default)]
144 pub(super) run_configuration: Option<RunConfigurationSnapshot>,
145 pub(super) plugin_descriptors: Vec<PluginDescriptor>,
146 pub(super) plugin_registration_closed: bool,
147 pub(super) commands: Vec<CommandRecord>,
148 #[serde(default)]
149 pub(super) command_attempts: Vec<CommandAttemptRecord>,
150 #[serde(default)]
151 pub(super) ingress: Vec<IngressRecord>,
152 pub(super) boundaries: Vec<BoundaryRecord>,
153 pub(super) final_time: SimTime,
154 pub(super) checkpoint_hash: String,
155 #[serde(default)]
156 pub(super) commitment_format_version: u32,
157 #[serde(default)]
158 pub(super) revision_format_version: u32,
159 #[serde(default)]
160 pub(super) final_revision: u64,
161}
162
163impl<'de> Deserialize<'de> for ReplayJournal {
164 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
165 where
166 D: serde::Deserializer<'de>,
167 {
168 let value = Value::deserialize(deserializer)?;
169 legacy_v4::deserialize_replay_value(&value).map_err(serde::de::Error::custom)
170 }
171}
172
173pub const CHECKPOINT_JOURNAL_FORMAT_VERSION: u32 = 1;
175const ARCHIVED_SEGMENT_MANIFEST_DOMAIN: &str = "canwu.evidence.archived-segment-manifest.v1";
176const ARCHIVED_RECEIPT_DOMAIN: &str = "canwu.evidence.archived-receipts.v1";
177const EVIDENCE_DEPENDENCY_DOMAIN: &str = "canwu.evidence.dependencies.v1";
178const KEYED_RESERVATION_DOMAIN: &str = "canwu.random.keyed-reservations.v1";
179pub const PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD: &str =
181 "canwu_payload_required_evidence_continuation";
182pub const PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION: u32 = 1;
184
185#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
186#[serde(rename_all = "snake_case")]
187pub enum EvidenceJournalKind {
188 Event,
189 Command,
190 CommandAttempt,
191 Ingress,
192 Boundary,
193 RandomDraw,
194}
195
196#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
197#[serde(tag = "type", rename_all = "snake_case")]
198pub enum EvidenceNestedLocator {
199 None,
200 BoundaryRecordChange { change_index: u64 },
201}
202
203#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
204pub struct EvidenceItemLocator {
205 pub journal: EvidenceJournalKind,
206 pub absolute_index: u64,
207 pub nested: EvidenceNestedLocator,
208}
209
210#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
211pub struct ArchivedEvidenceLocator {
212 pub segment_id: String,
213 pub item: EvidenceItemLocator,
214}
215
216#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
217pub struct ArchivedEvidenceReceipt {
218 pub evidence: EvidenceRef,
219 pub locator: ArchivedEvidenceLocator,
220 pub evidence_index_leaf: u64,
221 pub item_commitment: String,
222 pub merkle_path: Vec<String>,
223}
224
225#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
226pub struct EvidenceIndexEntry {
227 pub reference: EvidenceRef,
228 pub item: EvidenceItemLocator,
229 pub item_commitment: String,
230}
231
232#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
233pub struct EvidenceJournalRoots {
234 pub events: String,
235 pub commands: String,
236 pub command_attempts: String,
237 pub ingress: String,
238 pub boundaries: String,
239 pub random_draws: String,
240}
241
242#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
243pub struct ArchivedSegmentHeader {
244 pub segment_id: String,
245 pub start: EvidenceCursor,
246 pub end: EvidenceCursor,
247 pub journal_roots: EvidenceJournalRoots,
248 pub evidence_index_root: String,
249 pub evidence_index_entry_count: u64,
250}
251
252#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
253#[serde(rename_all = "snake_case")]
254pub enum EvidenceRequirement {
255 IdentityOnly,
257 PayloadRequired,
259}
260
261#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
269#[serde(deny_unknown_fields)]
270pub struct PayloadRequiredEvidenceContinuationV1 {
271 pub format_version: u32,
272 pub active: bool,
273 pub dependencies: Vec<EvidenceRef>,
274}
275
276impl PayloadRequiredEvidenceContinuationV1 {
277 #[must_use]
279 pub fn active(dependencies: Vec<EvidenceRef>) -> Self {
280 Self {
281 format_version: PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION,
282 active: true,
283 dependencies,
284 }
285 }
286
287 #[must_use]
289 pub const fn completed() -> Self {
290 Self {
291 format_version: PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION,
292 active: false,
293 dependencies: Vec::new(),
294 }
295 }
296}
297
298#[must_use]
301pub fn payload_required_evidence_continuation_property_v1() -> PayloadProperty {
302 PayloadProperty {
303 value_type: PayloadValueType::Object,
304 required: true,
305 }
306}
307#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
308pub struct EvidenceDependency {
309 pub reference: EvidenceRef,
310 pub requirement: EvidenceRequirement,
311}
312
313#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
314pub struct EvidenceArchiveIndex {
315 pub header: ArchivedSegmentHeader,
316 pub entries: Vec<EvidenceIndexEntry>,
317}
318
319pub trait ArchiveProvider {
320 fn load_evidence_segment(
321 &self,
322 segment_id: &str,
323 ) -> Result<Option<EvidenceJournalSegment>, CanwuError>;
324}
325
326pub trait ArchiveStore: ArchiveProvider {
327 fn store_evidence_segment(
328 &self,
329 segment: &EvidenceJournalSegment,
330 ) -> Result<ArchiveStoreOutcome, CanwuError>;
331}
332
333#[derive(Clone, Copy, Debug, Eq, PartialEq)]
334pub enum ArchiveStoreOutcome {
335 Stored,
336 AlreadyPresent,
337}
338
339#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
340pub struct EvidenceSealToken {
341 pub source_state_hash: String,
342 pub source_checkpoint_hash: String,
343 pub source_end: EvidenceCursor,
344 pub segment_id: String,
345 pub target_checkpoint_hash: String,
346 pub token_hash: String,
347}
348
349#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
350pub struct PreparedEvidenceSeal {
351 pub token: EvidenceSealToken,
352 pub segment: EvidenceJournalSegment,
353}
354
355#[derive(Serialize)]
356struct EvidenceIndexLeafMaterial<'a> {
357 format_version: u32,
358 reference: &'a EvidenceRef,
359 item: &'a EvidenceItemLocator,
360 item_commitment: &'a str,
361}
362
363#[derive(Serialize)]
364struct ArchivedSegmentHeaderMaterial<'a> {
365 start: EvidenceCursor,
366 end: EvidenceCursor,
367 journal_roots: &'a EvidenceJournalRoots,
368 evidence_index_root: &'a str,
369 evidence_index_entry_count: u64,
370}
371
372fn archive_error(message: impl Into<String>) -> CanwuError {
373 CanwuError::new(ErrorCode::InvalidArchive, message)
374}
375
376fn decode_hash(value: &str, label: &str) -> Result<[u8; 32], CanwuError> {
377 if value.len() != 64
378 || value
379 .bytes()
380 .any(|byte| !byte.is_ascii_hexdigit() || byte.is_ascii_uppercase())
381 {
382 return Err(archive_error(format!(
383 "{label} must be 32-byte lower-case hex"
384 )));
385 }
386 let mut bytes = [0_u8; 32];
387 for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() {
388 let digit = |byte: u8| match byte {
389 b'0'..=b'9' => Some(byte - b'0'),
390 b'a'..=b'f' => Some(byte - b'a' + 10),
391 _ => None,
392 };
393 bytes[index] = digit(pair[0])
394 .and_then(|high| digit(pair[1]).map(|low| (high << 4) | low))
395 .ok_or_else(|| archive_error(format!("{label} contains invalid hex")))?;
396 }
397 Ok(bytes)
398}
399
400fn archive_node(left: [u8; 32], right: [u8; 32]) -> [u8; 32] {
401 let mut hasher = blake3::Hasher::new();
402 hasher.update(b"canwu.evidence.index.node.v1");
403 hasher.update(&[0]);
404 hasher.update(&left);
405 hasher.update(&right);
406 *hasher.finalize().as_bytes()
407}
408
409fn archive_empty_root() -> String {
410 let mut hasher = blake3::Hasher::new();
411 hasher.update(b"canwu.evidence.index.empty.v1");
412 hasher.update(&[0]);
413 hasher.finalize().to_hex().to_string()
414}
415
416fn skipped_commitment_root<T: Serialize>(
417 domain: &str,
418 values: &[T],
419) -> Result<Option<String>, CanwuError> {
420 if values.is_empty() {
421 Ok(None)
422 } else {
423 super::canonical_hash(domain, values).map(Some)
424 }
425}
426
427fn validate_skipped_commitment_root<T: Serialize>(
428 root: Option<&str>,
429 domain: &str,
430 values: &[T],
431 label: &str,
432) -> Result<(), CanwuError> {
433 let expected = skipped_commitment_root(domain, values)?;
434 if root != expected.as_deref() {
435 return Err(invalid_snapshot_error(format!(
436 "compact continuation {label} does not match its canonical material"
437 )));
438 }
439 Ok(())
440}
441
442fn promote_dependency(
443 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
444 reference: EvidenceRef,
445 requirement: EvidenceRequirement,
446) {
447 dependencies
448 .entry(reference)
449 .and_modify(|current| *current = (*current).max(requirement))
450 .or_insert(requirement);
451}
452fn schema_declares_payload_required_continuation(schema: &DomainRecordSchema) -> bool {
453 matches!(
454 &schema.payload_schema,
455 PayloadSchema::Object { properties, .. }
456 if properties.get(PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD)
457 == Some(&payload_required_evidence_continuation_property_v1())
458 )
459}
460
461fn add_payload_required_continuation_dependencies(
462 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
463 record: &DomainRecord,
464 schema: &DomainRecordSchema,
465) -> Result<(), CanwuError> {
466 if !record.is_active() || !schema_declares_payload_required_continuation(schema) {
467 return Ok(());
468 }
469 let continuation = record
470 .payload
471 .get(PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD)
472 .ok_or_else(|| {
473 CanwuError::new(
474 ErrorCode::ArchiveNotReady,
475 format!(
476 "payload-required continuation record {} is missing its schema-declared field",
477 record.reference
478 ),
479 )
480 })?;
481 let continuation: PayloadRequiredEvidenceContinuationV1 =
482 serde_json::from_value(continuation.clone()).map_err(|error| {
483 CanwuError::new(
484 ErrorCode::ArchiveNotReady,
485 format!(
486 "payload-required continuation record {} is invalid: {error}",
487 record.reference
488 ),
489 )
490 })?;
491 if continuation.format_version != PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION {
492 return Err(CanwuError::new(
493 ErrorCode::ArchiveNotReady,
494 format!(
495 "payload-required continuation record {} uses unsupported format {}",
496 record.reference, continuation.format_version
497 ),
498 ));
499 }
500 if !continuation.active {
501 if continuation == PayloadRequiredEvidenceContinuationV1::completed() {
502 return Ok(());
503 }
504 return Err(CanwuError::new(
505 ErrorCode::ArchiveNotReady,
506 format!(
507 "completed payload-required continuation record {} retains dependencies",
508 record.reference
509 ),
510 ));
511 }
512 let continuation = PayloadRequiredEvidenceContinuationV1::active(continuation.dependencies);
513 if continuation.dependencies.is_empty()
514 || continuation
515 .dependencies
516 .windows(2)
517 .any(|pair| pair[0] >= pair[1])
518 {
519 return Err(CanwuError::new(
520 ErrorCode::ArchiveNotReady,
521 format!(
522 "active payload-required continuation record {} needs sorted unique dependencies",
523 record.reference
524 ),
525 ));
526 }
527 if continuation.dependencies.iter().any(|reference| {
528 matches!(
529 reference,
530 EvidenceRef::DomainRecordVersion(version)
531 if version.record == record.reference && version.version == record.version
532 )
533 }) {
534 return Err(CanwuError::new(
535 ErrorCode::ArchiveNotReady,
536 format!(
537 "payload-required continuation record {} cannot depend on its own version",
538 record.reference
539 ),
540 ));
541 }
542 for reference in continuation.dependencies {
543 if matches!(
544 &reference,
545 EvidenceRef::DomainRecordVersion(version)
546 if matches!(version.established_by, DomainRecordVersionSource::InitialScenario)
547 ) {
548 return Err(CanwuError::new(
549 ErrorCode::ArchiveNotReady,
550 format!(
551 "payload-required continuation record {} cannot depend on an initial-scenario payload",
552 record.reference
553 ),
554 ));
555 }
556 promote_dependency(
557 dependencies,
558 reference,
559 EvidenceRequirement::PayloadRequired,
560 );
561 }
562 Ok(())
563}
564
565fn required_archived_receipt_references(
566 dependencies: &[EvidenceDependency],
567 reservations: &[KeyedDrawReservation],
568) -> BTreeSet<EvidenceRef> {
569 let mut required = dependencies
570 .iter()
571 .map(|dependency| dependency.reference.clone())
572 .collect::<BTreeSet<_>>();
573 required.extend(
574 reservations
575 .iter()
576 .map(|reservation| reservation.draw_receipt.evidence.clone()),
577 );
578 required
579}
580
581fn retain_reachable_archived_evidence_receipts(
582 receipts: &mut BTreeMap<EvidenceRef, ArchivedEvidenceReceipt>,
583 dependencies: &[EvidenceDependency],
584 reservations: &[KeyedDrawReservation],
585) {
586 let required = required_archived_receipt_references(dependencies, reservations);
587 receipts.retain(|reference, _| required.contains(reference));
588}
589
590pub(crate) fn load_verified_archived_evidence_segment(
591 receipt: &ArchivedEvidenceReceipt,
592 provider: &dyn ArchiveProvider,
593) -> Result<EvidenceJournalSegment, CanwuError> {
594 let segment = provider
595 .load_evidence_segment(&receipt.locator.segment_id)?
596 .ok_or_else(|| {
597 CanwuError::new(
598 ErrorCode::EvidenceContentUnavailable,
599 "the archive provider did not return the required evidence segment",
600 )
601 })?;
602 let receipts = verify_archived_segment(&segment).map_err(|error| {
603 CanwuError::new(
604 ErrorCode::InvalidArchive,
605 format!(
606 "archive provider returned invalid content: {}",
607 error.message
608 ),
609 )
610 })?;
611 if !receipts.iter().any(|candidate| candidate == receipt) {
612 return Err(archive_error(
613 "archive provider segment does not reproduce the committed evidence receipt",
614 ));
615 }
616 Ok(segment)
617}
618
619fn add_cause_dependency(
620 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
621 cause: &CauseRef,
622) {
623 let reference = match cause {
624 CauseRef::Event(id) => Some(EvidenceRef::Event(*id)),
625 CauseRef::Command(id) => Some(EvidenceRef::Command(*id)),
626 CauseRef::Boundary(id) => Some(EvidenceRef::Boundary(*id)),
627 CauseRef::System(_) => None,
628 };
629 if let Some(reference) = reference {
630 promote_dependency(dependencies, reference, EvidenceRequirement::IdentityOnly);
631 }
632}
633
634fn add_command_outcome_dependencies(
635 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
636 outcome: &CommandOutcome,
637) {
638 let (attempt_id, command_id, emitted_events) = match outcome {
639 CommandOutcome::Accepted { receipt } => (
640 receipt.attempt_id,
641 Some(receipt.command_id),
642 receipt.emitted_events.as_slice(),
643 ),
644 CommandOutcome::Rejected { rejection } => (rejection.attempt_id, None, &[][..]),
645 };
646 if let Some(id) = attempt_id {
647 promote_dependency(
648 dependencies,
649 EvidenceRef::CommandAttempt(id),
650 EvidenceRequirement::IdentityOnly,
651 );
652 }
653 if let Some(id) = command_id {
654 promote_dependency(
655 dependencies,
656 EvidenceRef::Command(id),
657 EvidenceRequirement::IdentityOnly,
658 );
659 }
660 for id in emitted_events {
661 promote_dependency(
662 dependencies,
663 EvidenceRef::Event(*id),
664 EvidenceRequirement::IdentityOnly,
665 );
666 }
667}
668
669fn archive_leaf(entry: &EvidenceIndexEntry) -> Result<[u8; 32], CanwuError> {
670 let hash = super::canonical_hash(
671 "canwu.evidence.index.leaf.v1",
672 &EvidenceIndexLeafMaterial {
673 format_version: 1,
674 reference: &entry.reference,
675 item: &entry.item,
676 item_commitment: &entry.item_commitment,
677 },
678 )?;
679 decode_hash(&hash, "evidence-index leaf")
680}
681
682fn archive_merkle(
683 entries: &[EvidenceIndexEntry],
684) -> Result<(String, Vec<Vec<String>>), CanwuError> {
685 if entries.is_empty() {
686 return Ok((archive_empty_root(), Vec::new()));
687 }
688 let mut level: Vec<[u8; 32]> = entries.iter().map(archive_leaf).collect::<Result<_, _>>()?;
689 let mut positions: Vec<usize> = (0..entries.len()).collect();
690 let mut proofs = vec![Vec::new(); entries.len()];
691 while level.len() > 1 {
692 for (leaf, position) in positions.iter().copied().enumerate() {
693 let sibling = if position % 2 == 0 {
694 (position + 1).min(level.len() - 1)
695 } else {
696 position - 1
697 };
698 proofs[leaf].push(
699 blake3::Hash::from_bytes(level[sibling])
700 .to_hex()
701 .to_string(),
702 );
703 }
704 let mut next = Vec::with_capacity(level.len().div_ceil(2));
705 for pair in level.chunks(2) {
706 next.push(archive_node(pair[0], *pair.get(1).unwrap_or(&pair[0])));
707 }
708 level = next;
709 for position in &mut positions {
710 *position /= 2;
711 }
712 }
713 Ok((
714 blake3::Hash::from_bytes(level[0]).to_hex().to_string(),
715 proofs,
716 ))
717}
718
719fn item_commitment<T: Serialize>(
720 journal: EvidenceJournalKind,
721 item: &T,
722) -> Result<String, CanwuError> {
723 let domain = match journal {
724 EvidenceJournalKind::Event => "canwu.evidence.item.event.v1",
725 EvidenceJournalKind::Command => "canwu.evidence.item.command.v1",
726 EvidenceJournalKind::CommandAttempt => "canwu.evidence.item.command_attempt.v1",
727 EvidenceJournalKind::Ingress => "canwu.evidence.item.ingress.v1",
728 EvidenceJournalKind::Boundary => "canwu.evidence.item.boundary.v1",
729 EvidenceJournalKind::RandomDraw => "canwu.evidence.item.random_draw.v1",
730 };
731 super::canonical_hash(domain, item)
732}
733
734pub(crate) fn evidence_archive_index(
735 segment: &EvidenceJournalSegment,
736) -> Result<(EvidenceArchiveIndex, Vec<ArchivedEvidenceReceipt>), CanwuError> {
737 let roots = EvidenceJournalRoots {
738 events: super::canonical_hash("canwu.evidence.journal.events.v1", &segment.events)?,
739 commands: super::canonical_hash("canwu.evidence.journal.commands.v1", &segment.commands)?,
740 command_attempts: super::canonical_hash(
741 "canwu.evidence.journal.command_attempts.v1",
742 &segment.command_attempts,
743 )?,
744 ingress: super::canonical_hash("canwu.evidence.journal.ingress.v1", &segment.ingress)?,
745 boundaries: super::canonical_hash(
746 "canwu.evidence.journal.boundaries.v1",
747 &segment.boundaries,
748 )?,
749 random_draws: super::canonical_hash(
750 "canwu.evidence.journal.random_draws.v1",
751 &segment.random_draws,
752 )?,
753 };
754 let mut entries = Vec::new();
755 let mut add = |reference: EvidenceRef, journal, absolute_index, nested, commitment: String| {
756 entries.push(EvidenceIndexEntry {
757 reference,
758 item: EvidenceItemLocator {
759 journal,
760 absolute_index,
761 nested,
762 },
763 item_commitment: commitment,
764 });
765 };
766 for (offset, event) in segment.events.iter().enumerate() {
767 add(
768 EvidenceRef::Event(event.id),
769 EvidenceJournalKind::Event,
770 segment.start.event_count + offset as u64 + 1,
771 EvidenceNestedLocator::None,
772 item_commitment(EvidenceJournalKind::Event, event)?,
773 );
774 }
775 for (offset, command) in segment.commands.iter().enumerate() {
776 add(
777 EvidenceRef::Command(command.id),
778 EvidenceJournalKind::Command,
779 segment.start.command_count + offset as u64 + 1,
780 EvidenceNestedLocator::None,
781 item_commitment(EvidenceJournalKind::Command, command)?,
782 );
783 }
784 for (offset, attempt) in segment.command_attempts.iter().enumerate() {
785 add(
786 EvidenceRef::CommandAttempt(attempt.id),
787 EvidenceJournalKind::CommandAttempt,
788 segment.start.command_attempt_count + offset as u64 + 1,
789 EvidenceNestedLocator::None,
790 item_commitment(EvidenceJournalKind::CommandAttempt, attempt)?,
791 );
792 }
793 for (offset, ingress) in segment.ingress.iter().enumerate() {
794 add(
795 EvidenceRef::Ingress(ingress.id),
796 EvidenceJournalKind::Ingress,
797 segment.start.ingress_count + offset as u64 + 1,
798 EvidenceNestedLocator::None,
799 item_commitment(EvidenceJournalKind::Ingress, ingress)?,
800 );
801 }
802 for (offset, boundary) in segment.boundaries.iter().enumerate() {
803 let absolute_index = segment.start.boundary_count + offset as u64 + 1;
804 let commitment = item_commitment(EvidenceJournalKind::Boundary, boundary)?;
805 add(
806 EvidenceRef::Boundary(boundary.id),
807 EvidenceJournalKind::Boundary,
808 absolute_index,
809 EvidenceNestedLocator::None,
810 commitment.clone(),
811 );
812 for (change_index, change) in boundary.record_changes.iter().enumerate() {
813 add(
814 EvidenceRef::DomainRecordVersion(DomainRecordVersionRef {
815 record: change.current.reference.clone(),
816 version: change.current.version,
817 established_by: DomainRecordVersionSource::BoundaryChange {
818 boundary: boundary.id,
819 change_index: change_index as u64,
820 },
821 }),
822 EvidenceJournalKind::Boundary,
823 absolute_index,
824 EvidenceNestedLocator::BoundaryRecordChange {
825 change_index: change_index as u64,
826 },
827 commitment.clone(),
828 );
829 }
830 }
831 for (offset, draw) in segment.random_draws.iter().enumerate() {
832 add(
833 EvidenceRef::RandomDraw(draw.id),
834 EvidenceJournalKind::RandomDraw,
835 segment.start.random_draw_count + offset as u64 + 1,
836 EvidenceNestedLocator::None,
837 item_commitment(EvidenceJournalKind::RandomDraw, draw)?,
838 );
839 }
840 entries.sort();
841 if entries
842 .windows(2)
843 .any(|window| window[0].reference == window[1].reference)
844 {
845 return Err(archive_error(
846 "evidence archive contains duplicate references",
847 ));
848 }
849 let (evidence_index_root, proofs) = archive_merkle(&entries)?;
850 let entry_count = u64::try_from(entries.len())
851 .map_err(|_| archive_error("evidence-index entry count exceeds u64"))?;
852 let segment_id = super::canonical_hash(
853 "canwu.evidence.segment.v2",
854 &ArchivedSegmentHeaderMaterial {
855 start: segment.start,
856 end: segment.end,
857 journal_roots: &roots,
858 evidence_index_root: &evidence_index_root,
859 evidence_index_entry_count: entry_count,
860 },
861 )?;
862 let header = ArchivedSegmentHeader {
863 segment_id: segment_id.clone(),
864 start: segment.start,
865 end: segment.end,
866 journal_roots: roots,
867 evidence_index_root,
868 evidence_index_entry_count: entry_count,
869 };
870 let receipts = entries
871 .iter()
872 .zip(proofs)
873 .enumerate()
874 .map(|(index, (entry, merkle_path))| ArchivedEvidenceReceipt {
875 evidence: entry.reference.clone(),
876 locator: ArchivedEvidenceLocator {
877 segment_id: segment_id.clone(),
878 item: entry.item.clone(),
879 },
880 evidence_index_leaf: index as u64,
881 item_commitment: entry.item_commitment.clone(),
882 merkle_path,
883 })
884 .collect();
885 Ok((EvidenceArchiveIndex { header, entries }, receipts))
886}
887
888fn verify_archive_receipt(
889 receipt: &ArchivedEvidenceReceipt,
890 header: &ArchivedSegmentHeader,
891) -> Result<(), CanwuError> {
892 if receipt.locator.segment_id != header.segment_id
893 || receipt.evidence_index_leaf >= header.evidence_index_entry_count
894 {
895 return Err(archive_error(
896 "archived evidence receipt does not belong to its segment header",
897 ));
898 }
899 let entry = EvidenceIndexEntry {
900 reference: receipt.evidence.clone(),
901 item: receipt.locator.item.clone(),
902 item_commitment: receipt.item_commitment.clone(),
903 };
904 let mut hash = archive_leaf(&entry)?;
905 let mut position = receipt.evidence_index_leaf;
906 let mut width = header.evidence_index_entry_count;
907 let mut path_at = 0usize;
908 while width > 1 {
909 let sibling = receipt
910 .merkle_path
911 .get(path_at)
912 .ok_or_else(|| archive_error("archived evidence receipt Merkle path is too short"))?;
913 let sibling = decode_hash(sibling, "receipt Merkle sibling")?;
914 hash = if position.is_multiple_of(2) {
915 archive_node(hash, sibling)
916 } else {
917 archive_node(sibling, hash)
918 };
919 position /= 2;
920 width = width.div_ceil(2);
921 path_at += 1;
922 }
923 if path_at != receipt.merkle_path.len()
924 || blake3::Hash::from_bytes(hash).to_hex().as_str() != header.evidence_index_root
925 {
926 return Err(archive_error(
927 "archived evidence receipt Merkle proof is invalid",
928 ));
929 }
930 let (expected_journal, expected_nested) = match &receipt.evidence {
931 EvidenceRef::Event(_) => (EvidenceJournalKind::Event, EvidenceNestedLocator::None),
932 EvidenceRef::Command(_) => (EvidenceJournalKind::Command, EvidenceNestedLocator::None),
933 EvidenceRef::CommandAttempt(_) => (
934 EvidenceJournalKind::CommandAttempt,
935 EvidenceNestedLocator::None,
936 ),
937 EvidenceRef::Ingress(_) => (EvidenceJournalKind::Ingress, EvidenceNestedLocator::None),
938 EvidenceRef::Boundary(_) => (EvidenceJournalKind::Boundary, EvidenceNestedLocator::None),
939 EvidenceRef::RandomDraw(_) => {
940 (EvidenceJournalKind::RandomDraw, EvidenceNestedLocator::None)
941 }
942 EvidenceRef::DomainRecordVersion(version) => match version.established_by {
943 DomainRecordVersionSource::BoundaryChange { change_index, .. } => (
944 EvidenceJournalKind::Boundary,
945 EvidenceNestedLocator::BoundaryRecordChange { change_index },
946 ),
947 DomainRecordVersionSource::InitialScenario => {
948 return Err(archive_error(
949 "initial-scenario record versions cannot be archived",
950 ));
951 }
952 },
953 };
954 if receipt.locator.item.journal != expected_journal
955 || receipt.locator.item.nested != expected_nested
956 {
957 return Err(archive_error(
958 "archived evidence receipt uses an illegal typed locator",
959 ));
960 }
961 let (start, end) = match expected_journal {
962 EvidenceJournalKind::Event => (header.start.event_count, header.end.event_count),
963 EvidenceJournalKind::Command => (header.start.command_count, header.end.command_count),
964 EvidenceJournalKind::CommandAttempt => (
965 header.start.command_attempt_count,
966 header.end.command_attempt_count,
967 ),
968 EvidenceJournalKind::Ingress => (header.start.ingress_count, header.end.ingress_count),
969 EvidenceJournalKind::Boundary => (header.start.boundary_count, header.end.boundary_count),
970 EvidenceJournalKind::RandomDraw => {
971 (header.start.random_draw_count, header.end.random_draw_count)
972 }
973 };
974 if receipt.locator.item.absolute_index <= start || receipt.locator.item.absolute_index > end {
975 return Err(archive_error(
976 "archived evidence locator lies outside its segment cursor range",
977 ));
978 }
979 Ok(())
980}
981
982pub(crate) fn verify_archived_segment(
983 segment: &EvidenceJournalSegment,
984) -> Result<Vec<ArchivedEvidenceReceipt>, CanwuError> {
985 let stored = segment
986 .archive
987 .as_ref()
988 .ok_or_else(|| archive_error("archived segment is missing its evidence index"))?;
989 let mut plain = segment.clone();
990 plain.archive = None;
991 let (rebuilt, receipts) = evidence_archive_index(&plain)?;
992 if &rebuilt != stored {
993 return Err(archive_error(
994 "archived segment header or evidence index does not match its journal content",
995 ));
996 }
997 for receipt in &receipts {
998 verify_archive_receipt(receipt, &stored.header)?;
999 }
1000 Ok(receipts)
1001}
1002
1003#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
1004pub struct EvidenceCursor {
1006 pub event_count: u64,
1007 pub command_count: u64,
1008 pub command_attempt_count: u64,
1009 pub ingress_count: u64,
1010 pub boundary_count: u64,
1011 pub random_draw_count: u64,
1012}
1013
1014impl EvidenceCursor {
1015 fn from_evidence(evidence: &RuntimeEvidence) -> Result<Self, CanwuError> {
1016 let count = |len: usize, label: &str| {
1017 u64::try_from(len).map_err(|_| {
1018 CanwuError::new(
1019 ErrorCode::IdentifierExhausted,
1020 format!("{label} journal length exceeds the persistent cursor space"),
1021 )
1022 })
1023 };
1024 Ok(Self {
1025 event_count: evidence
1026 .archived
1027 .event_count
1028 .checked_add(count(evidence.events.len(), "event")?)
1029 .ok_or_else(|| invalid_snapshot_error("event journal cursor is exhausted"))?,
1030 command_count: evidence
1031 .archived
1032 .command_count
1033 .checked_add(count(evidence.commands.len(), "command")?)
1034 .ok_or_else(|| invalid_snapshot_error("command journal cursor is exhausted"))?,
1035 command_attempt_count: evidence
1036 .archived
1037 .command_attempt_count
1038 .checked_add(count(evidence.command_attempts.len(), "command-attempt")?)
1039 .ok_or_else(|| {
1040 invalid_snapshot_error("command-attempt journal cursor is exhausted")
1041 })?,
1042 ingress_count: evidence
1043 .archived
1044 .ingress_count
1045 .checked_add(count(evidence.ingress.len(), "ingress")?)
1046 .ok_or_else(|| invalid_snapshot_error("ingress journal cursor is exhausted"))?,
1047 boundary_count: evidence
1048 .archived
1049 .boundary_count
1050 .checked_add(count(evidence.boundaries.len(), "boundary")?)
1051 .ok_or_else(|| invalid_snapshot_error("boundary journal cursor is exhausted"))?,
1052 random_draw_count: evidence
1053 .archived
1054 .random_draw_count
1055 .checked_add(count(evidence.random_draws.len(), "random-draw")?)
1056 .ok_or_else(|| invalid_snapshot_error("random-draw journal cursor is exhausted"))?,
1057 })
1058 }
1059
1060 pub(super) fn checked_advance(
1061 self,
1062 segment: &EvidenceJournalSegment,
1063 ) -> Result<Self, CanwuError> {
1064 let advance = |value: u64, len: usize, label: &str| {
1065 value
1066 .checked_add(u64::try_from(len).map_err(|_| {
1067 invalid_snapshot_error(format!(
1068 "{label} journal segment exceeds the persistent cursor space"
1069 ))
1070 })?)
1071 .ok_or_else(|| {
1072 invalid_snapshot_error(format!(
1073 "{label} journal cursor exceeds the persistent cursor space"
1074 ))
1075 })
1076 };
1077 Ok(Self {
1078 event_count: advance(self.event_count, segment.events.len(), "event")?,
1079 command_count: advance(self.command_count, segment.commands.len(), "command")?,
1080 command_attempt_count: advance(
1081 self.command_attempt_count,
1082 segment.command_attempts.len(),
1083 "command-attempt",
1084 )?,
1085 ingress_count: advance(self.ingress_count, segment.ingress.len(), "ingress")?,
1086 boundary_count: advance(self.boundary_count, segment.boundaries.len(), "boundary")?,
1087 random_draw_count: advance(
1088 self.random_draw_count,
1089 segment.random_draws.len(),
1090 "random-draw",
1091 )?,
1092 })
1093 }
1094}
1095
1096#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1097pub struct SimulationCheckpoint {
1103 pub format_version: u32,
1104 pub journal_end: EvidenceCursor,
1105 pub state: SimulationSnapshot,
1106 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1107 pub archived_segment_headers: Vec<ArchivedSegmentHeader>,
1108 #[serde(default, skip_serializing_if = "Option::is_none")]
1109 pub archived_segment_manifest_root: Option<String>,
1110 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1111 pub archived_evidence_receipts: Vec<ArchivedEvidenceReceipt>,
1112 #[serde(default, skip_serializing_if = "Option::is_none")]
1113 pub archived_receipt_root: Option<String>,
1114 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1115 pub evidence_dependencies: Vec<EvidenceDependency>,
1116 #[serde(default, skip_serializing_if = "Option::is_none")]
1117 pub evidence_dependency_root: Option<String>,
1118 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1119 pub keyed_draw_reservations: Vec<KeyedDrawReservation>,
1120 #[serde(default, skip_serializing_if = "Option::is_none")]
1121 pub keyed_reservation_root: Option<String>,
1122}
1123
1124#[derive(Serialize)]
1125struct CompactCheckpointHashMaterial<'a> {
1126 state_checkpoint_hash: &'a str,
1127 journal_end: EvidenceCursor,
1128 archived_segment_manifest_root: Option<&'a str>,
1129 archived_receipt_root: Option<&'a str>,
1130 evidence_dependency_root: Option<&'a str>,
1131 keyed_reservation_root: Option<&'a str>,
1132}
1133
1134fn validate_compact_continuation(checkpoint: &SimulationCheckpoint) -> Result<(), CanwuError> {
1135 if checkpoint
1136 .archived_segment_headers
1137 .windows(2)
1138 .any(|headers| headers[0].end != headers[1].start)
1139 || checkpoint
1140 .archived_segment_headers
1141 .iter()
1142 .map(|header| &header.segment_id)
1143 .collect::<BTreeSet<_>>()
1144 .len()
1145 != checkpoint.archived_segment_headers.len()
1146 {
1147 return Err(invalid_snapshot_error(
1148 "compact archived-segment manifest is duplicated or noncontiguous",
1149 ));
1150 }
1151 if checkpoint
1152 .archived_evidence_receipts
1153 .windows(2)
1154 .any(|receipts| receipts[0].evidence >= receipts[1].evidence)
1155 {
1156 return Err(invalid_snapshot_error(
1157 "compact archived receipts must be sorted by unique evidence reference",
1158 ));
1159 }
1160 if checkpoint
1161 .evidence_dependencies
1162 .windows(2)
1163 .any(|dependencies| dependencies[0].reference >= dependencies[1].reference)
1164 {
1165 return Err(invalid_snapshot_error(
1166 "compact evidence dependencies must be sorted by unique evidence reference",
1167 ));
1168 }
1169 if checkpoint
1170 .keyed_draw_reservations
1171 .windows(2)
1172 .any(|reservations| {
1173 (&reservations[0].stream, &reservations[0].address)
1174 >= (&reservations[1].stream, &reservations[1].address)
1175 })
1176 {
1177 return Err(invalid_snapshot_error(
1178 "compact keyed reservations must be sorted by unique operation address",
1179 ));
1180 }
1181 validate_skipped_commitment_root(
1182 checkpoint.archived_segment_manifest_root.as_deref(),
1183 ARCHIVED_SEGMENT_MANIFEST_DOMAIN,
1184 &checkpoint.archived_segment_headers,
1185 "archived-segment manifest root",
1186 )?;
1187 validate_skipped_commitment_root(
1188 checkpoint.archived_receipt_root.as_deref(),
1189 ARCHIVED_RECEIPT_DOMAIN,
1190 &checkpoint.archived_evidence_receipts,
1191 "archived-receipt root",
1192 )?;
1193 validate_skipped_commitment_root(
1194 checkpoint.evidence_dependency_root.as_deref(),
1195 EVIDENCE_DEPENDENCY_DOMAIN,
1196 &checkpoint.evidence_dependencies,
1197 "evidence-dependency root",
1198 )?;
1199 validate_skipped_commitment_root(
1200 checkpoint.keyed_reservation_root.as_deref(),
1201 KEYED_RESERVATION_DOMAIN,
1202 &checkpoint.keyed_draw_reservations,
1203 "keyed-reservation root",
1204 )
1205}
1206
1207impl SimulationCheckpoint {
1208 pub fn reachable_archive_segment_ids(
1213 retained_checkpoints: &[Self],
1214 ) -> Result<BTreeSet<String>, CanwuError> {
1215 let mut reachable = BTreeSet::new();
1216 for checkpoint in retained_checkpoints {
1217 validate_compact_continuation(checkpoint)?;
1218 reachable.extend(
1219 checkpoint
1220 .archived_segment_headers
1221 .iter()
1222 .map(|header| header.segment_id.clone()),
1223 );
1224 }
1225 Ok(reachable)
1226 }
1227
1228 pub fn orphaned_archive_segment_ids(
1232 retained_checkpoints: &[Self],
1233 stored_segment_ids: &[String],
1234 ) -> Result<Vec<String>, CanwuError> {
1235 let reachable = Self::reachable_archive_segment_ids(retained_checkpoints)?;
1236 Ok(stored_segment_ids
1237 .iter()
1238 .filter(|segment_id| !reachable.contains(segment_id.as_str()))
1239 .cloned()
1240 .collect::<BTreeSet<_>>()
1241 .into_iter()
1242 .collect())
1243 }
1244}
1245
1246fn compact_checkpoint_hash(checkpoint: &SimulationCheckpoint) -> Result<String, CanwuError> {
1247 validate_compact_continuation(checkpoint)?;
1248 super::canonical_hash(
1249 "canwu.compact-checkpoint.v1",
1250 &CompactCheckpointHashMaterial {
1251 state_checkpoint_hash: &checkpoint.state.checkpoint_hash,
1252 journal_end: checkpoint.journal_end,
1253 archived_segment_manifest_root: checkpoint.archived_segment_manifest_root.as_deref(),
1254 archived_receipt_root: checkpoint.archived_receipt_root.as_deref(),
1255 evidence_dependency_root: checkpoint.evidence_dependency_root.as_deref(),
1256 keyed_reservation_root: checkpoint.keyed_reservation_root.as_deref(),
1257 },
1258 )
1259}
1260
1261#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1262pub struct EvidenceJournalSegment {
1264 pub format_version: u32,
1265 pub start: EvidenceCursor,
1266 pub end: EvidenceCursor,
1267 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1268 pub events: Vec<SimEvent>,
1269 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1270 pub commands: Vec<CommandRecord>,
1271 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1272 pub command_attempts: Vec<CommandAttemptRecord>,
1273 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1274 pub ingress: Vec<IngressRecord>,
1275 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1276 pub boundaries: Vec<BoundaryRecord>,
1277 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1278 pub random_draws: Vec<RandomDrawRecord>,
1279 #[serde(default, skip_serializing_if = "Option::is_none")]
1280 pub archive: Option<EvidenceArchiveIndex>,
1281}
1282
1283#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1284pub struct CheckpointJournal {
1286 pub checkpoint: SimulationCheckpoint,
1287 pub segments: Vec<EvidenceJournalSegment>,
1288}
1289
1290pub struct CompactedSimulation {
1298 simulation: Simulation,
1299 committed_seal_tokens: BTreeSet<String>,
1300}
1301
1302impl CompactedSimulation {
1303 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
1305 self.simulation.evidence_cursor()
1306 }
1307
1308 pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
1310 self.simulation.checkpoint()
1311 }
1312
1313 #[must_use]
1315 pub fn archived_evidence_receipt(
1316 &self,
1317 reference: &EvidenceRef,
1318 ) -> Option<&ArchivedEvidenceReceipt> {
1319 self.simulation
1320 .state
1321 .evidence
1322 .archived_evidence_receipts
1323 .get(reference)
1324 }
1325
1326 pub fn load_archived_evidence_segment(
1329 &self,
1330 reference: &EvidenceRef,
1331 provider: &dyn ArchiveProvider,
1332 ) -> Result<EvidenceJournalSegment, CanwuError> {
1333 let receipt = self.archived_evidence_receipt(reference).ok_or_else(|| {
1334 CanwuError::new(
1335 ErrorCode::EvidenceUnavailable,
1336 "no committed archive receipt exists for the requested evidence",
1337 )
1338 })?;
1339 load_verified_archived_evidence_segment(receipt, provider)
1340 }
1341
1342 pub fn seal_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
1349 if self
1350 .simulation
1351 .evidence_dependencies()?
1352 .iter()
1353 .any(|dependency| dependency.requirement == EvidenceRequirement::PayloadRequired)
1354 {
1355 return Err(CanwuError::new(
1356 ErrorCode::ArchiveNotReady,
1357 "payload-required continuations must use prepare/store/commit sealing",
1358 ));
1359 }
1360 self.simulation.seal_retained_evidence()
1361 }
1362
1363 pub fn prepare_evidence_seal(&self) -> Result<Option<PreparedEvidenceSeal>, CanwuError> {
1369 let source_state_hash = self.simulation.authoritative_state_hash()?;
1370 let source_checkpoint_hash = compact_checkpoint_hash(&self.simulation.checkpoint()?)?;
1371 let source_end = self.simulation.evidence_cursor()?;
1372 let mut candidate = self.simulation.fork();
1373 let Some(segment) = candidate.seal_retained_evidence()? else {
1374 return Ok(None);
1375 };
1376 let segment_id = segment
1377 .archive
1378 .as_ref()
1379 .ok_or_else(|| archive_error("prepared segment has no archive index"))?
1380 .header
1381 .segment_id
1382 .clone();
1383 let target_checkpoint_hash = compact_checkpoint_hash(&candidate.checkpoint()?)?;
1384 let token_hash = super::canonical_hash(
1385 "canwu.evidence.seal-token.v1",
1386 &(
1387 &source_state_hash,
1388 &source_checkpoint_hash,
1389 source_end,
1390 &segment_id,
1391 &target_checkpoint_hash,
1392 ),
1393 )?;
1394 Ok(Some(PreparedEvidenceSeal {
1395 token: EvidenceSealToken {
1396 source_state_hash,
1397 source_checkpoint_hash,
1398 source_end,
1399 segment_id,
1400 target_checkpoint_hash,
1401 token_hash,
1402 },
1403 segment,
1404 }))
1405 }
1406
1407 pub fn commit_evidence_seal(
1410 &mut self,
1411 token: &EvidenceSealToken,
1412 provider: &dyn ArchiveProvider,
1413 ) -> Result<(), CanwuError> {
1414 if self.committed_seal_tokens.contains(&token.token_hash) {
1415 return Ok(());
1416 }
1417 let expected_token_hash = super::canonical_hash(
1418 "canwu.evidence.seal-token.v1",
1419 &(
1420 &token.source_state_hash,
1421 &token.source_checkpoint_hash,
1422 token.source_end,
1423 &token.segment_id,
1424 &token.target_checkpoint_hash,
1425 ),
1426 )?;
1427 if expected_token_hash != token.token_hash {
1428 return Err(archive_error("evidence seal token hash is invalid"));
1429 }
1430 if self.simulation.authoritative_state_hash()? != token.source_state_hash
1431 || compact_checkpoint_hash(&self.simulation.checkpoint()?)?
1432 != token.source_checkpoint_hash
1433 || self.simulation.evidence_cursor()? != token.source_end
1434 {
1435 return Err(CanwuError::new(
1436 ErrorCode::StaleSealToken,
1437 "evidence seal token no longer names the live source cut",
1438 ));
1439 }
1440 let stored = provider
1441 .load_evidence_segment(&token.segment_id)?
1442 .ok_or_else(|| {
1443 CanwuError::new(
1444 ErrorCode::ArchiveNotReady,
1445 "prepared evidence segment is not available from the archive provider",
1446 )
1447 })?;
1448 let archive = stored
1449 .archive
1450 .as_ref()
1451 .ok_or_else(|| archive_error("stored evidence segment has no archive index"))?;
1452 if archive.header.segment_id != token.segment_id {
1453 return Err(archive_error(
1454 "stored evidence segment ID does not match the seal token",
1455 ));
1456 }
1457 verify_archived_segment(&stored)?;
1458
1459 let before = self.simulation.fork();
1460 let result = (|| {
1461 let actual = self.simulation.seal_retained_evidence()?.ok_or_else(|| {
1462 CanwuError::new(
1463 ErrorCode::StaleSealToken,
1464 "the prepared evidence tail is no longer retained",
1465 )
1466 })?;
1467 if actual != stored {
1468 return Err(archive_error(
1469 "stored evidence segment differs from the prepared live cut",
1470 ));
1471 }
1472 self.simulation
1473 .validate_payload_required_archive(provider)?;
1474 if compact_checkpoint_hash(&self.simulation.checkpoint()?)?
1475 != token.target_checkpoint_hash
1476 {
1477 return Err(archive_error(
1478 "committed checkpoint hash differs from the seal token",
1479 ));
1480 }
1481 Ok(())
1482 })();
1483 if let Err(error) = result {
1484 self.simulation = before;
1485 return Err(error);
1486 }
1487 self.committed_seal_tokens.insert(token.token_hash.clone());
1488 Ok(())
1489 }
1490
1491 #[must_use]
1492 pub const fn time(&self) -> SimTime {
1493 self.simulation.time()
1494 }
1495
1496 #[must_use]
1497 pub const fn revision(&self) -> u64 {
1498 self.simulation.revision()
1499 }
1500
1501 #[must_use]
1502 pub fn checkpoint_hash(&self) -> &str {
1503 self.simulation.checkpoint_hash()
1504 }
1505
1506 #[must_use]
1507 pub fn boundary_head_hash(&self) -> Option<&str> {
1508 self.simulation.boundary_head_hash()
1509 }
1510
1511 #[must_use]
1512 pub fn world(&self) -> WorldSnapshot {
1513 self.simulation.world()
1514 }
1515
1516 #[must_use]
1517 pub fn knowledge(&self) -> &KnowledgeSnapshot {
1518 self.simulation.knowledge()
1519 }
1520
1521 #[must_use]
1522 pub fn domain_record(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
1523 self.simulation.domain_record(reference)
1524 }
1525
1526 #[must_use]
1527 pub const fn decision_state(&self) -> &super::DecisionState {
1528 self.simulation.decision_state()
1529 }
1530
1531 #[must_use]
1532 pub fn decision_ticket(&self, id: super::DecisionTicketId) -> Option<&super::DecisionTicket> {
1533 self.simulation.decision_ticket(id)
1534 }
1535
1536 #[must_use]
1537 pub fn decision_traces(&self) -> &[super::DecisionTrace] {
1538 self.simulation.decision_traces()
1539 }
1540
1541 #[must_use]
1542 pub fn decision_attempts(&self) -> &[super::DecisionAttemptRecord] {
1543 self.simulation.decision_attempts()
1544 }
1545
1546 #[must_use]
1547 pub fn typed_domain_record<T: DomainRecordType>(
1548 &self,
1549 reference: &TypedDomainRecordRef<T>,
1550 ) -> Option<&DomainRecord> {
1551 self.simulation.typed_domain_record(reference)
1552 }
1553
1554 pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
1555 self.simulation.submit(envelope)
1556 }
1557
1558 pub fn process_command(
1559 &mut self,
1560 request: CommandRequest,
1561 ) -> Result<CommandOutcome, CanwuError> {
1562 self.simulation.process_command(request)
1563 }
1564
1565 pub fn enqueue_command(
1566 &mut self,
1567 due_at: SimTime,
1568 priority: i32,
1569 request: CommandRequest,
1570 ) -> Result<IngressReceipt, CanwuError> {
1571 self.simulation.enqueue_command(due_at, priority, request)
1572 }
1573
1574 pub fn enqueue_plugin_ingress(
1575 &mut self,
1576 request: PluginIngressRequest,
1577 ) -> Result<IngressReceipt, CanwuError> {
1578 self.simulation.enqueue_plugin_ingress(request)
1579 }
1580
1581 pub fn prepare_decision(
1582 &self,
1583 decision_request_id: super::DecisionRequestId,
1584 command_request_id: Option<super::CommandRequestId>,
1585 ticket_id: super::DecisionTicketId,
1586 policy: &dyn super::DecisionPolicy,
1587 ) -> Result<super::DecisionEvaluation, CanwuError> {
1588 self.simulation
1589 .prepare_decision(decision_request_id, command_request_id, ticket_id, policy)
1590 }
1591
1592 pub fn prepare_decision_at(
1593 &self,
1594 due_at: super::SimTime,
1595 decision_request_id: super::DecisionRequestId,
1596 command_request_id: Option<super::CommandRequestId>,
1597 ticket_id: super::DecisionTicketId,
1598 policy: &dyn super::DecisionPolicy,
1599 ) -> Result<super::DecisionEvaluation, CanwuError> {
1600 self.simulation.prepare_decision_at(
1601 due_at,
1602 decision_request_id,
1603 command_request_id,
1604 ticket_id,
1605 policy,
1606 )
1607 }
1608
1609 pub fn enqueue_decision(
1610 &mut self,
1611 due_at: super::SimTime,
1612 priority: i32,
1613 request: super::DecisionIngressRequest,
1614 ) -> Result<IngressReceipt, CanwuError> {
1615 self.simulation.enqueue_decision(due_at, priority, request)
1616 }
1617
1618 pub fn drive_decision(
1619 &mut self,
1620 due_at: super::SimTime,
1621 priority: i32,
1622 decision_request_id: super::DecisionRequestId,
1623 command_request_id: Option<super::CommandRequestId>,
1624 ticket_id: super::DecisionTicketId,
1625 policy: &dyn super::DecisionPolicy,
1626 ) -> Result<super::DecisionEvaluation, CanwuError> {
1627 self.simulation.drive_decision(
1628 due_at,
1629 priority,
1630 decision_request_id,
1631 command_request_id,
1632 ticket_id,
1633 policy,
1634 )
1635 }
1636
1637 pub fn schedule_calendar_boundary(
1638 &mut self,
1639 due_at: SimTime,
1640 cadences: Vec<SystemCadence>,
1641 ) -> Result<IngressReceipt, CanwuError> {
1642 self.simulation.schedule_calendar_boundary(due_at, cadences)
1643 }
1644
1645 pub fn advance(&mut self, duration: SimDuration) -> Result<Vec<SimEvent>, CanwuError> {
1646 self.simulation.advance(duration)
1647 }
1648
1649 pub fn advance_canonical(
1650 &mut self,
1651 duration: SimDuration,
1652 ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
1653 self.simulation.advance_canonical(duration)
1654 }
1655
1656 pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
1657 self.simulation.step_canonical()
1658 }
1659
1660 pub fn settle_boundary(
1661 &mut self,
1662 request: BoundaryRequest,
1663 ) -> Result<BoundaryReceipt, CanwuError> {
1664 self.simulation.settle_boundary(request)
1665 }
1666
1667 pub fn snapshot_with_segments(
1670 &self,
1671 mut segments: Vec<EvidenceJournalSegment>,
1672 ) -> Result<SimulationSnapshot, CanwuError> {
1673 let tail = self
1674 .simulation
1675 .journal_segment_since(self.simulation.state.evidence.archived)?;
1676 if tail.start != tail.end {
1677 segments.push(tail);
1678 }
1679 let snapshot =
1680 Simulation::snapshot_from_checkpoint_and_journal(self.checkpoint()?, segments)?;
1681 Simulation::from_snapshot(snapshot.clone())?;
1682 Ok(snapshot)
1683 }
1684
1685 pub fn replay_journal_with_segments(
1687 &self,
1688 segments: Vec<EvidenceJournalSegment>,
1689 ) -> Result<ReplayJournal, CanwuError> {
1690 let snapshot = self.snapshot_with_segments(segments)?;
1691 let simulation = Simulation::from_snapshot(snapshot)?;
1692 Ok(simulation.replay_journal())
1693 }
1694
1695 pub fn from_checkpoint_and_journal(
1698 checkpoint: SimulationCheckpoint,
1699 segments: Vec<EvidenceJournalSegment>,
1700 ) -> Result<Self, CanwuError> {
1701 Simulation::from_checkpoint_and_journal(checkpoint, segments)?.into_compacted()
1702 }
1703
1704 pub fn from_checkpoint_and_journal_with_plugins(
1706 checkpoint: SimulationCheckpoint,
1707 segments: Vec<EvidenceJournalSegment>,
1708 plugins: &[&dyn SimulationPlugin],
1709 ) -> Result<Self, CanwuError> {
1710 let mut simulation = Simulation::from_checkpoint_and_journal(checkpoint, segments)?;
1711 for plugin in plugins {
1712 simulation.register_plugin(*plugin)?;
1713 }
1714 simulation.ensure_runtime_ready()?;
1715 simulation.into_compacted()
1716 }
1717}
1718
1719impl Simulation {
1720 pub fn into_compacted(self) -> Result<CompactedSimulation, CanwuError> {
1726 self.ensure_runtime_ready()?;
1727 Ok(CompactedSimulation {
1728 simulation: self,
1729 committed_seal_tokens: BTreeSet::new(),
1730 })
1731 }
1732
1733 fn evidence_dependencies(&self) -> Result<Vec<EvidenceDependency>, CanwuError> {
1734 let mut dependencies = BTreeMap::new();
1735 let identity = EvidenceRequirement::IdentityOnly;
1736
1737 for records in self.state.current.knowledge.records.values() {
1738 for record in records.values() {
1739 for reference in &record.origin.evidence {
1740 promote_dependency(&mut dependencies, reference.clone(), identity);
1741 }
1742 }
1743 }
1744
1745 for record in self.state.current.domain_records.values() {
1746 let schema = self
1747 .plugins
1748 .record_schemas
1749 .get(&record.reference.kind)
1750 .map(|(_, schema)| schema)
1751 .ok_or_else(|| {
1752 CanwuError::new(
1753 ErrorCode::ArchiveNotReady,
1754 format!(
1755 "current domain record {} has no registered schema",
1756 record.reference
1757 ),
1758 )
1759 })?;
1760 add_payload_required_continuation_dependencies(&mut dependencies, record, schema)?;
1761 let retained = self
1762 .state
1763 .evidence
1764 .boundaries
1765 .iter()
1766 .rev()
1767 .find_map(|boundary| {
1768 boundary
1769 .record_changes
1770 .iter()
1771 .enumerate()
1772 .find(|(_, change)| {
1773 change.current.reference == record.reference
1774 && change.current.version == record.version
1775 })
1776 .map(|(change_index, _)| {
1777 EvidenceRef::DomainRecordVersion(DomainRecordVersionRef {
1778 record: record.reference.clone(),
1779 version: record.version,
1780 established_by: DomainRecordVersionSource::BoundaryChange {
1781 boundary: boundary.id,
1782 change_index: change_index as u64,
1783 },
1784 })
1785 })
1786 });
1787 let archived = self
1788 .state
1789 .evidence
1790 .archived_evidence_receipts
1791 .keys()
1792 .filter(|reference| {
1793 matches!(
1794 reference,
1795 EvidenceRef::DomainRecordVersion(version)
1796 if version.record == record.reference && version.version == record.version
1797 )
1798 })
1799 .cloned()
1800 .collect::<Vec<_>>();
1801 if archived.len() > 1 {
1802 return Err(CanwuError::new(
1803 ErrorCode::ArchiveNotReady,
1804 format!(
1805 "current domain record {} has ambiguous archived version provenance",
1806 record.reference
1807 ),
1808 ));
1809 }
1810 if let Some(reference) = retained.or_else(|| archived.into_iter().next()) {
1811 promote_dependency(&mut dependencies, reference, identity);
1812 continue;
1813 }
1814 let initial = self.bound_initial_scenario().is_some_and(|scenario| {
1815 scenario.domain_records.iter().any(|initial| {
1816 initial.reference == record.reference && initial.version == record.version
1817 })
1818 });
1819 if !initial {
1820 return Err(CanwuError::new(
1821 ErrorCode::ArchiveNotReady,
1822 format!(
1823 "current domain record {} has no retained, archived, or scenario version provenance",
1824 record.reference
1825 ),
1826 ));
1827 }
1828 }
1829
1830 for reservation in &self.state.evidence.keyed_draw_reservations {
1831 promote_dependency(
1832 &mut dependencies,
1833 reservation.operation_evidence.clone(),
1834 identity,
1835 );
1836 promote_dependency(
1837 &mut dependencies,
1838 reservation.draw_receipt.evidence.clone(),
1839 identity,
1840 );
1841 }
1842 for draw in &self.state.evidence.random_draws {
1843 if matches!(draw.address, RandomDrawAddress::OperationV1(_)) {
1844 let reference = draw.operation_evidence.clone().ok_or_else(|| {
1845 CanwuError::new(
1846 ErrorCode::ArchiveNotReady,
1847 "operation-keyed retained draw is missing its operation evidence",
1848 )
1849 })?;
1850 promote_dependency(&mut dependencies, reference, identity);
1851 promote_dependency(
1852 &mut dependencies,
1853 EvidenceRef::RandomDraw(draw.id),
1854 identity,
1855 );
1856 }
1857 }
1858
1859 for archived in self.state.evidence.archived_command_requests.values() {
1860 add_command_outcome_dependencies(&mut dependencies, &archived.outcome);
1861 }
1862 for archived in self.state.evidence.archived_ingress_requests.values() {
1863 promote_dependency(
1864 &mut dependencies,
1865 EvidenceRef::Ingress(archived.receipt.ingress_id),
1866 identity,
1867 );
1868 }
1869 for archived in self.state.evidence.archived_decision_requests.values() {
1870 promote_dependency(
1871 &mut dependencies,
1872 EvidenceRef::Ingress(archived.receipt.ingress_id),
1873 identity,
1874 );
1875 }
1876 for attempt in &self.state.evidence.command_attempts {
1877 if attempt.request_id.is_none() {
1878 continue;
1879 }
1880 promote_dependency(
1881 &mut dependencies,
1882 EvidenceRef::CommandAttempt(attempt.id),
1883 identity,
1884 );
1885 if let CommandAttemptOutcome::Accepted { command_id } = attempt.outcome {
1886 promote_dependency(
1887 &mut dependencies,
1888 EvidenceRef::Command(command_id),
1889 identity,
1890 );
1891 if let Some(command) = self
1892 .state
1893 .evidence
1894 .commands
1895 .iter()
1896 .find(|command| command.id == command_id)
1897 {
1898 for event in &command.emitted_events {
1899 promote_dependency(&mut dependencies, EvidenceRef::Event(*event), identity);
1900 }
1901 }
1902 }
1903 }
1904 for ingress in &self.state.evidence.ingress {
1905 if matches!(
1906 ingress.payload,
1907 IngressPayload::Command { .. } | IngressPayload::Decision { .. }
1908 ) {
1909 promote_dependency(
1910 &mut dependencies,
1911 EvidenceRef::Ingress(ingress.id),
1912 identity,
1913 );
1914 }
1915 }
1916
1917 for action in self.state.scheduler.actions.values() {
1918 match action {
1919 ScheduledAction::ArmyArrival { order_event, .. }
1920 | ScheduledAction::PersonArrival { order_event, .. } => promote_dependency(
1921 &mut dependencies,
1922 EvidenceRef::Event(*order_event),
1923 identity,
1924 ),
1925 ScheduledAction::KnowledgeReport { dispatch_event, .. } => promote_dependency(
1926 &mut dependencies,
1927 EvidenceRef::Event(*dispatch_event),
1928 identity,
1929 ),
1930 ScheduledAction::PluginDirective { cause, .. } => {
1931 add_cause_dependency(&mut dependencies, cause);
1932 }
1933 }
1934 }
1935
1936 Ok(dependencies
1937 .into_iter()
1938 .map(|(reference, requirement)| EvidenceDependency {
1939 reference,
1940 requirement,
1941 })
1942 .collect())
1943 }
1944
1945 fn validate_payload_required_archive(
1946 &self,
1947 provider: &dyn ArchiveProvider,
1948 ) -> Result<(), CanwuError> {
1949 for dependency in self
1950 .evidence_dependencies()?
1951 .into_iter()
1952 .filter(|dependency| dependency.requirement == EvidenceRequirement::PayloadRequired)
1953 {
1954 let receipt = self
1955 .state
1956 .evidence
1957 .archived_evidence_receipts
1958 .get(&dependency.reference)
1959 .ok_or_else(|| {
1960 CanwuError::new(
1961 ErrorCode::ArchiveNotReady,
1962 "payload-required evidence has no committed archive receipt",
1963 )
1964 })?;
1965 load_verified_archived_evidence_segment(receipt, provider)?;
1966 }
1967 Ok(())
1968 }
1969
1970 fn ensure_retained_evidence_is_sealable(&self) -> Result<(), CanwuError> {
1971 if !self.state.scheduler.pending_ingress.is_empty() {
1972 return Err(CanwuError::new(
1973 ErrorCode::ArchiveNotReady,
1974 "live evidence can be sealed only when the canonical ingress queue is empty",
1975 ));
1976 }
1977 let admitted_attempts: std::collections::BTreeSet<_> = self
1978 .state
1979 .evidence
1980 .boundaries
1981 .iter()
1982 .flat_map(|record| record.admitted_attempts.iter().copied())
1983 .collect();
1984 if admitted_attempts.len() != self.state.evidence.command_attempts.len()
1985 || self
1986 .state
1987 .evidence
1988 .command_attempts
1989 .iter()
1990 .any(|attempt| !admitted_attempts.contains(&attempt.id))
1991 {
1992 return Err(CanwuError::new(
1993 ErrorCode::ArchiveNotReady,
1994 "live evidence sealing requires every retained command attempt to belong to a completed boundary",
1995 ));
1996 }
1997 let admitted_commands: std::collections::BTreeSet<_> = self
1998 .state
1999 .evidence
2000 .boundaries
2001 .iter()
2002 .flat_map(|record| record.admitted_commands.iter().copied())
2003 .collect();
2004 if admitted_commands.len() != self.state.evidence.commands.len()
2005 || self
2006 .state
2007 .evidence
2008 .commands
2009 .iter()
2010 .any(|command| !admitted_commands.contains(&command.id))
2011 {
2012 return Err(CanwuError::new(
2013 ErrorCode::ArchiveNotReady,
2014 "live evidence sealing requires every retained command to belong to a completed boundary",
2015 ));
2016 }
2017 let admitted_ingress: std::collections::BTreeSet<_> = self
2018 .state
2019 .evidence
2020 .boundaries
2021 .iter()
2022 .flat_map(|record| record.admitted_ingress.iter().copied())
2023 .collect();
2024 if admitted_ingress.len() != self.state.evidence.ingress.len()
2025 || self
2026 .state
2027 .evidence
2028 .ingress
2029 .iter()
2030 .any(|record| !admitted_ingress.contains(&record.id))
2031 {
2032 return Err(CanwuError::new(
2033 ErrorCode::ArchiveNotReady,
2034 "live evidence sealing requires every retained ingress record to belong to a completed boundary",
2035 ));
2036 }
2037 let admitted_events: std::collections::BTreeSet<_> = self
2038 .state
2039 .evidence
2040 .boundaries
2041 .iter()
2042 .flat_map(|record| record.admitted_events.iter().copied())
2043 .collect();
2044 if self.state.counters.admitted_event_count
2045 != self
2046 .state
2047 .evidence
2048 .archived
2049 .event_count
2050 .checked_add(
2051 u64::try_from(self.state.evidence.events.len()).map_err(|_| {
2052 CanwuError::new(
2053 ErrorCode::ArchiveNotReady,
2054 "retained event count exceeds the live archive cursor range",
2055 )
2056 })?,
2057 )
2058 .ok_or_else(|| {
2059 CanwuError::new(
2060 ErrorCode::ArchiveNotReady,
2061 "retained event cursor is exhausted",
2062 )
2063 })?
2064 || admitted_events.len() != self.state.evidence.events.len()
2065 || self
2066 .state
2067 .evidence
2068 .events
2069 .iter()
2070 .any(|event| !admitted_events.contains(&event.id))
2071 {
2072 return Err(CanwuError::new(
2073 ErrorCode::ArchiveNotReady,
2074 "live evidence sealing requires every retained event to be admitted by a later completed boundary",
2075 ));
2076 }
2077 Ok(())
2078 }
2079
2080 fn seal_retained_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
2081 let start = self.state.evidence.archived;
2082 let end = self.evidence_cursor()?;
2083 if start == end {
2084 return Ok(None);
2085 }
2086 self.ensure_retained_evidence_is_sealable()?;
2087 let dependencies = self.evidence_dependencies()?;
2088 let checkpoint_hash = self.state.metadata.checkpoint_hash.clone();
2089 let commitment_roots = self.state.metadata.commitment_roots.clone();
2090 let commitment_cache = self.state.metadata.commitment_cache.clone();
2091 let prepared = (|| {
2092 self.refresh_checkpoint_hash()?;
2093
2094 let mut archived_command_requests = Vec::new();
2095 for attempt in &self.state.evidence.command_attempts {
2096 let Some(request_id) = attempt.request_id else {
2097 continue;
2098 };
2099 let outcome = self.command_outcome_from_attempt(attempt)?;
2100 archived_command_requests.push((
2101 request_id,
2102 ArchivedCommandRequestOutcome {
2103 input_hash: super::canonical_hash(
2104 "canwu.archive.command.request.v1",
2105 &(attempt.expected_revision, &attempt.envelope),
2106 )?,
2107 outcome,
2108 },
2109 ));
2110 }
2111
2112 let mut archived_ingress_requests = Vec::new();
2113 let mut archived_decision_requests = Vec::new();
2114 let mut archived_decision_command_requests = Vec::new();
2115 for record in &self.state.evidence.ingress {
2116 let receipt = IngressReceipt {
2117 ingress_id: record.id,
2118 issued_at: record.issued_at,
2119 due_at: record.due_at,
2120 };
2121 match &record.payload {
2122 IngressPayload::Command { request } => {
2123 archived_ingress_requests.push((
2124 request.request_id,
2125 ArchivedIngressRequest {
2126 input_hash: super::canonical_hash(
2127 "canwu.archive.ingress.command.v1",
2128 &(record.due_at, record.priority, request.as_ref()),
2129 )?,
2130 receipt,
2131 },
2132 ));
2133 }
2134 IngressPayload::Decision { request } => {
2135 if let Some(command) = &request.command {
2136 archived_decision_command_requests.push(command.request_id);
2137 }
2138 archived_decision_requests.push((
2139 request.request_id,
2140 ArchivedIngressRequest {
2141 input_hash: super::canonical_hash(
2142 "canwu.ingress.decision-request.v1",
2143 &(record.due_at, record.priority, request.as_ref()),
2144 )?,
2145 receipt,
2146 },
2147 ));
2148 }
2149 IngressPayload::Plugin { .. } | IngressPayload::Calendar { .. } => {}
2150 }
2151 }
2152 Ok::<_, CanwuError>((
2153 archived_command_requests,
2154 archived_ingress_requests,
2155 archived_decision_requests,
2156 archived_decision_command_requests,
2157 ))
2158 })();
2159 let (
2160 archived_command_requests,
2161 archived_ingress_requests,
2162 archived_decision_requests,
2163 archived_decision_command_requests,
2164 ) = match prepared {
2165 Ok(prepared) => prepared,
2166 Err(error) => {
2167 self.state.metadata.checkpoint_hash = checkpoint_hash;
2168 self.state.metadata.commitment_roots = commitment_roots;
2169 self.state.metadata.commitment_cache = commitment_cache;
2170 return Err(error);
2171 }
2172 };
2173
2174 let mut segment = EvidenceJournalSegment {
2175 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
2176 start,
2177 end,
2178 events: self.state.evidence.events.clone(),
2179 commands: self.state.evidence.commands.clone(),
2180 command_attempts: self.state.evidence.command_attempts.clone(),
2181 ingress: self.state.evidence.ingress.clone(),
2182 boundaries: self.state.evidence.boundaries.clone(),
2183 random_draws: self.state.evidence.random_draws.clone(),
2184 archive: None,
2185 };
2186 let (archive, receipts) = evidence_archive_index(&segment)?;
2187 segment.archive = Some(archive.clone());
2188 verify_archived_segment(&segment)?;
2189
2190 let receipt_map: BTreeMap<_, _> = receipts
2191 .iter()
2192 .cloned()
2193 .map(|receipt| (receipt.evidence.clone(), receipt))
2194 .collect();
2195 let mut reservations = Vec::new();
2196 for draw in &segment.random_draws {
2197 let RandomDrawAddress::OperationV1(address) = &draw.address else {
2198 continue;
2199 };
2200 let operation_evidence = draw.operation_evidence.clone().ok_or_else(|| {
2201 archive_error("operation-keyed archived draw is missing operation evidence")
2202 })?;
2203 let draw_receipt = receipt_map
2204 .get(&EvidenceRef::RandomDraw(draw.id))
2205 .cloned()
2206 .ok_or_else(|| {
2207 archive_error("operation-keyed archived draw is missing its receipt")
2208 })?;
2209 reservations.push(KeyedDrawReservation {
2210 stream: draw.stream.clone(),
2211 address: address.clone(),
2212 upper_exclusive: draw.upper_exclusive,
2213 purpose_hash: super::random::purpose_hash_hex_v1(&draw.purpose)?,
2214 result: draw.value,
2215 draw_id: draw.id,
2216 operation_evidence,
2217 draw_receipt,
2218 });
2219 }
2220 reservations.sort_by(|left, right| {
2221 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2222 });
2223
2224 self.state.evidence.archived_boundary_head = self
2225 .state
2226 .evidence
2227 .boundaries
2228 .last()
2229 .map(|record| record.hash.clone())
2230 .or_else(|| self.state.evidence.archived_boundary_head.clone());
2231 self.state.evidence.archived_legacy_commands |= self
2232 .state
2233 .evidence
2234 .commands
2235 .iter()
2236 .any(|record| record.attempt_id.is_none());
2237 self.state.evidence.archived_tracked_attempts |=
2238 !self.state.evidence.command_attempts.is_empty()
2239 || !self.state.evidence.ingress.is_empty();
2240 self.state.evidence.archived_unqueued_command_history |= has_unqueued_command_history(
2241 &self.state.evidence.commands,
2242 &self.state.evidence.command_attempts,
2243 &self.state.evidence.ingress,
2244 );
2245 self.state
2246 .evidence
2247 .archived_command_requests
2248 .extend(archived_command_requests);
2249 self.state
2250 .evidence
2251 .archived_ingress_requests
2252 .extend(archived_ingress_requests);
2253 self.state
2254 .evidence
2255 .archived_decision_requests
2256 .extend(archived_decision_requests);
2257 self.state
2258 .evidence
2259 .archived_decision_command_requests
2260 .extend(archived_decision_command_requests);
2261 self.state.evidence.archived = end;
2262 self.state.evidence.events.clear();
2263 self.state.evidence.commands.clear();
2264 self.state.evidence.command_attempts.clear();
2265 self.state.evidence.ingress.clear();
2266 self.state.evidence.boundaries.clear();
2267 self.state.evidence.random_draws.clear();
2268 self.state
2269 .evidence
2270 .archived_segment_headers
2271 .push(archive.header);
2272 for receipt in receipts {
2273 if self
2274 .state
2275 .evidence
2276 .archived_evidence_receipts
2277 .insert(receipt.evidence.clone(), receipt)
2278 .is_some()
2279 {
2280 return Err(archive_error("archived evidence receipt was duplicated"));
2281 }
2282 }
2283 self.state
2284 .evidence
2285 .keyed_draw_reservations
2286 .extend(reservations);
2287 self.state
2288 .evidence
2289 .keyed_draw_reservations
2290 .sort_by(|left, right| {
2291 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2292 });
2293 if self
2294 .state
2295 .evidence
2296 .keyed_draw_reservations
2297 .windows(2)
2298 .any(|reservations| {
2299 (&reservations[0].stream, &reservations[0].address)
2300 == (&reservations[1].stream, &reservations[1].address)
2301 })
2302 {
2303 return Err(archive_error("archived keyed reservation was duplicated"));
2304 }
2305 retain_reachable_archived_evidence_receipts(
2306 &mut self.state.evidence.archived_evidence_receipts,
2307 &dependencies,
2308 &self.state.evidence.keyed_draw_reservations,
2309 );
2310 for dependency in dependencies
2311 .iter()
2312 .filter(|dependency| dependency.requirement == EvidenceRequirement::PayloadRequired)
2313 {
2314 if !self
2315 .state
2316 .evidence
2317 .archived_evidence_receipts
2318 .contains_key(&dependency.reference)
2319 {
2320 return Err(CanwuError::new(
2321 ErrorCode::ArchiveNotReady,
2322 "payload-required evidence was not present in the sealed archive prefix",
2323 ));
2324 }
2325 }
2326 Ok(Some(segment))
2327 }
2328
2329 pub(super) fn checkpoint_state(&self) -> SimulationSnapshot {
2330 SimulationSnapshot {
2331 engine_version: ENGINE_VERSION.to_owned(),
2332 snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
2333 run_manifest: Some(self.state.metadata.run_manifest.clone()),
2334 run_manifest_hash: self.state.metadata.run_manifest_hash.clone(),
2335 run_configuration: Some(self.state.metadata.run_configuration.clone()),
2336 checkpoint_hash: self.state.metadata.checkpoint_hash.clone(),
2337 commitment_format_version: self.state.metadata.commitment_format_version,
2338 commitment_roots: self.state.metadata.commitment_roots.clone(),
2339 revision_format_version: STATE_REVISION_FORMAT_VERSION,
2340 state_revision: self.state.counters.state_revision,
2341 replay_revision_format_version: self.state.metadata.replay_revision_format_version,
2342 admission_cursor_format_version: ADMISSION_CURSOR_FORMAT_VERSION,
2343 admitted_attempt_count: self.state.counters.admitted_attempt_count,
2344 admitted_command_count: self.state.counters.admitted_command_count,
2345 admitted_event_count: self.state.counters.admitted_event_count,
2346 initial_time: self.state.scheduler.initial_time,
2347 initial_scenario: self.bound_initial_scenario().cloned(),
2348 now: self.state.scheduler.now,
2349 plugin_registration_closed: self.state.metadata.plugin_registration_closed,
2350 world: self.world(),
2351 knowledge: self.state.current.knowledge.clone(),
2352 events: Vec::new(),
2353 commands: Vec::new(),
2354 command_attempts: Vec::new(),
2355 ingress: Vec::new(),
2356 boundaries: Vec::new(),
2357 plugin_components: self
2358 .state
2359 .current
2360 .plugin_components
2361 .values()
2362 .cloned()
2363 .collect(),
2364 domain_records: self
2365 .state
2366 .current
2367 .domain_records
2368 .values()
2369 .cloned()
2370 .collect(),
2371 decisions: self.state.current.decisions.clone(),
2372 plugin_descriptors: self.plugins.descriptors().cloned().collect(),
2373 schema: self.schema.clone(),
2374 root_seed: self.state.current.root_seed,
2375 random_streams: self
2376 .state
2377 .current
2378 .random_streams
2379 .values()
2380 .cloned()
2381 .collect(),
2382 random_draws: Vec::new(),
2383 scheduled: self
2384 .state
2385 .scheduler
2386 .actions
2387 .iter()
2388 .map(|(key, action)| ScheduledRecord {
2389 key: key.clone(),
2390 action: action.clone(),
2391 })
2392 .collect(),
2393 legacy_rng: None,
2394 next_event_id: self.state.counters.next_event_id,
2395 next_command_id: self.state.counters.next_command_id,
2396 next_command_attempt_id: self.state.counters.next_command_attempt_id,
2397 next_ingress_id: self.state.counters.next_ingress_id,
2398 next_boundary_id: self.state.counters.next_boundary_id,
2399 next_random_draw_id: self.state.counters.next_random_draw_id,
2400 next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
2401 next_schedule_sequence: self.state.counters.next_schedule_sequence,
2402 next_correlation_id: self.state.counters.next_correlation_id,
2403 next_decision_trace_id: self.state.counters.next_decision_trace_id,
2404 }
2405 }
2406
2407 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
2409 EvidenceCursor::from_evidence(&self.state.evidence)
2410 }
2411
2412 pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
2414 let archived_segment_headers = self.state.evidence.archived_segment_headers.clone();
2415 let evidence_dependencies = self.evidence_dependencies()?;
2416 let mut keyed_draw_reservations = self.state.evidence.keyed_draw_reservations.clone();
2417 keyed_draw_reservations.sort_by(|left, right| {
2418 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2419 });
2420 let required_receipts =
2421 required_archived_receipt_references(&evidence_dependencies, &keyed_draw_reservations);
2422 let archived_evidence_receipts: Vec<_> = self
2423 .state
2424 .evidence
2425 .archived_evidence_receipts
2426 .iter()
2427 .filter(|(reference, _)| required_receipts.contains(*reference))
2428 .map(|(_, receipt)| receipt.clone())
2429 .collect();
2430 let checkpoint = SimulationCheckpoint {
2431 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
2432 journal_end: self.evidence_cursor()?,
2433 state: self.checkpoint_state(),
2434 archived_segment_manifest_root: skipped_commitment_root(
2435 ARCHIVED_SEGMENT_MANIFEST_DOMAIN,
2436 &archived_segment_headers,
2437 )?,
2438 archived_segment_headers,
2439 archived_receipt_root: skipped_commitment_root(
2440 ARCHIVED_RECEIPT_DOMAIN,
2441 &archived_evidence_receipts,
2442 )?,
2443 archived_evidence_receipts,
2444 evidence_dependency_root: skipped_commitment_root(
2445 EVIDENCE_DEPENDENCY_DOMAIN,
2446 &evidence_dependencies,
2447 )?,
2448 evidence_dependencies,
2449 keyed_reservation_root: skipped_commitment_root(
2450 KEYED_RESERVATION_DOMAIN,
2451 &keyed_draw_reservations,
2452 )?,
2453 keyed_draw_reservations,
2454 };
2455 validate_compact_continuation(&checkpoint)?;
2456 Ok(checkpoint)
2457 }
2458
2459 pub fn journal_segment_since(
2461 &self,
2462 start: EvidenceCursor,
2463 ) -> Result<EvidenceJournalSegment, CanwuError> {
2464 let end = self.evidence_cursor()?;
2465 let cut = |value: u64, archived: u64, len: usize, label: &str| {
2466 let value = value.checked_sub(archived).ok_or_else(|| {
2467 CanwuError::new(
2468 ErrorCode::InvalidSnapshot,
2469 format!("{label} journal cursor precedes the retained live evidence window"),
2470 )
2471 })?;
2472 let value = usize::try_from(value).map_err(|_| {
2473 CanwuError::new(
2474 ErrorCode::InvalidSnapshot,
2475 format!("{label} journal cursor is not representable on this platform"),
2476 )
2477 })?;
2478 if value > len {
2479 return Err(CanwuError::new(
2480 ErrorCode::InvalidSnapshot,
2481 format!("{label} journal cursor exceeds the current evidence tail"),
2482 ));
2483 }
2484 Ok(value)
2485 };
2486 let archived = self.state.evidence.archived;
2487 let event_start = cut(
2488 start.event_count,
2489 archived.event_count,
2490 self.state.evidence.events.len(),
2491 "event",
2492 )?;
2493 let command_start = cut(
2494 start.command_count,
2495 archived.command_count,
2496 self.state.evidence.commands.len(),
2497 "command",
2498 )?;
2499 let attempt_start = cut(
2500 start.command_attempt_count,
2501 archived.command_attempt_count,
2502 self.state.evidence.command_attempts.len(),
2503 "command-attempt",
2504 )?;
2505 let ingress_start = cut(
2506 start.ingress_count,
2507 archived.ingress_count,
2508 self.state.evidence.ingress.len(),
2509 "ingress",
2510 )?;
2511 let boundary_start = cut(
2512 start.boundary_count,
2513 archived.boundary_count,
2514 self.state.evidence.boundaries.len(),
2515 "boundary",
2516 )?;
2517 let draw_start = cut(
2518 start.random_draw_count,
2519 archived.random_draw_count,
2520 self.state.evidence.random_draws.len(),
2521 "random-draw",
2522 )?;
2523 Ok(EvidenceJournalSegment {
2524 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
2525 start,
2526 end,
2527 events: self.state.evidence.events[event_start..].to_vec(),
2528 commands: self.state.evidence.commands[command_start..].to_vec(),
2529 command_attempts: self.state.evidence.command_attempts[attempt_start..].to_vec(),
2530 ingress: self.state.evidence.ingress[ingress_start..].to_vec(),
2531 boundaries: self.state.evidence.boundaries[boundary_start..].to_vec(),
2532 random_draws: self.state.evidence.random_draws[draw_start..].to_vec(),
2533 archive: None,
2534 })
2535 }
2536
2537 pub fn checkpoint_journal(&self) -> Result<CheckpointJournal, CanwuError> {
2539 if self.state.evidence.archived != EvidenceCursor::default() {
2540 return Err(CanwuError::new(
2541 ErrorCode::InvalidSnapshot,
2542 "a compact live runtime requires its previously sealed evidence segments to build a portable save",
2543 ));
2544 }
2545 let segment = self.journal_segment_since(EvidenceCursor::default())?;
2546 Ok(CheckpointJournal {
2547 checkpoint: self.checkpoint()?,
2548 segments: (segment.start != segment.end)
2549 .then_some(segment)
2550 .into_iter()
2551 .collect(),
2552 })
2553 }
2554
2555 pub fn checkpoint_journal_json(&self) -> Result<String, CanwuError> {
2557 serde_json::to_string_pretty(&self.checkpoint_journal()?).map_err(|error| {
2558 CanwuError::new(
2559 ErrorCode::InvalidSnapshot,
2560 format!("could not serialize checkpoint journal: {error}"),
2561 )
2562 })
2563 }
2564
2565 fn snapshot_from_checkpoint_and_journal(
2566 checkpoint: SimulationCheckpoint,
2567 segments: Vec<EvidenceJournalSegment>,
2568 ) -> Result<SimulationSnapshot, CanwuError> {
2569 if checkpoint.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
2570 return Err(invalid_snapshot_error(format!(
2571 "checkpoint-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
2572 checkpoint.format_version
2573 )));
2574 }
2575 validate_compact_continuation(&checkpoint)?;
2576 let expected_headers = checkpoint.archived_segment_headers;
2577 let expected_receipts = checkpoint.archived_evidence_receipts;
2578 let expected_dependencies = checkpoint.evidence_dependencies;
2579 let expected_reservations = checkpoint.keyed_draw_reservations;
2580 let mut snapshot = checkpoint.state;
2581 if snapshot.snapshot_format_version != SNAPSHOT_FORMAT_VERSION {
2582 return Err(invalid_snapshot_error(format!(
2583 "checkpoint-journal format {CHECKPOINT_JOURNAL_FORMAT_VERSION} requires snapshot format {SNAPSHOT_FORMAT_VERSION}"
2584 )));
2585 }
2586 if !snapshot.events.is_empty()
2587 || !snapshot.commands.is_empty()
2588 || !snapshot.command_attempts.is_empty()
2589 || !snapshot.ingress.is_empty()
2590 || !snapshot.boundaries.is_empty()
2591 || !snapshot.random_draws.is_empty()
2592 {
2593 return Err(invalid_snapshot_error(
2594 "checkpoint current state must not duplicate append-only evidence",
2595 ));
2596 }
2597
2598 let mut cursor = EvidenceCursor::default();
2599 let mut rebuilt_headers = Vec::new();
2600 let mut rebuilt_receipts = BTreeMap::new();
2601 let mut rebuilt_reservations = Vec::new();
2602 for segment in segments {
2603 if segment.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
2604 return Err(invalid_snapshot_error(format!(
2605 "evidence-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
2606 segment.format_version
2607 )));
2608 }
2609 if segment.start != cursor {
2610 return Err(invalid_snapshot_error(
2611 "evidence-journal segments must form one contiguous global prefix",
2612 ));
2613 }
2614 let end = cursor.checked_advance(&segment)?;
2615 if end == cursor {
2616 return Err(invalid_snapshot_error(
2617 "evidence-journal segments must advance at least one journal cursor",
2618 ));
2619 }
2620 if segment.end != end {
2621 return Err(invalid_snapshot_error(
2622 "evidence-journal segment end does not match its encoded records",
2623 ));
2624 }
2625 if let Some(archive) = &segment.archive {
2626 let receipts = verify_archived_segment(&segment).map_err(|error| {
2627 invalid_snapshot_error(format!(
2628 "archived evidence segment is invalid: {}",
2629 error.message
2630 ))
2631 })?;
2632 rebuilt_headers.push(archive.header.clone());
2633 for receipt in receipts {
2634 if rebuilt_receipts
2635 .insert(receipt.evidence.clone(), receipt)
2636 .is_some()
2637 {
2638 return Err(invalid_snapshot_error(
2639 "archived evidence segments contain duplicate receipts",
2640 ));
2641 }
2642 }
2643 for draw in &segment.random_draws {
2644 let RandomDrawAddress::OperationV1(address) = &draw.address else {
2645 continue;
2646 };
2647 let operation_evidence = draw.operation_evidence.clone().ok_or_else(|| {
2648 invalid_snapshot_error("archived keyed draw is missing operation evidence")
2649 })?;
2650 let draw_receipt = rebuilt_receipts
2651 .get(&EvidenceRef::RandomDraw(draw.id))
2652 .cloned()
2653 .ok_or_else(|| {
2654 invalid_snapshot_error("archived keyed draw receipt is missing")
2655 })?;
2656 rebuilt_reservations.push(KeyedDrawReservation {
2657 stream: draw.stream.clone(),
2658 address: address.clone(),
2659 upper_exclusive: draw.upper_exclusive,
2660 purpose_hash: super::random::purpose_hash_hex_v1(&draw.purpose)?,
2661 result: draw.value,
2662 draw_id: draw.id,
2663 operation_evidence,
2664 draw_receipt,
2665 });
2666 }
2667 }
2668 snapshot.events.extend(segment.events);
2669 snapshot.commands.extend(segment.commands);
2670 snapshot.command_attempts.extend(segment.command_attempts);
2671 snapshot.ingress.extend(segment.ingress);
2672 snapshot.boundaries.extend(segment.boundaries);
2673 snapshot.random_draws.extend(segment.random_draws);
2674 cursor = end;
2675 }
2676 if cursor != checkpoint.journal_end {
2677 return Err(invalid_snapshot_error(
2678 "evidence-journal segments do not reach the checkpoint journal cut",
2679 ));
2680 }
2681 rebuilt_reservations.sort_by(|left, right| {
2682 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2683 });
2684 let rebuilt_dependencies =
2685 Simulation::from_snapshot(snapshot.clone())?.evidence_dependencies()?;
2686 retain_reachable_archived_evidence_receipts(
2687 &mut rebuilt_receipts,
2688 &rebuilt_dependencies,
2689 &rebuilt_reservations,
2690 );
2691 let rebuilt_receipts: Vec<_> = rebuilt_receipts.into_values().collect();
2692 if rebuilt_headers != expected_headers
2693 || rebuilt_receipts != expected_receipts
2694 || rebuilt_dependencies != expected_dependencies
2695 || rebuilt_reservations != expected_reservations
2696 {
2697 return Err(invalid_snapshot_error(
2698 "checkpoint compact continuation does not match its reconstructed authoritative indexes",
2699 ));
2700 }
2701 Ok(snapshot)
2702 }
2703
2704 pub fn from_checkpoint_and_journal(
2706 checkpoint: SimulationCheckpoint,
2707 segments: Vec<EvidenceJournalSegment>,
2708 ) -> Result<Self, CanwuError> {
2709 Self::from_snapshot(Self::snapshot_from_checkpoint_and_journal(
2710 checkpoint, segments,
2711 )?)
2712 }
2713
2714 pub fn from_checkpoint_journal(bundle: CheckpointJournal) -> Result<Self, CanwuError> {
2716 Self::from_checkpoint_and_journal(bundle.checkpoint, bundle.segments)
2717 }
2718
2719 pub fn from_checkpoint_journal_with_plugins(
2721 bundle: CheckpointJournal,
2722 plugins: &[&dyn SimulationPlugin],
2723 ) -> Result<Self, CanwuError> {
2724 let mut simulation = Self::from_checkpoint_journal(bundle)?;
2725 for plugin in plugins {
2726 simulation.register_plugin(*plugin)?;
2727 }
2728 simulation.ensure_runtime_ready()?;
2729 Ok(simulation)
2730 }
2731
2732 pub fn from_checkpoint_journal_json(json: &str) -> Result<Self, CanwuError> {
2734 let bundle = super::legacy_v4::deserialize_checkpoint_journal_json(json)?;
2735 Self::from_checkpoint_journal(bundle)
2736 }
2737
2738 pub fn from_checkpoint_journal_json_with_plugins(
2740 json: &str,
2741 plugins: &[&dyn SimulationPlugin],
2742 ) -> Result<Self, CanwuError> {
2743 let bundle = super::legacy_v4::deserialize_checkpoint_journal_json(json)?;
2744 Self::from_checkpoint_journal_with_plugins(bundle, plugins)
2745 }
2746}
2747
2748#[cfg(test)]
2749mod tests {
2750 #![allow(clippy::unnecessary_literal_bound, clippy::unnecessary_wraps)]
2751 use super::super::{
2752 BoundaryContext, BoundaryDirective, BoundaryPhase, BoundaryProposal,
2753 BoundarySystemContract, Command, CommandContext, CommandRequestId, DomainRecordClass,
2754 DomainRecordDraft, DomainRecordMutation, Issuer, KnowledgeHolderRef, KnowledgeOrigin,
2755 KnowledgeRecordDraft, KnowledgeRecordKind, KnowledgeSchemaId, KnowledgeWriteGrant,
2756 PluginActionDescriptor, PluginKnowledgeSchema, PluginRegistrar, SimulationView, StateKey,
2757 StateVisibility, SystemDirective, demo_scenario,
2758 };
2759 use super::*;
2760 use canwu_core::{BoundaryId, CommandId, DomainRecordKind, PersonId};
2761 use serde_json::{Map, Value, json};
2762 use std::cell::RefCell;
2763
2764 #[derive(Default)]
2765 struct TestArchive {
2766 segments: RefCell<BTreeMap<String, EvidenceJournalSegment>>,
2767 }
2768
2769 impl TestArchive {
2770 fn segment_ids(&self) -> Vec<String> {
2771 self.segments.borrow().keys().cloned().collect()
2772 }
2773 }
2774
2775 impl ArchiveProvider for TestArchive {
2776 fn load_evidence_segment(
2777 &self,
2778 segment_id: &str,
2779 ) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
2780 Ok(self.segments.borrow().get(segment_id).cloned())
2781 }
2782 }
2783
2784 impl ArchiveStore for TestArchive {
2785 fn store_evidence_segment(
2786 &self,
2787 segment: &EvidenceJournalSegment,
2788 ) -> Result<ArchiveStoreOutcome, CanwuError> {
2789 let segment_id = segment
2790 .archive
2791 .as_ref()
2792 .ok_or_else(|| archive_error("test archive segment has no index"))?
2793 .header
2794 .segment_id
2795 .clone();
2796 let mut segments = self.segments.borrow_mut();
2797 if let Some(existing) = segments.get(&segment_id) {
2798 return if existing == segment {
2799 Ok(ArchiveStoreOutcome::AlreadyPresent)
2800 } else {
2801 Err(archive_error(
2802 "content-addressed test segment ID has conflicting bytes",
2803 ))
2804 };
2805 }
2806 segments.insert(segment_id, segment.clone());
2807 Ok(ArchiveStoreOutcome::Stored)
2808 }
2809 }
2810
2811 fn archived_identity_schema() -> KnowledgeSchemaId {
2812 KnowledgeSchemaId::new(
2813 KnowledgeRecordKind::new("fixture.archive", "archived_command_notice"),
2814 1,
2815 )
2816 }
2817
2818 #[allow(clippy::unnecessary_wraps)]
2819 fn retain_archive_source_command(
2820 _view: &SimulationView<'_>,
2821 _context: &CommandContext,
2822 _payload: &Value,
2823 ) -> Result<Vec<SystemDirective>, CanwuError> {
2824 Ok(Vec::new())
2825 }
2826
2827 #[allow(clippy::unnecessary_wraps)]
2828 fn publish_archived_command_identity(
2829 _view: &SimulationView<'_>,
2830 context: &BoundaryContext,
2831 ) -> Result<BoundaryProposal, CanwuError> {
2832 let record = KnowledgeRecordDraft {
2833 schema: archived_identity_schema(),
2834 subjects: Vec::new(),
2835 payload: json!({ "boundary": context.boundary_id.get() }),
2836 as_of: None,
2837 confidence_per_mille: 1_000,
2838 origin: KnowledgeOrigin {
2839 method: "archived_command_identity_v1".to_owned(),
2840 evidence: vec![EvidenceRef::Command(CommandId::new(1))],
2841 },
2842 supersedes: Vec::new(),
2843 contradicts: Vec::new(),
2844 };
2845 Ok(BoundaryProposal {
2846 directives: vec![BoundaryDirective::PublishKnowledge {
2847 holder: KnowledgeHolderRef::Person(PersonId::new(1)),
2848 visibility: StateVisibility::SameBoundary,
2849 producer_correlation: Some(format!(
2850 "archived-command-boundary-{}",
2851 context.boundary_id.get()
2852 )),
2853 records: vec![record],
2854 summary: "Publish knowledge from a retained or archived command identity"
2855 .to_owned(),
2856 }],
2857 ..BoundaryProposal::default()
2858 })
2859 }
2860
2861 struct ArchivedIdentityPublicationPlugin;
2862
2863 impl SimulationPlugin for ArchivedIdentityPublicationPlugin {
2864 fn name(&self) -> &str {
2865 "fixture-archived-identity-publication"
2866 }
2867
2868 fn version(&self) -> &str {
2869 "test-v1"
2870 }
2871
2872 fn semantic_hash(&self) -> &str {
2873 "7100000000000000000000000000000000000000000000000000000000000000"
2874 }
2875
2876 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
2877 registrar.register_knowledge_schema(PluginKnowledgeSchema {
2878 id: archived_identity_schema(),
2879 schema_hash: "7200000000000000000000000000000000000000000000000000000000000000"
2880 .to_owned(),
2881 writable: true,
2882 payload_schema: PayloadSchema::Any,
2883 subjects: Vec::new(),
2884 })?;
2885 registrar.register_command(
2886 PluginActionDescriptor {
2887 name: "retain_archive_source_v1".to_owned(),
2888 description: "Persist one neutral command identity for archive testing"
2889 .to_owned(),
2890 payload_schema: PayloadSchema::Any,
2891 reads: Vec::new(),
2892 writes: Vec::new(),
2893 },
2894 retain_archive_source_command,
2895 )?;
2896 let mut publisher = BoundarySystemContract::new(
2897 "publish-archived-command-identity",
2898 BoundaryPhase::PerspectiveAndReportMaterialization,
2899 SystemCadence::Daily,
2900 );
2901 publisher.knowledge_writes = vec![KnowledgeWriteGrant {
2902 schema: archived_identity_schema(),
2903 visibilities: vec![StateVisibility::SameBoundary],
2904 }];
2905 registrar.register_boundary_system(publisher, publish_archived_command_identity)
2906 }
2907 }
2908
2909 fn archived_identity_two_segment_fixture() -> (SimulationCheckpoint, Vec<EvidenceJournalSegment>)
2910 {
2911 let (scenario, _) = demo_scenario();
2912 let plugin = ArchivedIdentityPublicationPlugin;
2913 let mut simulation = Simulation::new(711, scenario).expect("fixture scenario should load");
2914 simulation
2915 .register_plugin(&plugin)
2916 .expect("archived-identity plugin should register");
2917 simulation
2918 .enqueue_command(
2919 SimTime::EPOCH,
2920 0,
2921 CommandRequest::new(
2922 CommandRequestId::new(1),
2923 simulation.revision(),
2924 CommandEnvelope::new(
2925 Issuer::System("archive-identity-fixture".to_owned()),
2926 Command::Plugin {
2927 plugin: plugin.name().to_owned(),
2928 command: "retain_archive_source_v1".to_owned(),
2929 payload: json!({ "format_version": 1 }),
2930 },
2931 )
2932 .at_time(SimTime::EPOCH),
2933 ),
2934 )
2935 .expect("archive source command should queue");
2936 simulation
2937 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
2938 .expect("boundary one should retain the command-backed publication");
2939 simulation
2940 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
2941 .expect("the following cut should admit boundary-one events before sealing");
2942
2943 let mut compact = simulation
2944 .into_compacted()
2945 .expect("the fixture should enter compact mode");
2946 let first_segment = compact
2947 .seal_evidence()
2948 .expect("boundary-one evidence should seal")
2949 .expect("boundary one should produce an archive segment");
2950 let first_checkpoint = compact.checkpoint().expect("checkpoint one should build");
2951 assert!(
2952 first_checkpoint
2953 .archived_evidence_receipts
2954 .iter()
2955 .any(|receipt| receipt.evidence == EvidenceRef::Command(CommandId::new(1)))
2956 );
2957
2958 let mut restored = CompactedSimulation::from_checkpoint_and_journal_with_plugins(
2959 first_checkpoint,
2960 vec![first_segment.clone()],
2961 &[&plugin],
2962 )
2963 .expect("checkpoint one should restore with its exact archive prefix");
2964 let restored_prefix = restored
2965 .seal_evidence()
2966 .expect("restored boundary-one evidence should reseal")
2967 .expect("restored boundary-one evidence should remain non-empty");
2968 assert_eq!(restored_prefix, first_segment);
2969 assert!(
2970 restored
2971 .archived_evidence_receipt(&EvidenceRef::Command(CommandId::new(1)))
2972 .is_some(),
2973 "the second publication must consume an archived identity, not retained payload"
2974 );
2975 let archived_reads = [StateKey::core_commands()];
2976 let archived_view = restored
2977 .simulation
2978 .plugin_view("archive-identity-probe", &archived_reads);
2979 let error = archived_view
2980 .command(CommandId::new(1))
2981 .expect_err("archived identity must not expose retained command payload");
2982 assert_eq!(error.code, ErrorCode::EvidenceContentUnavailable);
2983 restored
2984 .settle_boundary(
2985 BoundaryRequest::at(SimTime::EPOCH + SimDuration::days(1))
2986 .with_cadence(SystemCadence::Daily),
2987 )
2988 .expect("boundary two should accept the archived command identity");
2989 let holder = KnowledgeHolderRef::Person(PersonId::new(1));
2990 let records = restored
2991 .knowledge()
2992 .for_holder(&holder)
2993 .expect("the holder should have both publications");
2994 assert_eq!(records.len(), 2);
2995 assert!(records.values().all(|record| {
2996 record.origin.evidence == vec![EvidenceRef::Command(CommandId::new(1))]
2997 }));
2998 restored
2999 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH + SimDuration::days(1)))
3000 .expect("the following cut should admit boundary-two events before sealing");
3001
3002 let second_segment = restored
3003 .seal_evidence()
3004 .expect("boundary-two evidence should seal")
3005 .expect("boundary two should produce an archive segment");
3006 let second_checkpoint = restored.checkpoint().expect("checkpoint two should build");
3007 (second_checkpoint, vec![restored_prefix, second_segment])
3008 }
3009
3010 #[test]
3011 fn compact_restore_preserves_archived_identity_and_rejects_noncontiguous_segments() {
3012 let plugin = ArchivedIdentityPublicationPlugin;
3013 let (checkpoint, segments) = archived_identity_two_segment_fixture();
3014 let restored = CompactedSimulation::from_checkpoint_and_journal_with_plugins(
3015 checkpoint.clone(),
3016 segments.clone(),
3017 &[&plugin],
3018 )
3019 .expect("the complete two-segment archive should restore");
3020 assert_eq!(
3021 restored
3022 .knowledge()
3023 .for_holder(&KnowledgeHolderRef::Person(PersonId::new(1)))
3024 .expect("published holder ledger")
3025 .len(),
3026 2
3027 );
3028
3029 let first = segments[0].clone();
3030 let second = segments[1].clone();
3031 for (label, tampered) in [
3032 ("omission", vec![first.clone()]),
3033 (
3034 "overlap",
3035 vec![first.clone(), first.clone(), second.clone()],
3036 ),
3037 ("reorder", vec![second, first]),
3038 ] {
3039 let error = Simulation::from_checkpoint_and_journal(checkpoint.clone(), tampered)
3040 .err()
3041 .expect(label);
3042 assert_eq!(error.code, ErrorCode::InvalidSnapshot, "{label}");
3043 }
3044 }
3045
3046 struct PayloadContinuationPlugin;
3047
3048 fn continuation_record_ref() -> DomainRecordRef {
3049 DomainRecordRef {
3050 kind: DomainRecordKind::new("fixture.archive", "payload_continuation"),
3051 id: "primary".to_owned(),
3052 }
3053 }
3054
3055 fn continuation_payload(continuation: PayloadRequiredEvidenceContinuationV1) -> Value {
3056 Value::Object(Map::from_iter([(
3057 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD.to_owned(),
3058 serde_json::to_value(continuation).expect("fixture continuation should encode"),
3059 )]))
3060 }
3061
3062 fn mutate_payload_continuation(
3063 _view: &SimulationView<'_>,
3064 context: &BoundaryContext,
3065 ) -> Result<BoundaryProposal, CanwuError> {
3066 let directive = match context.boundary_id.get() {
3067 1 => Some(BoundaryDirective::MutateRecord {
3068 mutation: DomainRecordMutation::Create {
3069 record: DomainRecordDraft::new(
3070 continuation_record_ref(),
3071 continuation_payload(PayloadRequiredEvidenceContinuationV1::active(vec![
3072 EvidenceRef::Boundary(BoundaryId::new(1)),
3073 ])),
3074 ),
3075 },
3076 summary: "Create an active payload continuation".to_owned(),
3077 }),
3078 3 => Some(BoundaryDirective::MutateRecord {
3079 mutation: DomainRecordMutation::Update {
3080 record: DomainRecordDraft::new(
3081 continuation_record_ref(),
3082 continuation_payload(PayloadRequiredEvidenceContinuationV1::completed()),
3083 ),
3084 expected_version: 1,
3085 },
3086 summary: "Complete the payload continuation".to_owned(),
3087 }),
3088 _ => None,
3089 };
3090 Ok(BoundaryProposal {
3091 directives: directive.into_iter().collect(),
3092 ..BoundaryProposal::default()
3093 })
3094 }
3095
3096 impl SimulationPlugin for PayloadContinuationPlugin {
3097 fn name(&self) -> &str {
3098 "fixture-payload-continuation"
3099 }
3100
3101 fn version(&self) -> &str {
3102 "test-v1"
3103 }
3104
3105 fn semantic_hash(&self) -> &str {
3106 "7000000000000000000000000000000000000000000000000000000000000000"
3107 }
3108
3109 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
3110 let mut properties = BTreeMap::new();
3111 properties.insert(
3112 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD.to_owned(),
3113 payload_required_evidence_continuation_property_v1(),
3114 );
3115 let mut schema =
3116 DomainRecordSchema::new(continuation_record_ref().kind, DomainRecordClass::Record);
3117 schema.payload_schema = PayloadSchema::Object {
3118 properties,
3119 allow_additional: false,
3120 };
3121 let state = schema.state_key();
3122 registrar.register_record_schema(schema)?;
3123 let mut contract = BoundarySystemContract::new(
3124 "payload-continuation",
3125 BoundaryPhase::DomainDeltaProposal,
3126 SystemCadence::Daily,
3127 );
3128 contract.writes = vec![state];
3129 contract.visibility = StateVisibility::SameBoundary;
3130 registrar.register_boundary_system(contract, mutate_payload_continuation)
3131 }
3132 }
3133
3134 fn payload_continuation_runtime() -> CompactedSimulation {
3135 let (scenario, _) = demo_scenario();
3136 let mut simulation = Simulation::new(701, scenario).expect("fixture scenario should load");
3137 simulation
3138 .register_plugin(&PayloadContinuationPlugin)
3139 .expect("payload-continuation plugin should register");
3140 simulation
3141 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
3142 .expect("the active continuation should be created");
3143 simulation
3144 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
3145 .expect("a later boundary should admit the record-change event");
3146 simulation
3147 .into_compacted()
3148 .expect("the fixture should enter compact mode")
3149 }
3150
3151 fn store_prepared(
3152 compact: &CompactedSimulation,
3153 archive: &TestArchive,
3154 ) -> PreparedEvidenceSeal {
3155 let prepared = compact
3156 .prepare_evidence_seal()
3157 .expect("the fixture should prepare a seal")
3158 .expect("the retained tail should be non-empty");
3159 assert_eq!(
3160 archive
3161 .store_evidence_segment(&prepared.segment)
3162 .expect("the prepared segment should store"),
3163 ArchiveStoreOutcome::Stored
3164 );
3165 prepared
3166 }
3167
3168 #[test]
3169 fn payload_required_receipts_are_exactly_reachable_and_prune_after_completion() {
3170 let mut compact = payload_continuation_runtime();
3171 assert_eq!(
3172 compact
3173 .seal_evidence()
3174 .expect_err("direct sealing must reject payload continuations")
3175 .code,
3176 ErrorCode::ArchiveNotReady
3177 );
3178
3179 let archive = TestArchive::default();
3180 let first = store_prepared(&compact, &archive);
3181 compact
3182 .commit_evidence_seal(&first.token, &archive)
3183 .expect("provider-backed sealing should commit");
3184 let first_checkpoint = compact.checkpoint().expect("checkpoint should build");
3185 let first_references = first_checkpoint
3186 .archived_evidence_receipts
3187 .iter()
3188 .map(|receipt| receipt.evidence.clone())
3189 .collect::<BTreeSet<_>>();
3190 let first_dependencies = first_checkpoint
3191 .evidence_dependencies
3192 .iter()
3193 .map(|dependency| (dependency.reference.clone(), dependency.requirement))
3194 .collect::<BTreeMap<_, _>>();
3195 assert_eq!(
3196 first_dependencies.get(&EvidenceRef::Boundary(BoundaryId::new(1))),
3197 Some(&EvidenceRequirement::PayloadRequired)
3198 );
3199 assert_eq!(first_references.len(), first_dependencies.len());
3200 assert_eq!(
3201 first_references,
3202 first_dependencies.keys().cloned().collect::<BTreeSet<_>>()
3203 );
3204 assert!(
3205 first
3206 .segment
3207 .archive
3208 .as_ref()
3209 .expect("prepared segment should have an archive index")
3210 .entries
3211 .len()
3212 > first_references.len(),
3213 "the full segment index must outlive the reachable-only receipt set"
3214 );
3215
3216 compact
3217 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
3218 .expect("the continuation should complete");
3219 compact
3220 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
3221 .expect("a later boundary should admit the completion event");
3222 let second = store_prepared(&compact, &archive);
3223 compact
3224 .commit_evidence_seal(&second.token, &archive)
3225 .expect("the completed continuation should seal");
3226 let second_checkpoint = compact.checkpoint().expect("checkpoint should build");
3227 assert_eq!(second_checkpoint.archived_segment_headers.len(), 2);
3228 assert_eq!(second_checkpoint.archived_evidence_receipts.len(), 1);
3229 assert_eq!(second_checkpoint.evidence_dependencies.len(), 1);
3230 assert_eq!(
3231 second_checkpoint.archived_evidence_receipts[0].evidence,
3232 second_checkpoint.evidence_dependencies[0].reference
3233 );
3234 assert_eq!(
3235 second_checkpoint.evidence_dependencies[0].requirement,
3236 EvidenceRequirement::IdentityOnly
3237 );
3238 assert!(
3239 !second_checkpoint
3240 .archived_evidence_receipts
3241 .iter()
3242 .any(|receipt| receipt.evidence == EvidenceRef::Boundary(BoundaryId::new(1)))
3243 );
3244
3245 Simulation::from_checkpoint_and_journal(
3246 second_checkpoint,
3247 vec![first.segment, second.segment],
3248 )
3249 .expect("reconstruction must filter full segment indexes to the compact receipt frontier");
3250 }
3251
3252 #[test]
3253 fn payload_required_commit_fails_closed_when_an_older_segment_is_missing() {
3254 let mut compact = payload_continuation_runtime();
3255 let complete_archive = TestArchive::default();
3256 let first = store_prepared(&compact, &complete_archive);
3257 compact
3258 .commit_evidence_seal(&first.token, &complete_archive)
3259 .expect("the initial provider-backed seal should commit");
3260
3261 compact
3262 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
3263 .expect("a new tail should preserve the active continuation");
3264 let incomplete_archive = TestArchive::default();
3265 let second = store_prepared(&compact, &incomplete_archive);
3266 let before = compact.checkpoint().expect("checkpoint should build");
3267 let error = compact
3268 .commit_evidence_seal(&second.token, &incomplete_archive)
3269 .expect_err("the provider must retain every payload-required segment");
3270 assert_eq!(error.code, ErrorCode::EvidenceContentUnavailable);
3271 assert_eq!(
3272 compact.checkpoint().expect("checkpoint should build"),
3273 before
3274 );
3275
3276 incomplete_archive
3277 .store_evidence_segment(&first.segment)
3278 .expect("restoring the older required segment should succeed");
3279 compact
3280 .commit_evidence_seal(&second.token, &incomplete_archive)
3281 .expect("the exact provider set should permit commit");
3282 }
3283
3284 #[test]
3285 fn host_orphan_candidates_follow_all_retained_manifests_without_deleting() {
3286 let mut compact = payload_continuation_runtime();
3287 let archive = TestArchive::default();
3288 let prepared = store_prepared(&compact, &archive);
3289 let stored_ids = archive.segment_ids();
3290 let before = compact.checkpoint().expect("checkpoint should build");
3291 assert_eq!(
3292 SimulationCheckpoint::orphaned_archive_segment_ids(
3293 std::slice::from_ref(&before),
3294 &stored_ids,
3295 )
3296 .expect("manifest reachability should validate"),
3297 vec![prepared.token.segment_id.clone()]
3298 );
3299 assert!(
3300 archive
3301 .load_evidence_segment(&prepared.token.segment_id)
3302 .expect("the store should remain readable")
3303 .is_some(),
3304 "the conformance API must not delete host content"
3305 );
3306
3307 compact
3308 .commit_evidence_seal(&prepared.token, &archive)
3309 .expect("the stored segment should commit");
3310 let after = compact.checkpoint().expect("checkpoint should build");
3311 assert!(
3312 SimulationCheckpoint::orphaned_archive_segment_ids(
3313 std::slice::from_ref(&after),
3314 &stored_ids,
3315 )
3316 .expect("committed manifest reachability should validate")
3317 .is_empty()
3318 );
3319 assert_eq!(
3320 SimulationCheckpoint::reachable_archive_segment_ids(&[before, after])
3321 .expect("all retained manifests should be scanned"),
3322 BTreeSet::from([prepared.token.segment_id.clone()])
3323 );
3324 }
3325}