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