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, DecisionArchiveBlob,
6 DecisionArchiveBucketPage, DecisionArchiveProvider, DecisionHistoryKey,
7 DecisionHistoryLocation, DecisionState, DeterministicRng, DomainRecord, DomainRecordPageRoots,
8 DomainRecordRef, DomainRecordSchema, DomainRecordType, DomainRecordVersionRef,
9 DomainRecordVersionSource, ENGINE_VERSION, ErrorCode, EvidenceRef, IngressPayload,
10 IngressReceipt, IngressRecord, KeyedDrawReservation, KnowledgeSnapshot, OutboxEntry,
11 PayloadProperty, PayloadSchema, PayloadValueType, PluginComponentRecord, PluginDescriptor,
12 PluginIngressRequest, PreparedDecisionArchive, PreparedStateDelta, RandomDrawAddress,
13 RandomDrawRecord, RandomStreamState, RunConfigurationSnapshot, RunManifest, RuntimeEvidence,
14 SNAPSHOT_FORMAT_VERSION, STATE_REVISION_FORMAT_VERSION, Scenario, ScheduledAction,
15 ScheduledRecord, SchemaRegistry, SimDuration, SimEvent, SimTime, Simulation, SimulationPlugin,
16 StatePageBlob, StatePageProvider, StatePageRetentionLedger, StatePageStore, SystemCadence,
17 TypedDomainRecordRef, WorldSnapshot, canonical_byte_hash, has_unqueued_command_history,
18 invalid_snapshot_error, is_canonical_hash, is_one_u64, is_zero_u32, is_zero_u64, one_u64,
19 prepare_state_delta, verify_state_delta,
20};
21use serde::{Deserialize, Serialize};
22use std::cell::RefCell;
23use std::collections::{BTreeMap, BTreeSet};
24
25#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
26#[serde(deny_unknown_fields)]
27pub struct SimulationSnapshot {
28 pub engine_version: String,
29 pub snapshot_format_version: u32,
30 #[serde(default)]
31 pub run_manifest: Option<RunManifest>,
32 #[serde(default)]
33 pub run_manifest_hash: String,
34 #[serde(default, skip_serializing_if = "Option::is_none")]
35 pub run_configuration: Option<RunConfigurationSnapshot>,
36 #[serde(default)]
37 pub checkpoint_hash: String,
38 #[serde(default, skip_serializing_if = "is_zero_u32")]
39 pub commitment_format_version: u32,
41 #[serde(default, skip_serializing_if = "Option::is_none")]
42 pub commitment_roots: Option<CommitmentRoots>,
44 #[serde(default)]
45 pub revision_format_version: u32,
47 #[serde(default, skip_serializing_if = "is_zero_u64")]
48 pub state_revision: u64,
50 #[serde(default, skip_serializing_if = "is_zero_u32")]
51 pub replay_revision_format_version: u32,
53 #[serde(default, skip_serializing_if = "is_zero_u32")]
54 pub admission_cursor_format_version: u32,
56 #[serde(default, skip_serializing_if = "is_zero_u64")]
57 pub admitted_attempt_count: u64,
59 #[serde(default, skip_serializing_if = "is_zero_u64")]
60 pub admitted_command_count: u64,
62 #[serde(default, skip_serializing_if = "is_zero_u64")]
63 pub admitted_event_count: u64,
65 pub initial_time: SimTime,
66 #[serde(default, skip_serializing_if = "Option::is_none")]
67 pub initial_scenario: Option<Scenario>,
68 pub now: SimTime,
69 pub plugin_registration_closed: bool,
70 #[serde(default, skip_serializing_if = "Vec::is_empty")]
71 pub entities: Vec<super::EntityRef>,
72 #[serde(default, skip_serializing_if = "legacy_world_is_empty")]
73 pub world: WorldSnapshot,
74 pub knowledge: KnowledgeSnapshot,
75 pub events: Vec<SimEvent>,
76 pub commands: Vec<CommandRecord>,
77 #[serde(default, skip_serializing_if = "Vec::is_empty")]
78 pub command_attempts: Vec<CommandAttemptRecord>,
79 #[serde(default, skip_serializing_if = "Vec::is_empty")]
80 pub ingress: Vec<IngressRecord>,
81 #[serde(default)]
82 pub boundaries: Vec<BoundaryRecord>,
83 pub plugin_components: Vec<PluginComponentRecord>,
84 #[serde(default, skip_serializing_if = "Vec::is_empty")]
85 pub domain_records: Vec<DomainRecord>,
86 #[serde(default, skip_serializing_if = "DecisionState::is_empty")]
87 pub decisions: DecisionState,
88 pub plugin_descriptors: Vec<PluginDescriptor>,
89 pub schema: SchemaRegistry,
90 #[serde(default)]
91 pub root_seed: u64,
92 #[serde(default)]
94 pub authority_root_seed: u64,
95 #[serde(default)]
96 pub random_streams: Vec<RandomStreamState>,
97 #[serde(default)]
98 pub random_draws: Vec<RandomDrawRecord>,
99 pub(super) scheduled: Vec<ScheduledRecord>,
100 #[serde(default, rename = "rng", skip_serializing_if = "Option::is_none")]
101 pub(super) legacy_rng: Option<DeterministicRng>,
102 pub(super) next_event_id: u64,
103 pub(super) next_command_id: u64,
104 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
105 pub(super) next_command_attempt_id: u64,
106 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
107 pub(super) next_ingress_id: u64,
108 #[serde(default)]
109 pub(super) next_boundary_id: u64,
110 #[serde(default)]
111 pub(super) next_random_draw_id: u64,
112 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
113 pub(super) next_knowledge_record_id: u64,
114 pub(super) next_schedule_sequence: u64,
115 pub(super) next_correlation_id: u64,
116 #[serde(default = "one_u64", skip_serializing_if = "is_one_u64")]
117 pub(super) next_decision_trace_id: u64,
118}
119
120pub const PAGED_CHECKPOINT_FORMAT_VERSION: u32 = 4;
121const MAX_PAGED_DECISION_DIRECTORY_ENTRIES: usize = 1_024;
122
123#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
124#[serde(deny_unknown_fields)]
125pub struct PagedSimulationCheckpoint {
126 pub format_version: u32,
127 pub root_page_id: String,
128 pub checkpoint_hash: String,
129}
130
131#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
132#[serde(deny_unknown_fields)]
133pub struct PreparedPagedSimulationCheckpoint {
134 pub checkpoint: PagedSimulationCheckpoint,
135 pub delta: PreparedStateDelta,
136}
137
138#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
139#[serde(deny_unknown_fields)]
140pub struct PortablePagedSimulationCheckpoint {
141 pub checkpoint: PagedSimulationCheckpoint,
142 pub pages: Vec<StatePageBlob>,
143}
144
145#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
148#[serde(deny_unknown_fields)]
149pub struct ArchiveReachabilityManifest {
150 pub state_page_ids: BTreeSet<String>,
151 pub evidence_segment_ids: BTreeSet<String>,
152 pub decision_blob_ids: BTreeSet<String>,
153 pub plugin_objects: BTreeMap<String, BTreeSet<String>>,
154}
155
156impl ArchiveReachabilityManifest {
157 pub fn insert_plugin_object(&mut self, namespace: impl Into<String>, id: impl Into<String>) {
158 self.plugin_objects
159 .entry(namespace.into())
160 .or_default()
161 .insert(id.into());
162 }
163
164 pub fn merge(&mut self, other: Self) {
165 self.state_page_ids.extend(other.state_page_ids);
166 self.evidence_segment_ids.extend(other.evidence_segment_ids);
167 self.decision_blob_ids.extend(other.decision_blob_ids);
168 for (namespace, ids) in other.plugin_objects {
169 self.plugin_objects
170 .entry(namespace)
171 .or_default()
172 .extend(ids);
173 }
174 }
175}
176
177#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
178#[serde(deny_unknown_fields)]
179struct PagedCheckpointEnvelope {
180 format_version: u32,
181 checkpoint_without_paged_state: SimulationCheckpoint,
182 domain_records: DomainRecordPageRoots,
183 decision_manifest_page_id: String,
184}
185
186#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
187#[serde(deny_unknown_fields)]
188struct PagedDecisionManifest {
189 format_version: u32,
190 hot_page_id: String,
191 archive_receipt_root: String,
192 archive_receipt_count: u64,
193 archive_directory_page_ids: Vec<String>,
194}
195
196#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
197#[serde(deny_unknown_fields)]
198struct PagedDecisionDirectoryPage {
199 format_version: u32,
200 archive_bucket_pages: Vec<(super::DecisionArchivePageKey, String)>,
201}
202
203impl PagedDecisionDirectoryPage {
204 fn validate(&self) -> Result<(), CanwuError> {
205 if self.format_version != PAGED_CHECKPOINT_FORMAT_VERSION
206 || self.archive_bucket_pages.is_empty()
207 || self.archive_bucket_pages.len() > MAX_PAGED_DECISION_DIRECTORY_ENTRIES
208 || self
209 .archive_bucket_pages
210 .windows(2)
211 .any(|pair| pair[0].0 >= pair[1].0)
212 {
213 return Err(invalid_snapshot_error(
214 "paged decision directory page is malformed",
215 ));
216 }
217 for (_, page_id) in &self.archive_bucket_pages {
218 if !is_canonical_hash(page_id) {
219 return Err(invalid_snapshot_error(
220 "paged decision archive bucket page ID is malformed",
221 ));
222 }
223 }
224 Ok(())
225 }
226}
227
228fn validate_paged_decision_manifest(manifest: &PagedDecisionManifest) -> Result<(), CanwuError> {
229 let unique_directory_ids = manifest
230 .archive_directory_page_ids
231 .iter()
232 .collect::<BTreeSet<_>>();
233 if manifest.format_version != PAGED_CHECKPOINT_FORMAT_VERSION
234 || !is_canonical_hash(&manifest.hot_page_id)
235 || !is_canonical_hash(&manifest.archive_receipt_root)
236 || manifest
237 .archive_directory_page_ids
238 .iter()
239 .any(|page_id| !is_canonical_hash(page_id))
240 || unique_directory_ids.len() != manifest.archive_directory_page_ids.len()
241 || (manifest.archive_receipt_count == 0) != manifest.archive_directory_page_ids.is_empty()
242 {
243 return Err(invalid_snapshot_error(
244 "paged decision manifest directory is malformed",
245 ));
246 }
247 Ok(())
248}
249
250fn assemble_paged_decision_directory(
251 pages: Vec<PagedDecisionDirectoryPage>,
252) -> Result<BTreeMap<super::DecisionArchivePageKey, String>, CanwuError> {
253 let page_count = pages.len();
254 let mut archive_bucket_pages = BTreeMap::new();
255 let mut previous_page_key = None;
256 for (directory_ordinal, directory_page) in pages.into_iter().enumerate() {
257 directory_page.validate()?;
258 let expected_len = if directory_ordinal + 1 == page_count {
259 1..=MAX_PAGED_DECISION_DIRECTORY_ENTRIES
260 } else {
261 MAX_PAGED_DECISION_DIRECTORY_ENTRIES..=MAX_PAGED_DECISION_DIRECTORY_ENTRIES
262 };
263 if !expected_len.contains(&directory_page.archive_bucket_pages.len()) {
264 return Err(invalid_snapshot_error(
265 "paged decision directory uses noncanonical page chunking",
266 ));
267 }
268 for (page_key, page_id) in directory_page.archive_bucket_pages {
269 if previous_page_key.is_some_and(|previous| previous >= page_key) {
270 return Err(invalid_snapshot_error(
271 "paged decision directory is not globally strictly ordered",
272 ));
273 }
274 previous_page_key = Some(page_key);
275 if archive_bucket_pages.insert(page_key, page_id).is_some() {
276 return Err(invalid_snapshot_error(
277 "paged decision directory contains a duplicate bucket key",
278 ));
279 }
280 }
281 }
282 Ok(archive_bucket_pages)
283}
284
285impl PreparedPagedSimulationCheckpoint {
286 pub fn store_and_verify(&self, store: &dyn StatePageStore) -> Result<(), CanwuError> {
287 if self.checkpoint.format_version != PAGED_CHECKPOINT_FORMAT_VERSION
288 || self.delta.target_root != self.checkpoint.root_page_id
289 {
290 return Err(invalid_snapshot_error(
291 "paged checkpoint descriptor disagrees with its prepared delta",
292 ));
293 }
294 for page in &self.delta.new_pages {
295 let _ = store.store_state_page(page)?;
296 }
297 verify_state_delta(&self.delta, store)?;
298 let envelope = store
299 .load_state_page(&self.checkpoint.root_page_id)?
300 .ok_or_else(|| {
301 CanwuError::new(
302 ErrorCode::StatePageUnavailable,
303 "paged checkpoint root is unavailable after storage",
304 )
305 })?;
306 envelope.validate()?;
307 Ok(())
308 }
309}
310
311#[derive(Default)]
312struct EmbeddedPageProvider {
313 pages: BTreeMap<String, StatePageBlob>,
314}
315
316impl StatePageProvider for EmbeddedPageProvider {
317 fn load_state_page(&self, page_id: &str) -> Result<Option<StatePageBlob>, CanwuError> {
318 Ok(self.pages.get(page_id).cloned())
319 }
320}
321
322#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
323#[serde(deny_unknown_fields)]
324pub struct PagedCheckpointScaleMetrics {
325 pub decision_entries: u64,
326 pub decision_locator: super::DecisionLocatorScaleMetrics,
327 pub state_pages: u64,
328 pub decision_directory_pages: u64,
329 pub max_state_page_bytes: u64,
330 pub initial_delta_pages: u64,
331 pub repeat_delta_pages: u64,
332 pub single_page_change_delta_pages: u64,
333 pub initial_provider_calls: u64,
334 pub repeat_provider_calls: u64,
335 pub single_page_change_provider_calls: u64,
336 pub exact_restart_queries: u64,
337 pub restored_root_matches: bool,
338 pub replayed_root_matches: bool,
339 pub root_page_id: String,
340}
341
342struct Format8ScalePageStore {
343 pages: RefCell<BTreeMap<String, StatePageBlob>>,
344 decision_blobs: RefCell<BTreeMap<String, DecisionArchiveBlob>>,
345 provider_calls: RefCell<u64>,
346}
347
348impl StatePageProvider for Format8ScalePageStore {
349 fn load_state_page(&self, page_id: &str) -> Result<Option<StatePageBlob>, CanwuError> {
350 *self.provider_calls.borrow_mut() += 1;
351 Ok(self.pages.borrow().get(page_id).cloned())
352 }
353}
354
355impl StatePageStore for Format8ScalePageStore {
356 fn store_state_page(&self, page: &StatePageBlob) -> Result<ArchiveStoreOutcome, CanwuError> {
357 page.validate()?;
358 let mut pages = self.pages.borrow_mut();
359 if let Some(existing) = pages.get(&page.page_id) {
360 if existing != page {
361 return Err(invalid_snapshot_error(
362 "Format-8 scale store page ID contains different bytes",
363 ));
364 }
365 return Ok(ArchiveStoreOutcome::AlreadyPresent);
366 }
367 pages.insert(page.page_id.clone(), page.clone());
368 Ok(ArchiveStoreOutcome::Stored)
369 }
370}
371
372impl DecisionArchiveProvider for Format8ScalePageStore {
373 fn load_decision_archive(
374 &self,
375 locator: &str,
376 ) -> Result<Option<DecisionArchiveBlob>, super::DecisionError> {
377 Ok(self.decision_blobs.borrow().get(locator).cloned())
378 }
379
380 fn load_decision_archive_bucket_page(
381 &self,
382 page_id: &str,
383 ) -> Result<Option<DecisionArchiveBucketPage>, super::DecisionError> {
384 let Some(page) = self.pages.borrow().get(page_id).cloned() else {
385 return Ok(None);
386 };
387 serde_json::from_slice(&page.bytes)
388 .map(Some)
389 .map_err(|error| {
390 super::DecisionError::new(
391 super::DecisionErrorCode::DecisionHistoryUnavailable,
392 format!("cannot decode decision bucket state page: {error}"),
393 )
394 })
395 }
396}
397
398fn legacy_world_is_empty(world: &WorldSnapshot) -> bool {
399 world.people.is_empty()
400 && world.governments.is_empty()
401 && world.territories.is_empty()
402 && world.routes.is_empty()
403 && world.armies.is_empty()
404 && world.letters.is_empty()
405}
406
407#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
408#[serde(deny_unknown_fields)]
409pub struct ReplayJournal {
411 pub engine_version: String,
412 pub snapshot_format_version: u32,
413 pub root_seed: u64,
414 pub initial_scenario: Scenario,
417 pub authority_root_seed: u64,
418 pub run_manifest: RunManifest,
419 pub run_manifest_hash: String,
420 pub run_configuration: RunConfigurationSnapshot,
421 pub plugin_descriptors: Vec<PluginDescriptor>,
422 pub plugin_registration_closed: bool,
423 pub commands: Vec<CommandRecord>,
424 pub command_attempts: Vec<CommandAttemptRecord>,
425 #[serde(default, skip_serializing_if = "Vec::is_empty")]
426 pub ingress: Vec<IngressRecord>,
427 pub boundaries: Vec<BoundaryRecord>,
428 pub final_time: SimTime,
429 pub checkpoint_hash: String,
430 pub commitment_format_version: u32,
432 pub revision_format_version: u32,
434 pub final_revision: u64,
436}
437
438pub const CHECKPOINT_JOURNAL_FORMAT_VERSION: u32 = 4;
440const ARCHIVED_SEGMENT_MANIFEST_DOMAIN: &str = "canwu.evidence.archived-segment-manifest.v1";
441const ARCHIVED_RECEIPT_DOMAIN: &str = "canwu.evidence.archived-receipts.v2";
442const EVIDENCE_DEPENDENCY_DOMAIN: &str = "canwu.evidence.dependencies.v1";
443const KEYED_RESERVATION_DOMAIN: &str = "canwu.random.keyed-reservations.v1";
444pub const PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD: &str =
446 "canwu_payload_required_evidence_continuation";
447pub const PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION: u32 = 1;
449pub const IDENTITY_EVIDENCE_DEPENDENCIES_FIELD: &str = "canwu_identity_evidence_dependencies";
451pub const IDENTITY_EVIDENCE_DEPENDENCIES_FORMAT_VERSION: u32 = 1;
453
454#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
455#[serde(rename_all = "snake_case")]
456pub enum EvidenceJournalKind {
457 Event,
458 Command,
459 CommandAttempt,
460 Ingress,
461 Boundary,
462 RandomDraw,
463}
464
465#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
466#[serde(tag = "type", rename_all = "snake_case")]
467pub enum EvidenceNestedLocator {
468 None,
469 BoundaryRecordChange { change_index: u64 },
470}
471
472#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
473pub struct EvidenceItemLocator {
474 pub journal: EvidenceJournalKind,
475 pub absolute_index: u64,
476 pub nested: EvidenceNestedLocator,
477}
478
479#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
480pub struct ArchivedEvidenceLocator {
481 pub segment_id: String,
482 pub item: EvidenceItemLocator,
483}
484
485#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
486pub struct ArchivedEvidenceReceipt {
487 pub evidence: EvidenceRef,
488 pub locator: ArchivedEvidenceLocator,
489 pub evidence_index_leaf: u64,
490 pub item_commitment: String,
491 #[serde(default, skip_serializing_if = "Option::is_none")]
492 pub plugin_ingress_provenance: Option<ArchivedPluginIngressProvenance>,
493 pub merkle_path: Vec<String>,
494}
495
496#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
497pub struct ArchivedPluginIngressProvenance {
498 pub plugin: String,
499 pub packet_type: String,
500 pub producer_boundary: super::BoundaryId,
501}
502
503#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
504pub struct EvidenceIndexEntry {
505 pub reference: EvidenceRef,
506 pub item: EvidenceItemLocator,
507 pub item_commitment: String,
508 #[serde(default, skip_serializing_if = "Option::is_none")]
509 pub plugin_ingress_provenance: Option<ArchivedPluginIngressProvenance>,
510}
511
512#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
513pub struct EvidenceJournalRoots {
514 pub events: String,
515 pub commands: String,
516 pub command_attempts: String,
517 pub ingress: String,
518 pub boundaries: String,
519 pub random_draws: String,
520}
521
522#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
523pub struct ArchivedSegmentHeader {
524 pub segment_id: String,
525 pub start: EvidenceCursor,
526 pub end: EvidenceCursor,
527 pub journal_roots: EvidenceJournalRoots,
528 pub evidence_index_root: String,
529 pub evidence_index_entry_count: u64,
530}
531
532#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
533#[serde(rename_all = "snake_case")]
534pub enum EvidenceRequirement {
535 IdentityOnly,
537 PayloadRequired,
539}
540
541#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
549#[serde(deny_unknown_fields)]
550pub struct PayloadRequiredEvidenceContinuationV1 {
551 pub format_version: u32,
552 pub active: bool,
553 pub dependencies: Vec<EvidenceRef>,
554}
555
556impl PayloadRequiredEvidenceContinuationV1 {
557 #[must_use]
559 pub fn active(dependencies: Vec<EvidenceRef>) -> Self {
560 Self {
561 format_version: PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION,
562 active: true,
563 dependencies,
564 }
565 }
566
567 #[must_use]
569 pub const fn completed() -> Self {
570 Self {
571 format_version: PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION,
572 active: false,
573 dependencies: Vec::new(),
574 }
575 }
576}
577
578#[must_use]
581pub fn payload_required_evidence_continuation_property_v1() -> PayloadProperty {
582 PayloadProperty {
583 value_type: PayloadValueType::Object,
584 required: true,
585 }
586}
587
588#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
594#[serde(deny_unknown_fields)]
595pub struct IdentityEvidenceDependenciesV1 {
596 pub format_version: u32,
597 pub dependencies: Vec<EvidenceRef>,
598}
599
600impl IdentityEvidenceDependenciesV1 {
601 #[must_use]
603 pub const fn new(dependencies: Vec<EvidenceRef>) -> Self {
604 Self {
605 format_version: IDENTITY_EVIDENCE_DEPENDENCIES_FORMAT_VERSION,
606 dependencies,
607 }
608 }
609}
610
611#[must_use]
614pub fn identity_evidence_dependencies_property_v1() -> PayloadProperty {
615 PayloadProperty {
616 value_type: PayloadValueType::Object,
617 required: true,
618 }
619}
620#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
621pub struct EvidenceDependency {
622 pub reference: EvidenceRef,
623 pub requirement: EvidenceRequirement,
624}
625
626#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
627pub struct EvidenceArchiveIndex {
628 pub header: ArchivedSegmentHeader,
629 pub entries: Vec<EvidenceIndexEntry>,
630}
631
632pub trait ArchiveProvider {
633 fn load_evidence_segment(
634 &self,
635 segment_id: &str,
636 ) -> Result<Option<EvidenceJournalSegment>, CanwuError>;
637}
638
639pub trait ArchiveStore: ArchiveProvider {
640 fn store_evidence_segment(
641 &self,
642 segment: &EvidenceJournalSegment,
643 ) -> Result<ArchiveStoreOutcome, CanwuError>;
644}
645
646#[derive(Clone, Copy, Debug, Eq, PartialEq)]
647pub enum ArchiveStoreOutcome {
648 Stored,
649 AlreadyPresent,
650}
651
652#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
653pub struct EvidenceSealToken {
654 pub source_state_hash: String,
655 pub source_checkpoint_hash: String,
656 pub source_end: EvidenceCursor,
657 pub segment_id: String,
658 pub target_checkpoint_hash: String,
659 pub token_hash: String,
660}
661
662#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
663pub struct PreparedEvidenceSeal {
664 pub token: EvidenceSealToken,
665 pub segment: EvidenceJournalSegment,
666}
667
668#[derive(Serialize)]
669struct EvidenceIndexLeafMaterial<'a> {
670 format_version: u32,
671 reference: &'a EvidenceRef,
672 item: &'a EvidenceItemLocator,
673 item_commitment: &'a str,
674 plugin_ingress_provenance: &'a Option<ArchivedPluginIngressProvenance>,
675}
676
677#[derive(Serialize)]
678struct ArchivedSegmentHeaderMaterial<'a> {
679 start: EvidenceCursor,
680 end: EvidenceCursor,
681 journal_roots: &'a EvidenceJournalRoots,
682 evidence_index_root: &'a str,
683 evidence_index_entry_count: u64,
684}
685
686fn archive_error(message: impl Into<String>) -> CanwuError {
687 CanwuError::new(ErrorCode::InvalidArchive, message)
688}
689
690fn decode_hash(value: &str, label: &str) -> Result<[u8; 32], CanwuError> {
691 if value.len() != 64
692 || value
693 .bytes()
694 .any(|byte| !byte.is_ascii_hexdigit() || byte.is_ascii_uppercase())
695 {
696 return Err(archive_error(format!(
697 "{label} must be 32-byte lower-case hex"
698 )));
699 }
700 let mut bytes = [0_u8; 32];
701 for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() {
702 let digit = |byte: u8| match byte {
703 b'0'..=b'9' => Some(byte - b'0'),
704 b'a'..=b'f' => Some(byte - b'a' + 10),
705 _ => None,
706 };
707 bytes[index] = digit(pair[0])
708 .and_then(|high| digit(pair[1]).map(|low| (high << 4) | low))
709 .ok_or_else(|| archive_error(format!("{label} contains invalid hex")))?;
710 }
711 Ok(bytes)
712}
713
714fn archive_node(left: [u8; 32], right: [u8; 32]) -> [u8; 32] {
715 let mut hasher = blake3::Hasher::new();
716 hasher.update(b"canwu.evidence.index.node.v1");
717 hasher.update(&[0]);
718 hasher.update(&left);
719 hasher.update(&right);
720 *hasher.finalize().as_bytes()
721}
722
723fn archive_empty_root() -> String {
724 let mut hasher = blake3::Hasher::new();
725 hasher.update(b"canwu.evidence.index.empty.v1");
726 hasher.update(&[0]);
727 hasher.finalize().to_hex().to_string()
728}
729
730fn skipped_commitment_root<T: Serialize>(
731 domain: &str,
732 values: &[T],
733) -> Result<Option<String>, CanwuError> {
734 if values.is_empty() {
735 Ok(None)
736 } else {
737 super::canonical_hash(domain, values).map(Some)
738 }
739}
740
741fn validate_skipped_commitment_root<T: Serialize>(
742 root: Option<&str>,
743 domain: &str,
744 values: &[T],
745 label: &str,
746) -> Result<(), CanwuError> {
747 let expected = skipped_commitment_root(domain, values)?;
748 if root != expected.as_deref() {
749 return Err(invalid_snapshot_error(format!(
750 "compact continuation {label} does not match its canonical material"
751 )));
752 }
753 Ok(())
754}
755
756fn promote_dependency(
757 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
758 reference: EvidenceRef,
759 requirement: EvidenceRequirement,
760) {
761 dependencies
762 .entry(reference)
763 .and_modify(|current| *current = (*current).max(requirement))
764 .or_insert(requirement);
765}
766fn schema_declares_payload_required_continuation(schema: &DomainRecordSchema) -> bool {
767 matches!(
768 &schema.payload_schema,
769 PayloadSchema::Object { properties, .. }
770 if properties.get(PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD)
771 == Some(&payload_required_evidence_continuation_property_v1())
772 )
773}
774
775fn schema_declares_identity_evidence_dependencies(schema: &DomainRecordSchema) -> bool {
776 matches!(
777 &schema.payload_schema,
778 PayloadSchema::Object { properties, .. }
779 if properties.get(IDENTITY_EVIDENCE_DEPENDENCIES_FIELD)
780 == Some(&identity_evidence_dependencies_property_v1())
781 )
782}
783
784fn add_identity_evidence_dependencies(
785 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
786 record: &DomainRecord,
787 schema: &DomainRecordSchema,
788) -> Result<(), CanwuError> {
789 if !record.is_active() || !schema_declares_identity_evidence_dependencies(schema) {
790 return Ok(());
791 }
792 let declaration = record
793 .payload
794 .get(IDENTITY_EVIDENCE_DEPENDENCIES_FIELD)
795 .ok_or_else(|| {
796 CanwuError::new(
797 ErrorCode::ArchiveNotReady,
798 format!(
799 "identity-evidence record {} is missing its schema-declared field",
800 record.reference
801 ),
802 )
803 })?;
804 let declaration: IdentityEvidenceDependenciesV1 = serde_json::from_value(declaration.clone())
805 .map_err(|error| {
806 CanwuError::new(
807 ErrorCode::ArchiveNotReady,
808 format!(
809 "identity-evidence record {} is invalid: {error}",
810 record.reference
811 ),
812 )
813 })?;
814 if declaration.format_version != IDENTITY_EVIDENCE_DEPENDENCIES_FORMAT_VERSION {
815 return Err(CanwuError::new(
816 ErrorCode::ArchiveNotReady,
817 format!(
818 "identity-evidence record {} uses unsupported format {}",
819 record.reference, declaration.format_version
820 ),
821 ));
822 }
823 if declaration
824 .dependencies
825 .windows(2)
826 .any(|pair| pair[0] >= pair[1])
827 {
828 return Err(CanwuError::new(
829 ErrorCode::ArchiveNotReady,
830 format!(
831 "identity-evidence record {} needs sorted unique dependencies",
832 record.reference
833 ),
834 ));
835 }
836 if declaration.dependencies.iter().any(|reference| {
837 matches!(
838 reference,
839 EvidenceRef::DomainRecordVersion(version)
840 if version.record == record.reference && version.version == record.version
841 )
842 }) {
843 return Err(CanwuError::new(
844 ErrorCode::ArchiveNotReady,
845 format!(
846 "identity-evidence record {} cannot depend on its own version",
847 record.reference
848 ),
849 ));
850 }
851 for reference in declaration.dependencies {
852 if matches!(
853 &reference,
854 EvidenceRef::DomainRecordVersion(version)
855 if matches!(version.established_by, DomainRecordVersionSource::InitialScenario)
856 ) {
857 return Err(CanwuError::new(
858 ErrorCode::ArchiveNotReady,
859 format!(
860 "identity-evidence record {} cannot archive initial-scenario evidence",
861 record.reference
862 ),
863 ));
864 }
865 promote_dependency(dependencies, reference, EvidenceRequirement::IdentityOnly);
866 }
867 Ok(())
868}
869
870fn add_payload_required_continuation_dependencies(
871 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
872 record: &DomainRecord,
873 schema: &DomainRecordSchema,
874) -> Result<(), CanwuError> {
875 if !record.is_active() || !schema_declares_payload_required_continuation(schema) {
876 return Ok(());
877 }
878 let continuation = record
879 .payload
880 .get(PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD)
881 .ok_or_else(|| {
882 CanwuError::new(
883 ErrorCode::ArchiveNotReady,
884 format!(
885 "payload-required continuation record {} is missing its schema-declared field",
886 record.reference
887 ),
888 )
889 })?;
890 let continuation: PayloadRequiredEvidenceContinuationV1 =
891 serde_json::from_value(continuation.clone()).map_err(|error| {
892 CanwuError::new(
893 ErrorCode::ArchiveNotReady,
894 format!(
895 "payload-required continuation record {} is invalid: {error}",
896 record.reference
897 ),
898 )
899 })?;
900 if continuation.format_version != PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION {
901 return Err(CanwuError::new(
902 ErrorCode::ArchiveNotReady,
903 format!(
904 "payload-required continuation record {} uses unsupported format {}",
905 record.reference, continuation.format_version
906 ),
907 ));
908 }
909 if !continuation.active {
910 if continuation == PayloadRequiredEvidenceContinuationV1::completed() {
911 return Ok(());
912 }
913 return Err(CanwuError::new(
914 ErrorCode::ArchiveNotReady,
915 format!(
916 "completed payload-required continuation record {} retains dependencies",
917 record.reference
918 ),
919 ));
920 }
921 let continuation = PayloadRequiredEvidenceContinuationV1::active(continuation.dependencies);
922 if continuation.dependencies.is_empty()
923 || continuation
924 .dependencies
925 .windows(2)
926 .any(|pair| pair[0] >= pair[1])
927 {
928 return Err(CanwuError::new(
929 ErrorCode::ArchiveNotReady,
930 format!(
931 "active payload-required continuation record {} needs sorted unique dependencies",
932 record.reference
933 ),
934 ));
935 }
936 if continuation.dependencies.iter().any(|reference| {
937 matches!(
938 reference,
939 EvidenceRef::DomainRecordVersion(version)
940 if version.record == record.reference && version.version == record.version
941 )
942 }) {
943 return Err(CanwuError::new(
944 ErrorCode::ArchiveNotReady,
945 format!(
946 "payload-required continuation record {} cannot depend on its own version",
947 record.reference
948 ),
949 ));
950 }
951 for reference in continuation.dependencies {
952 if matches!(
953 &reference,
954 EvidenceRef::DomainRecordVersion(version)
955 if matches!(version.established_by, DomainRecordVersionSource::InitialScenario)
956 ) {
957 return Err(CanwuError::new(
958 ErrorCode::ArchiveNotReady,
959 format!(
960 "payload-required continuation record {} cannot depend on an initial-scenario payload",
961 record.reference
962 ),
963 ));
964 }
965 promote_dependency(
966 dependencies,
967 reference,
968 EvidenceRequirement::PayloadRequired,
969 );
970 }
971 Ok(())
972}
973
974fn required_archived_receipt_references(
975 dependencies: &[EvidenceDependency],
976 reservations: &[KeyedDrawReservation],
977) -> BTreeSet<EvidenceRef> {
978 let mut required = dependencies
979 .iter()
980 .map(|dependency| dependency.reference.clone())
981 .collect::<BTreeSet<_>>();
982 required.extend(
983 reservations
984 .iter()
985 .map(|reservation| reservation.draw_receipt.evidence.clone()),
986 );
987 required
988}
989
990fn retain_reachable_archived_evidence_receipts(
991 receipts: &mut BTreeMap<EvidenceRef, ArchivedEvidenceReceipt>,
992 dependencies: &[EvidenceDependency],
993 reservations: &[KeyedDrawReservation],
994) {
995 let required = required_archived_receipt_references(dependencies, reservations);
996 receipts.retain(|reference, _| required.contains(reference));
997}
998
999pub(crate) fn load_verified_archived_evidence_segment(
1000 receipt: &ArchivedEvidenceReceipt,
1001 provider: &dyn ArchiveProvider,
1002) -> Result<EvidenceJournalSegment, CanwuError> {
1003 let segment = provider
1004 .load_evidence_segment(&receipt.locator.segment_id)?
1005 .ok_or_else(|| {
1006 CanwuError::new(
1007 ErrorCode::EvidenceContentUnavailable,
1008 "the archive provider did not return the required evidence segment",
1009 )
1010 })?;
1011 let receipts = verify_archived_segment(&segment).map_err(|error| {
1012 CanwuError::new(
1013 ErrorCode::InvalidArchive,
1014 format!(
1015 "archive provider returned invalid content: {}",
1016 error.message
1017 ),
1018 )
1019 })?;
1020 if !receipts.iter().any(|candidate| candidate == receipt) {
1021 return Err(archive_error(
1022 "archive provider segment does not reproduce the committed evidence receipt",
1023 ));
1024 }
1025 Ok(segment)
1026}
1027
1028fn add_cause_dependency(
1029 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
1030 cause: &CauseRef,
1031) {
1032 let reference = match cause {
1033 CauseRef::Event(id) => Some(EvidenceRef::Event(*id)),
1034 CauseRef::Command(id) => Some(EvidenceRef::Command(*id)),
1035 CauseRef::Boundary(id) => Some(EvidenceRef::Boundary(*id)),
1036 CauseRef::System(_) => None,
1037 };
1038 if let Some(reference) = reference {
1039 promote_dependency(dependencies, reference, EvidenceRequirement::IdentityOnly);
1040 }
1041}
1042
1043fn add_command_outcome_dependencies(
1044 dependencies: &mut BTreeMap<EvidenceRef, EvidenceRequirement>,
1045 outcome: &CommandOutcome,
1046) {
1047 let (attempt_id, command_id, emitted_events) = match outcome {
1048 CommandOutcome::Accepted { receipt } => (
1049 receipt.attempt_id,
1050 Some(receipt.command_id),
1051 receipt.emitted_events.as_slice(),
1052 ),
1053 CommandOutcome::Rejected { rejection } => (rejection.attempt_id, None, &[][..]),
1054 };
1055 if let Some(id) = attempt_id {
1056 promote_dependency(
1057 dependencies,
1058 EvidenceRef::CommandAttempt(id),
1059 EvidenceRequirement::IdentityOnly,
1060 );
1061 }
1062 if let Some(id) = command_id {
1063 promote_dependency(
1064 dependencies,
1065 EvidenceRef::Command(id),
1066 EvidenceRequirement::IdentityOnly,
1067 );
1068 }
1069 for id in emitted_events {
1070 promote_dependency(
1071 dependencies,
1072 EvidenceRef::Event(*id),
1073 EvidenceRequirement::IdentityOnly,
1074 );
1075 }
1076}
1077
1078fn archive_leaf(entry: &EvidenceIndexEntry) -> Result<[u8; 32], CanwuError> {
1079 let hash = super::canonical_hash(
1080 "canwu.evidence.index.leaf.v2",
1081 &EvidenceIndexLeafMaterial {
1082 format_version: 2,
1083 reference: &entry.reference,
1084 item: &entry.item,
1085 item_commitment: &entry.item_commitment,
1086 plugin_ingress_provenance: &entry.plugin_ingress_provenance,
1087 },
1088 )?;
1089 decode_hash(&hash, "evidence-index leaf")
1090}
1091
1092fn archive_merkle(
1093 entries: &[EvidenceIndexEntry],
1094) -> Result<(String, Vec<Vec<String>>), CanwuError> {
1095 if entries.is_empty() {
1096 return Ok((archive_empty_root(), Vec::new()));
1097 }
1098 let mut level: Vec<[u8; 32]> = entries.iter().map(archive_leaf).collect::<Result<_, _>>()?;
1099 let mut positions: Vec<usize> = (0..entries.len()).collect();
1100 let mut proofs = vec![Vec::new(); entries.len()];
1101 while level.len() > 1 {
1102 for (leaf, position) in positions.iter().copied().enumerate() {
1103 let sibling = if position % 2 == 0 {
1104 (position + 1).min(level.len() - 1)
1105 } else {
1106 position - 1
1107 };
1108 proofs[leaf].push(
1109 blake3::Hash::from_bytes(level[sibling])
1110 .to_hex()
1111 .to_string(),
1112 );
1113 }
1114 let mut next = Vec::with_capacity(level.len().div_ceil(2));
1115 for pair in level.chunks(2) {
1116 next.push(archive_node(pair[0], *pair.get(1).unwrap_or(&pair[0])));
1117 }
1118 level = next;
1119 for position in &mut positions {
1120 *position /= 2;
1121 }
1122 }
1123 Ok((
1124 blake3::Hash::from_bytes(level[0]).to_hex().to_string(),
1125 proofs,
1126 ))
1127}
1128
1129fn item_commitment<T: Serialize>(
1130 journal: EvidenceJournalKind,
1131 item: &T,
1132) -> Result<String, CanwuError> {
1133 let domain = match journal {
1134 EvidenceJournalKind::Event => "canwu.evidence.item.event.v1",
1135 EvidenceJournalKind::Command => "canwu.evidence.item.command.v1",
1136 EvidenceJournalKind::CommandAttempt => "canwu.evidence.item.command_attempt.v1",
1137 EvidenceJournalKind::Ingress => "canwu.evidence.item.ingress.v1",
1138 EvidenceJournalKind::Boundary => "canwu.evidence.item.boundary.v1",
1139 EvidenceJournalKind::RandomDraw => "canwu.evidence.item.random_draw.v1",
1140 };
1141 super::canonical_hash(domain, item)
1142}
1143
1144pub(crate) fn evidence_archive_index(
1145 segment: &EvidenceJournalSegment,
1146) -> Result<(EvidenceArchiveIndex, Vec<ArchivedEvidenceReceipt>), CanwuError> {
1147 let roots = EvidenceJournalRoots {
1148 events: super::canonical_hash("canwu.evidence.journal.events.v1", &segment.events)?,
1149 commands: super::canonical_hash("canwu.evidence.journal.commands.v1", &segment.commands)?,
1150 command_attempts: super::canonical_hash(
1151 "canwu.evidence.journal.command_attempts.v1",
1152 &segment.command_attempts,
1153 )?,
1154 ingress: super::canonical_hash("canwu.evidence.journal.ingress.v1", &segment.ingress)?,
1155 boundaries: super::canonical_hash(
1156 "canwu.evidence.journal.boundaries.v1",
1157 &segment.boundaries,
1158 )?,
1159 random_draws: super::canonical_hash(
1160 "canwu.evidence.journal.random_draws.v1",
1161 &segment.random_draws,
1162 )?,
1163 };
1164 let mut entries = Vec::new();
1165 let mut add = |reference: EvidenceRef,
1166 journal,
1167 absolute_index,
1168 nested,
1169 commitment: String,
1170 plugin_ingress_provenance| {
1171 entries.push(EvidenceIndexEntry {
1172 reference,
1173 item: EvidenceItemLocator {
1174 journal,
1175 absolute_index,
1176 nested,
1177 },
1178 item_commitment: commitment,
1179 plugin_ingress_provenance,
1180 });
1181 };
1182 for (offset, event) in segment.events.iter().enumerate() {
1183 add(
1184 EvidenceRef::Event(event.id),
1185 EvidenceJournalKind::Event,
1186 segment.start.event_count + offset as u64 + 1,
1187 EvidenceNestedLocator::None,
1188 item_commitment(EvidenceJournalKind::Event, event)?,
1189 None,
1190 );
1191 }
1192 for (offset, command) in segment.commands.iter().enumerate() {
1193 add(
1194 EvidenceRef::Command(command.id),
1195 EvidenceJournalKind::Command,
1196 segment.start.command_count + offset as u64 + 1,
1197 EvidenceNestedLocator::None,
1198 item_commitment(EvidenceJournalKind::Command, command)?,
1199 None,
1200 );
1201 }
1202 for (offset, attempt) in segment.command_attempts.iter().enumerate() {
1203 add(
1204 EvidenceRef::CommandAttempt(attempt.id),
1205 EvidenceJournalKind::CommandAttempt,
1206 segment.start.command_attempt_count + offset as u64 + 1,
1207 EvidenceNestedLocator::None,
1208 item_commitment(EvidenceJournalKind::CommandAttempt, attempt)?,
1209 None,
1210 );
1211 }
1212 for (offset, ingress) in segment.ingress.iter().enumerate() {
1213 let plugin_ingress_provenance = match (&ingress.payload, ingress.cause.as_ref()) {
1214 (
1215 IngressPayload::Plugin {
1216 plugin,
1217 packet_type,
1218 ..
1219 },
1220 Some(CauseRef::Boundary(producer_boundary)),
1221 ) if segment.boundaries.iter().any(|boundary| {
1222 boundary.id == *producer_boundary
1223 && boundary.generated_ingress.iter().any(|generation| {
1224 generation.ingress == ingress.id && generation.plugin == *plugin
1225 })
1226 }) =>
1227 {
1228 Some(ArchivedPluginIngressProvenance {
1229 plugin: plugin.clone(),
1230 packet_type: packet_type.clone(),
1231 producer_boundary: *producer_boundary,
1232 })
1233 }
1234 _ => None,
1235 };
1236 add(
1237 EvidenceRef::Ingress(ingress.id),
1238 EvidenceJournalKind::Ingress,
1239 segment.start.ingress_count + offset as u64 + 1,
1240 EvidenceNestedLocator::None,
1241 item_commitment(EvidenceJournalKind::Ingress, ingress)?,
1242 plugin_ingress_provenance,
1243 );
1244 }
1245 for (offset, boundary) in segment.boundaries.iter().enumerate() {
1246 let absolute_index = segment.start.boundary_count + offset as u64 + 1;
1247 let commitment = item_commitment(EvidenceJournalKind::Boundary, boundary)?;
1248 add(
1249 EvidenceRef::Boundary(boundary.id),
1250 EvidenceJournalKind::Boundary,
1251 absolute_index,
1252 EvidenceNestedLocator::None,
1253 commitment.clone(),
1254 None,
1255 );
1256 for (change_index, change) in boundary.record_changes.iter().enumerate() {
1257 add(
1258 EvidenceRef::DomainRecordVersion(DomainRecordVersionRef {
1259 record: change.current.reference.clone(),
1260 version: change.current.version,
1261 established_by: DomainRecordVersionSource::BoundaryChange {
1262 boundary: boundary.id,
1263 change_index: change_index as u64,
1264 },
1265 }),
1266 EvidenceJournalKind::Boundary,
1267 absolute_index,
1268 EvidenceNestedLocator::BoundaryRecordChange {
1269 change_index: change_index as u64,
1270 },
1271 commitment.clone(),
1272 None,
1273 );
1274 }
1275 }
1276 for (offset, draw) in segment.random_draws.iter().enumerate() {
1277 add(
1278 EvidenceRef::RandomDraw(draw.id),
1279 EvidenceJournalKind::RandomDraw,
1280 segment.start.random_draw_count + offset as u64 + 1,
1281 EvidenceNestedLocator::None,
1282 item_commitment(EvidenceJournalKind::RandomDraw, draw)?,
1283 None,
1284 );
1285 }
1286 entries.sort();
1287 if entries
1288 .windows(2)
1289 .any(|window| window[0].reference == window[1].reference)
1290 {
1291 return Err(archive_error(
1292 "evidence archive contains duplicate references",
1293 ));
1294 }
1295 let (evidence_index_root, proofs) = archive_merkle(&entries)?;
1296 let entry_count = u64::try_from(entries.len())
1297 .map_err(|_| archive_error("evidence-index entry count exceeds u64"))?;
1298 let segment_id = super::canonical_hash(
1299 "canwu.evidence.segment.v3",
1300 &ArchivedSegmentHeaderMaterial {
1301 start: segment.start,
1302 end: segment.end,
1303 journal_roots: &roots,
1304 evidence_index_root: &evidence_index_root,
1305 evidence_index_entry_count: entry_count,
1306 },
1307 )?;
1308 let header = ArchivedSegmentHeader {
1309 segment_id: segment_id.clone(),
1310 start: segment.start,
1311 end: segment.end,
1312 journal_roots: roots,
1313 evidence_index_root,
1314 evidence_index_entry_count: entry_count,
1315 };
1316 let receipts = entries
1317 .iter()
1318 .zip(proofs)
1319 .enumerate()
1320 .map(|(index, (entry, merkle_path))| ArchivedEvidenceReceipt {
1321 evidence: entry.reference.clone(),
1322 locator: ArchivedEvidenceLocator {
1323 segment_id: segment_id.clone(),
1324 item: entry.item.clone(),
1325 },
1326 evidence_index_leaf: index as u64,
1327 item_commitment: entry.item_commitment.clone(),
1328 plugin_ingress_provenance: entry.plugin_ingress_provenance.clone(),
1329 merkle_path,
1330 })
1331 .collect();
1332 Ok((EvidenceArchiveIndex { header, entries }, receipts))
1333}
1334
1335fn verify_archive_receipt(
1336 receipt: &ArchivedEvidenceReceipt,
1337 header: &ArchivedSegmentHeader,
1338) -> Result<(), CanwuError> {
1339 if receipt.locator.segment_id != header.segment_id
1340 || receipt.evidence_index_leaf >= header.evidence_index_entry_count
1341 {
1342 return Err(archive_error(
1343 "archived evidence receipt does not belong to its segment header",
1344 ));
1345 }
1346 let entry = EvidenceIndexEntry {
1347 reference: receipt.evidence.clone(),
1348 item: receipt.locator.item.clone(),
1349 item_commitment: receipt.item_commitment.clone(),
1350 plugin_ingress_provenance: receipt.plugin_ingress_provenance.clone(),
1351 };
1352 let mut hash = archive_leaf(&entry)?;
1353 let mut position = receipt.evidence_index_leaf;
1354 let mut width = header.evidence_index_entry_count;
1355 let mut path_at = 0usize;
1356 while width > 1 {
1357 let sibling = receipt
1358 .merkle_path
1359 .get(path_at)
1360 .ok_or_else(|| archive_error("archived evidence receipt Merkle path is too short"))?;
1361 let sibling = decode_hash(sibling, "receipt Merkle sibling")?;
1362 hash = if position.is_multiple_of(2) {
1363 archive_node(hash, sibling)
1364 } else {
1365 archive_node(sibling, hash)
1366 };
1367 position /= 2;
1368 width = width.div_ceil(2);
1369 path_at += 1;
1370 }
1371 if path_at != receipt.merkle_path.len()
1372 || blake3::Hash::from_bytes(hash).to_hex().as_str() != header.evidence_index_root
1373 {
1374 return Err(archive_error(
1375 "archived evidence receipt Merkle proof is invalid",
1376 ));
1377 }
1378 let (expected_journal, expected_nested) = match &receipt.evidence {
1379 EvidenceRef::Event(_) => (EvidenceJournalKind::Event, EvidenceNestedLocator::None),
1380 EvidenceRef::Command(_) => (EvidenceJournalKind::Command, EvidenceNestedLocator::None),
1381 EvidenceRef::CommandAttempt(_) => (
1382 EvidenceJournalKind::CommandAttempt,
1383 EvidenceNestedLocator::None,
1384 ),
1385 EvidenceRef::Ingress(_) => (EvidenceJournalKind::Ingress, EvidenceNestedLocator::None),
1386 EvidenceRef::Boundary(_) => (EvidenceJournalKind::Boundary, EvidenceNestedLocator::None),
1387 EvidenceRef::RandomDraw(_) => {
1388 (EvidenceJournalKind::RandomDraw, EvidenceNestedLocator::None)
1389 }
1390 EvidenceRef::DomainRecordVersion(version) => match version.established_by {
1391 DomainRecordVersionSource::BoundaryChange { change_index, .. } => (
1392 EvidenceJournalKind::Boundary,
1393 EvidenceNestedLocator::BoundaryRecordChange { change_index },
1394 ),
1395 DomainRecordVersionSource::InitialScenario => {
1396 return Err(archive_error(
1397 "initial-scenario record versions cannot be archived",
1398 ));
1399 }
1400 },
1401 };
1402 if receipt.locator.item.journal != expected_journal
1403 || receipt.locator.item.nested != expected_nested
1404 {
1405 return Err(archive_error(
1406 "archived evidence receipt uses an illegal typed locator",
1407 ));
1408 }
1409 let (start, end) = match expected_journal {
1410 EvidenceJournalKind::Event => (header.start.event_count, header.end.event_count),
1411 EvidenceJournalKind::Command => (header.start.command_count, header.end.command_count),
1412 EvidenceJournalKind::CommandAttempt => (
1413 header.start.command_attempt_count,
1414 header.end.command_attempt_count,
1415 ),
1416 EvidenceJournalKind::Ingress => (header.start.ingress_count, header.end.ingress_count),
1417 EvidenceJournalKind::Boundary => (header.start.boundary_count, header.end.boundary_count),
1418 EvidenceJournalKind::RandomDraw => {
1419 (header.start.random_draw_count, header.end.random_draw_count)
1420 }
1421 };
1422 if receipt.locator.item.absolute_index <= start || receipt.locator.item.absolute_index > end {
1423 return Err(archive_error(
1424 "archived evidence locator lies outside its segment cursor range",
1425 ));
1426 }
1427 Ok(())
1428}
1429
1430pub(crate) fn verify_archived_segment(
1431 segment: &EvidenceJournalSegment,
1432) -> Result<Vec<ArchivedEvidenceReceipt>, CanwuError> {
1433 let stored = segment
1434 .archive
1435 .as_ref()
1436 .ok_or_else(|| archive_error("archived segment is missing its evidence index"))?;
1437 let mut plain = segment.clone();
1438 plain.archive = None;
1439 let (rebuilt, receipts) = evidence_archive_index(&plain)?;
1440 if &rebuilt != stored {
1441 return Err(archive_error(
1442 "archived segment header or evidence index does not match its journal content",
1443 ));
1444 }
1445 for receipt in &receipts {
1446 verify_archive_receipt(receipt, &stored.header)?;
1447 }
1448 Ok(receipts)
1449}
1450
1451#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
1452pub struct EvidenceCursor {
1454 pub event_count: u64,
1455 pub command_count: u64,
1456 pub command_attempt_count: u64,
1457 pub ingress_count: u64,
1458 pub boundary_count: u64,
1459 pub random_draw_count: u64,
1460}
1461
1462impl EvidenceCursor {
1463 fn from_evidence(evidence: &RuntimeEvidence) -> Result<Self, CanwuError> {
1464 let count = |len: usize, label: &str| {
1465 u64::try_from(len).map_err(|_| {
1466 CanwuError::new(
1467 ErrorCode::IdentifierExhausted,
1468 format!("{label} journal length exceeds the persistent cursor space"),
1469 )
1470 })
1471 };
1472 Ok(Self {
1473 event_count: evidence
1474 .archived
1475 .event_count
1476 .checked_add(count(evidence.events.len(), "event")?)
1477 .ok_or_else(|| invalid_snapshot_error("event journal cursor is exhausted"))?,
1478 command_count: evidence
1479 .archived
1480 .command_count
1481 .checked_add(count(evidence.commands.len(), "command")?)
1482 .ok_or_else(|| invalid_snapshot_error("command journal cursor is exhausted"))?,
1483 command_attempt_count: evidence
1484 .archived
1485 .command_attempt_count
1486 .checked_add(count(evidence.command_attempts.len(), "command-attempt")?)
1487 .ok_or_else(|| {
1488 invalid_snapshot_error("command-attempt journal cursor is exhausted")
1489 })?,
1490 ingress_count: evidence
1491 .archived
1492 .ingress_count
1493 .checked_add(count(evidence.ingress.len(), "ingress")?)
1494 .ok_or_else(|| invalid_snapshot_error("ingress journal cursor is exhausted"))?,
1495 boundary_count: evidence
1496 .archived
1497 .boundary_count
1498 .checked_add(count(evidence.boundaries.len(), "boundary")?)
1499 .ok_or_else(|| invalid_snapshot_error("boundary journal cursor is exhausted"))?,
1500 random_draw_count: evidence
1501 .archived
1502 .random_draw_count
1503 .checked_add(count(evidence.random_draws.len(), "random-draw")?)
1504 .ok_or_else(|| invalid_snapshot_error("random-draw journal cursor is exhausted"))?,
1505 })
1506 }
1507
1508 pub(super) fn checked_advance(
1509 self,
1510 segment: &EvidenceJournalSegment,
1511 ) -> Result<Self, CanwuError> {
1512 let advance = |value: u64, len: usize, label: &str| {
1513 value
1514 .checked_add(u64::try_from(len).map_err(|_| {
1515 invalid_snapshot_error(format!(
1516 "{label} journal segment exceeds the persistent cursor space"
1517 ))
1518 })?)
1519 .ok_or_else(|| {
1520 invalid_snapshot_error(format!(
1521 "{label} journal cursor exceeds the persistent cursor space"
1522 ))
1523 })
1524 };
1525 Ok(Self {
1526 event_count: advance(self.event_count, segment.events.len(), "event")?,
1527 command_count: advance(self.command_count, segment.commands.len(), "command")?,
1528 command_attempt_count: advance(
1529 self.command_attempt_count,
1530 segment.command_attempts.len(),
1531 "command-attempt",
1532 )?,
1533 ingress_count: advance(self.ingress_count, segment.ingress.len(), "ingress")?,
1534 boundary_count: advance(self.boundary_count, segment.boundaries.len(), "boundary")?,
1535 random_draw_count: advance(
1536 self.random_draw_count,
1537 segment.random_draws.len(),
1538 "random-draw",
1539 )?,
1540 })
1541 }
1542}
1543
1544#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1545pub struct SimulationCheckpoint {
1551 pub format_version: u32,
1552 pub journal_end: EvidenceCursor,
1553 pub state: SimulationSnapshot,
1554 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1555 pub archived_segment_headers: Vec<ArchivedSegmentHeader>,
1556 #[serde(default, skip_serializing_if = "Option::is_none")]
1557 pub archived_segment_manifest_root: Option<String>,
1558 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1559 pub archived_evidence_receipts: Vec<ArchivedEvidenceReceipt>,
1560 #[serde(default, skip_serializing_if = "Option::is_none")]
1561 pub archived_receipt_root: Option<String>,
1562 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1563 pub evidence_dependencies: Vec<EvidenceDependency>,
1564 #[serde(default, skip_serializing_if = "Option::is_none")]
1565 pub evidence_dependency_root: Option<String>,
1566 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1567 pub keyed_draw_reservations: Vec<KeyedDrawReservation>,
1568 #[serde(default, skip_serializing_if = "Option::is_none")]
1569 pub keyed_reservation_root: Option<String>,
1570}
1571
1572#[derive(Serialize)]
1573struct CompactCheckpointHashMaterial<'a> {
1574 state_checkpoint_hash: &'a str,
1575 journal_end: EvidenceCursor,
1576 archived_segment_manifest_root: Option<&'a str>,
1577 archived_receipt_root: Option<&'a str>,
1578 evidence_dependency_root: Option<&'a str>,
1579 keyed_reservation_root: Option<&'a str>,
1580}
1581
1582fn validate_compact_continuation(checkpoint: &SimulationCheckpoint) -> Result<(), CanwuError> {
1583 if checkpoint
1584 .archived_segment_headers
1585 .windows(2)
1586 .any(|headers| headers[0].end != headers[1].start)
1587 || checkpoint
1588 .archived_segment_headers
1589 .iter()
1590 .map(|header| &header.segment_id)
1591 .collect::<BTreeSet<_>>()
1592 .len()
1593 != checkpoint.archived_segment_headers.len()
1594 {
1595 return Err(invalid_snapshot_error(
1596 "compact archived-segment manifest is duplicated or noncontiguous",
1597 ));
1598 }
1599 if checkpoint
1600 .archived_evidence_receipts
1601 .windows(2)
1602 .any(|receipts| receipts[0].evidence >= receipts[1].evidence)
1603 {
1604 return Err(invalid_snapshot_error(
1605 "compact archived receipts must be sorted by unique evidence reference",
1606 ));
1607 }
1608 if checkpoint
1609 .evidence_dependencies
1610 .windows(2)
1611 .any(|dependencies| dependencies[0].reference >= dependencies[1].reference)
1612 {
1613 return Err(invalid_snapshot_error(
1614 "compact evidence dependencies must be sorted by unique evidence reference",
1615 ));
1616 }
1617 if checkpoint
1618 .keyed_draw_reservations
1619 .windows(2)
1620 .any(|reservations| {
1621 (&reservations[0].stream, &reservations[0].address)
1622 >= (&reservations[1].stream, &reservations[1].address)
1623 })
1624 {
1625 return Err(invalid_snapshot_error(
1626 "compact keyed reservations must be sorted by unique operation address",
1627 ));
1628 }
1629 validate_skipped_commitment_root(
1630 checkpoint.archived_segment_manifest_root.as_deref(),
1631 ARCHIVED_SEGMENT_MANIFEST_DOMAIN,
1632 &checkpoint.archived_segment_headers,
1633 "archived-segment manifest root",
1634 )?;
1635 validate_skipped_commitment_root(
1636 checkpoint.archived_receipt_root.as_deref(),
1637 ARCHIVED_RECEIPT_DOMAIN,
1638 &checkpoint.archived_evidence_receipts,
1639 "archived-receipt root",
1640 )?;
1641 validate_skipped_commitment_root(
1642 checkpoint.evidence_dependency_root.as_deref(),
1643 EVIDENCE_DEPENDENCY_DOMAIN,
1644 &checkpoint.evidence_dependencies,
1645 "evidence-dependency root",
1646 )?;
1647 validate_skipped_commitment_root(
1648 checkpoint.keyed_reservation_root.as_deref(),
1649 KEYED_RESERVATION_DOMAIN,
1650 &checkpoint.keyed_draw_reservations,
1651 "keyed-reservation root",
1652 )
1653}
1654
1655impl SimulationCheckpoint {
1656 pub fn reachable_archive_segment_ids(
1661 retained_checkpoints: &[Self],
1662 ) -> Result<BTreeSet<String>, CanwuError> {
1663 let mut reachable = BTreeSet::new();
1664 for checkpoint in retained_checkpoints {
1665 validate_compact_continuation(checkpoint)?;
1666 reachable.extend(
1667 checkpoint
1668 .archived_segment_headers
1669 .iter()
1670 .map(|header| header.segment_id.clone()),
1671 );
1672 }
1673 Ok(reachable)
1674 }
1675
1676 pub fn orphaned_archive_segment_ids(
1680 retained_checkpoints: &[Self],
1681 stored_segment_ids: &[String],
1682 ) -> Result<Vec<String>, CanwuError> {
1683 let reachable = Self::reachable_archive_segment_ids(retained_checkpoints)?;
1684 Ok(stored_segment_ids
1685 .iter()
1686 .filter(|segment_id| !reachable.contains(segment_id.as_str()))
1687 .cloned()
1688 .collect::<BTreeSet<_>>()
1689 .into_iter()
1690 .collect())
1691 }
1692}
1693
1694fn compact_checkpoint_hash(checkpoint: &SimulationCheckpoint) -> Result<String, CanwuError> {
1695 validate_compact_continuation(checkpoint)?;
1696 super::canonical_hash(
1697 "canwu.compact-checkpoint.v1",
1698 &CompactCheckpointHashMaterial {
1699 state_checkpoint_hash: &checkpoint.state.checkpoint_hash,
1700 journal_end: checkpoint.journal_end,
1701 archived_segment_manifest_root: checkpoint.archived_segment_manifest_root.as_deref(),
1702 archived_receipt_root: checkpoint.archived_receipt_root.as_deref(),
1703 evidence_dependency_root: checkpoint.evidence_dependency_root.as_deref(),
1704 keyed_reservation_root: checkpoint.keyed_reservation_root.as_deref(),
1705 },
1706 )
1707}
1708
1709#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1710pub struct EvidenceJournalSegment {
1712 pub format_version: u32,
1713 pub start: EvidenceCursor,
1714 pub end: EvidenceCursor,
1715 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1716 pub events: Vec<SimEvent>,
1717 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1718 pub commands: Vec<CommandRecord>,
1719 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1720 pub command_attempts: Vec<CommandAttemptRecord>,
1721 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1722 pub ingress: Vec<IngressRecord>,
1723 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1724 pub boundaries: Vec<BoundaryRecord>,
1725 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1726 pub random_draws: Vec<RandomDrawRecord>,
1727 #[serde(default, skip_serializing_if = "Option::is_none")]
1728 pub archive: Option<EvidenceArchiveIndex>,
1729}
1730
1731#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1732pub struct CheckpointJournal {
1734 pub checkpoint: SimulationCheckpoint,
1735 pub segments: Vec<EvidenceJournalSegment>,
1736}
1737
1738pub struct CompactedSimulation {
1746 simulation: Simulation,
1747 committed_seal_tokens: BTreeSet<String>,
1748}
1749
1750impl CompactedSimulation {
1751 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
1753 self.simulation.evidence_cursor()
1754 }
1755
1756 pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
1758 self.simulation.checkpoint()
1759 }
1760
1761 pub fn journal_segment_since(
1764 &self,
1765 start: EvidenceCursor,
1766 ) -> Result<EvidenceJournalSegment, CanwuError> {
1767 self.simulation.journal_segment_since(start)
1768 }
1769
1770 pub fn prepare_paged_checkpoint(
1773 &self,
1774 source: Option<&PagedSimulationCheckpoint>,
1775 provider: &dyn StatePageProvider,
1776 ) -> Result<PreparedPagedSimulationCheckpoint, CanwuError> {
1777 self.simulation.prepare_paged_checkpoint(source, provider)
1778 }
1779
1780 #[must_use]
1783 pub fn boundaries(&self) -> &[BoundaryRecord] {
1784 self.simulation.boundaries()
1785 }
1786
1787 pub fn outbox_entries(&self) -> Result<Vec<OutboxEntry>, CanwuError> {
1791 self.simulation.outbox_entries()
1792 }
1793
1794 pub fn outbox_entries_for_segment(
1799 &self,
1800 segment: &EvidenceJournalSegment,
1801 ) -> Result<Vec<OutboxEntry>, CanwuError> {
1802 Simulation::outbox_entries_for_boundaries(
1803 &self.simulation.state.metadata.run_manifest_hash,
1804 &segment.boundaries,
1805 )
1806 }
1807
1808 #[must_use]
1810 pub fn archived_evidence_receipt(
1811 &self,
1812 reference: &EvidenceRef,
1813 ) -> Option<&ArchivedEvidenceReceipt> {
1814 self.simulation
1815 .state
1816 .evidence
1817 .archived_evidence_receipts
1818 .get(reference)
1819 }
1820
1821 pub fn load_archived_evidence_segment(
1824 &self,
1825 reference: &EvidenceRef,
1826 provider: &dyn ArchiveProvider,
1827 ) -> Result<EvidenceJournalSegment, CanwuError> {
1828 let receipt = self.archived_evidence_receipt(reference).ok_or_else(|| {
1829 CanwuError::new(
1830 ErrorCode::EvidenceUnavailable,
1831 "no committed archive receipt exists for the requested evidence",
1832 )
1833 })?;
1834 load_verified_archived_evidence_segment(receipt, provider)
1835 }
1836
1837 pub fn prepare_decision_archive(
1841 &self,
1842 keys: &[DecisionHistoryKey],
1843 ) -> Result<PreparedDecisionArchive, CanwuError> {
1844 self.simulation
1845 .state
1846 .current
1847 .decisions
1848 .prepare_decision_archive(keys)
1849 .map_err(super::decision::decision_error)
1850 }
1851
1852 pub fn commit_decision_archive(
1857 &mut self,
1858 prepared: &PreparedDecisionArchive,
1859 provider: &dyn DecisionArchiveProvider,
1860 ) -> Result<IngressReceipt, CanwuError> {
1861 let verified = self
1862 .simulation
1863 .state
1864 .current
1865 .decisions
1866 .verify_decision_archive(prepared, provider)
1867 .map_err(super::decision::decision_error)?;
1868 let at = self.simulation.time();
1869 self.simulation
1870 .enqueue_decision_archive_commit(at, 0, verified)
1871 }
1872
1873 pub fn seal_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
1880 if self
1881 .simulation
1882 .evidence_dependencies()?
1883 .iter()
1884 .any(|dependency| dependency.requirement == EvidenceRequirement::PayloadRequired)
1885 {
1886 return Err(CanwuError::new(
1887 ErrorCode::ArchiveNotReady,
1888 "payload-required continuations must use prepare/store/commit sealing",
1889 ));
1890 }
1891 let before = self.simulation.fork();
1892 match self.simulation.seal_retained_evidence() {
1893 Ok(segment) => Ok(segment),
1894 Err(error) => {
1895 self.simulation = before;
1896 Err(error)
1897 }
1898 }
1899 }
1900
1901 pub fn prepare_evidence_seal(&self) -> Result<Option<PreparedEvidenceSeal>, CanwuError> {
1907 let source_state_hash = self.simulation.authoritative_state_hash()?;
1908 let source_checkpoint_hash = compact_checkpoint_hash(&self.simulation.checkpoint()?)?;
1909 let source_end = self.simulation.evidence_cursor()?;
1910 let mut candidate = self.simulation.fork();
1911 let Some(segment) = candidate.seal_retained_evidence()? else {
1912 return Ok(None);
1913 };
1914 let segment_id = segment
1915 .archive
1916 .as_ref()
1917 .ok_or_else(|| archive_error("prepared segment has no archive index"))?
1918 .header
1919 .segment_id
1920 .clone();
1921 let target_checkpoint_hash = compact_checkpoint_hash(&candidate.checkpoint()?)?;
1922 let token_hash = super::canonical_hash(
1923 "canwu.evidence.seal-token.v1",
1924 &(
1925 &source_state_hash,
1926 &source_checkpoint_hash,
1927 source_end,
1928 &segment_id,
1929 &target_checkpoint_hash,
1930 ),
1931 )?;
1932 Ok(Some(PreparedEvidenceSeal {
1933 token: EvidenceSealToken {
1934 source_state_hash,
1935 source_checkpoint_hash,
1936 source_end,
1937 segment_id,
1938 target_checkpoint_hash,
1939 token_hash,
1940 },
1941 segment,
1942 }))
1943 }
1944
1945 pub fn commit_evidence_seal(
1948 &mut self,
1949 token: &EvidenceSealToken,
1950 provider: &dyn ArchiveProvider,
1951 ) -> Result<(), CanwuError> {
1952 if self.committed_seal_tokens.contains(&token.token_hash) {
1953 return Ok(());
1954 }
1955 let expected_token_hash = super::canonical_hash(
1956 "canwu.evidence.seal-token.v1",
1957 &(
1958 &token.source_state_hash,
1959 &token.source_checkpoint_hash,
1960 token.source_end,
1961 &token.segment_id,
1962 &token.target_checkpoint_hash,
1963 ),
1964 )?;
1965 if expected_token_hash != token.token_hash {
1966 return Err(archive_error("evidence seal token hash is invalid"));
1967 }
1968 if self.simulation.authoritative_state_hash()? != token.source_state_hash
1969 || compact_checkpoint_hash(&self.simulation.checkpoint()?)?
1970 != token.source_checkpoint_hash
1971 || self.simulation.evidence_cursor()? != token.source_end
1972 {
1973 return Err(CanwuError::new(
1974 ErrorCode::StaleSealToken,
1975 "evidence seal token no longer names the live source cut",
1976 ));
1977 }
1978 let stored = provider
1979 .load_evidence_segment(&token.segment_id)?
1980 .ok_or_else(|| {
1981 CanwuError::new(
1982 ErrorCode::ArchiveNotReady,
1983 "prepared evidence segment is not available from the archive provider",
1984 )
1985 })?;
1986 let archive = stored
1987 .archive
1988 .as_ref()
1989 .ok_or_else(|| archive_error("stored evidence segment has no archive index"))?;
1990 if archive.header.segment_id != token.segment_id {
1991 return Err(archive_error(
1992 "stored evidence segment ID does not match the seal token",
1993 ));
1994 }
1995 verify_archived_segment(&stored)?;
1996
1997 let before = self.simulation.fork();
1998 let result = (|| {
1999 let actual = self.simulation.seal_retained_evidence()?.ok_or_else(|| {
2000 CanwuError::new(
2001 ErrorCode::StaleSealToken,
2002 "the prepared evidence tail is no longer retained",
2003 )
2004 })?;
2005 if actual != stored {
2006 return Err(archive_error(
2007 "stored evidence segment differs from the prepared live cut",
2008 ));
2009 }
2010 self.simulation
2011 .validate_payload_required_archive(provider)?;
2012 if compact_checkpoint_hash(&self.simulation.checkpoint()?)?
2013 != token.target_checkpoint_hash
2014 {
2015 return Err(archive_error(
2016 "committed checkpoint hash differs from the seal token",
2017 ));
2018 }
2019 Ok(())
2020 })();
2021 if let Err(error) = result {
2022 self.simulation = before;
2023 return Err(error);
2024 }
2025 self.committed_seal_tokens.insert(token.token_hash.clone());
2026 Ok(())
2027 }
2028
2029 #[must_use]
2030 pub const fn time(&self) -> SimTime {
2031 self.simulation.time()
2032 }
2033
2034 #[must_use]
2035 pub const fn revision(&self) -> u64 {
2036 self.simulation.revision()
2037 }
2038
2039 #[must_use]
2040 pub fn checkpoint_hash(&self) -> &str {
2041 self.simulation.checkpoint_hash()
2042 }
2043
2044 #[must_use]
2045 pub fn boundary_head_hash(&self) -> Option<&str> {
2046 self.simulation.boundary_head_hash()
2047 }
2048
2049 pub fn entities(&self) -> impl Iterator<Item = &super::EntityRef> {
2050 self.simulation.entities()
2051 }
2052
2053 #[must_use]
2054 pub fn entity_exists(&self, entity: &super::EntityRef) -> bool {
2055 self.simulation.entity_exists(entity)
2056 }
2057
2058 #[must_use]
2059 pub fn world(&self) -> super::WorldSnapshot {
2060 self.simulation.world()
2061 }
2062
2063 #[must_use]
2064 pub fn knowledge(&self) -> &KnowledgeSnapshot {
2065 self.simulation.knowledge()
2066 }
2067
2068 #[must_use]
2069 pub fn domain_record(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
2070 self.simulation.domain_record(reference)
2071 }
2072
2073 #[must_use]
2074 pub fn decision_ticket(&self, id: super::DecisionTicketId) -> Option<&super::DecisionTicket> {
2075 self.simulation.decision_ticket(id)
2076 }
2077
2078 #[must_use]
2079 pub fn decision_controller(&self, id: &str) -> Option<&super::DecisionControllerBinding> {
2080 self.simulation.decision_controller(id)
2081 }
2082
2083 #[must_use]
2084 pub fn decision_trace(&self, id: super::DecisionTraceId) -> Option<&super::DecisionTrace> {
2085 self.simulation.decision_trace(id)
2086 }
2087
2088 #[must_use]
2089 pub fn decision_attempt(
2090 &self,
2091 id: super::DecisionRequestId,
2092 ) -> Option<&super::DecisionAttemptRecord> {
2093 self.simulation.decision_attempt(id)
2094 }
2095
2096 #[must_use]
2097 pub fn decision_hot_state(&self) -> super::DecisionHotState {
2098 self.simulation.decision_hot_state()
2099 }
2100
2101 #[must_use]
2102 pub fn decision_history_location(
2103 &self,
2104 key: &super::DecisionHistoryKey,
2105 ) -> super::DecisionHistoryLocation {
2106 self.simulation.decision_history_location(key)
2107 }
2108
2109 pub fn decision_history_location_with_provider(
2110 &self,
2111 key: &super::DecisionHistoryKey,
2112 provider: &dyn super::DecisionArchiveProvider,
2113 ) -> Result<super::DecisionHistoryLocation, CanwuError> {
2114 self.simulation
2115 .decision_history_location_with_provider(key, provider)
2116 }
2117
2118 #[must_use]
2119 pub fn typed_domain_record<T: DomainRecordType>(
2120 &self,
2121 reference: &TypedDomainRecordRef<T>,
2122 ) -> Option<&DomainRecord> {
2123 self.simulation.typed_domain_record(reference)
2124 }
2125
2126 pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
2127 self.simulation.submit(envelope)
2128 }
2129
2130 pub fn process_command(
2131 &mut self,
2132 request: CommandRequest,
2133 ) -> Result<CommandOutcome, CanwuError> {
2134 self.simulation.process_command(request)
2135 }
2136
2137 pub fn enqueue_command(
2138 &mut self,
2139 due_at: SimTime,
2140 priority: i32,
2141 request: CommandRequest,
2142 ) -> Result<IngressReceipt, CanwuError> {
2143 self.simulation.enqueue_command(due_at, priority, request)
2144 }
2145
2146 pub fn enqueue_plugin_ingress(
2147 &mut self,
2148 request: PluginIngressRequest,
2149 ) -> Result<IngressReceipt, CanwuError> {
2150 self.simulation.enqueue_plugin_ingress(request)
2151 }
2152
2153 pub fn enqueue_permitted_plugin_ingress(
2154 &mut self,
2155 request: PluginIngressRequest,
2156 permit: &super::PluginIngressPermit,
2157 ) -> Result<IngressReceipt, CanwuError> {
2158 self.simulation
2159 .enqueue_permitted_plugin_ingress(request, permit)
2160 }
2161
2162 pub fn prepare_decision(
2163 &self,
2164 decision_request_id: super::DecisionRequestId,
2165 command_request_id: Option<super::CommandRequestId>,
2166 ticket_id: super::DecisionTicketId,
2167 policy: &dyn super::DecisionPolicy,
2168 ) -> Result<super::DecisionEvaluation, CanwuError> {
2169 self.simulation
2170 .prepare_decision(decision_request_id, command_request_id, ticket_id, policy)
2171 }
2172
2173 pub fn prepare_decision_at(
2174 &self,
2175 due_at: super::SimTime,
2176 decision_request_id: super::DecisionRequestId,
2177 command_request_id: Option<super::CommandRequestId>,
2178 ticket_id: super::DecisionTicketId,
2179 policy: &dyn super::DecisionPolicy,
2180 ) -> Result<super::DecisionEvaluation, CanwuError> {
2181 self.simulation.prepare_decision_at(
2182 due_at,
2183 decision_request_id,
2184 command_request_id,
2185 ticket_id,
2186 policy,
2187 )
2188 }
2189
2190 pub fn enqueue_decision(
2191 &mut self,
2192 due_at: super::SimTime,
2193 priority: i32,
2194 request: super::DecisionIngressRequest,
2195 ) -> Result<IngressReceipt, CanwuError> {
2196 self.simulation.enqueue_decision(due_at, priority, request)
2197 }
2198
2199 pub fn drive_decision(
2200 &mut self,
2201 due_at: super::SimTime,
2202 priority: i32,
2203 decision_request_id: super::DecisionRequestId,
2204 command_request_id: Option<super::CommandRequestId>,
2205 ticket_id: super::DecisionTicketId,
2206 policy: &dyn super::DecisionPolicy,
2207 ) -> Result<super::DecisionEvaluation, CanwuError> {
2208 self.simulation.drive_decision(
2209 due_at,
2210 priority,
2211 decision_request_id,
2212 command_request_id,
2213 ticket_id,
2214 policy,
2215 )
2216 }
2217
2218 pub fn schedule_calendar_boundary(
2219 &mut self,
2220 due_at: SimTime,
2221 cadences: Vec<SystemCadence>,
2222 ) -> Result<IngressReceipt, CanwuError> {
2223 self.simulation.schedule_calendar_boundary(due_at, cadences)
2224 }
2225
2226 pub fn advance(&mut self, duration: SimDuration) -> Result<Vec<SimEvent>, CanwuError> {
2227 self.simulation.advance(duration)
2228 }
2229
2230 pub fn advance_canonical(
2231 &mut self,
2232 duration: SimDuration,
2233 ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
2234 self.simulation.advance_canonical(duration)
2235 }
2236
2237 pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
2238 self.simulation.step_canonical()
2239 }
2240
2241 pub fn settle_boundary(
2242 &mut self,
2243 request: BoundaryRequest,
2244 ) -> Result<BoundaryReceipt, CanwuError> {
2245 self.simulation.settle_boundary(request)
2246 }
2247
2248 pub fn snapshot_with_segments(
2251 &self,
2252 mut segments: Vec<EvidenceJournalSegment>,
2253 ) -> Result<SimulationSnapshot, CanwuError> {
2254 let tail = self
2255 .simulation
2256 .journal_segment_since(self.simulation.state.evidence.archived)?;
2257 if tail.start != tail.end {
2258 segments.push(tail);
2259 }
2260 let snapshot =
2261 Simulation::snapshot_from_checkpoint_and_journal(self.checkpoint()?, segments)?;
2262 Simulation::from_snapshot(snapshot.clone())?;
2263 Ok(snapshot)
2264 }
2265
2266 pub fn replay_journal_with_segments(
2268 &self,
2269 segments: Vec<EvidenceJournalSegment>,
2270 ) -> Result<ReplayJournal, CanwuError> {
2271 let snapshot = self.snapshot_with_segments(segments)?;
2272 let simulation = Simulation::from_snapshot(snapshot)?;
2273 Ok(simulation.replay_journal())
2274 }
2275
2276 pub fn from_checkpoint_and_journal(
2279 checkpoint: SimulationCheckpoint,
2280 segments: Vec<EvidenceJournalSegment>,
2281 ) -> Result<Self, CanwuError> {
2282 Simulation::from_checkpoint_and_journal(checkpoint, segments)?.into_compacted()
2283 }
2284
2285 pub fn from_checkpoint_and_journal_with_plugins(
2287 checkpoint: SimulationCheckpoint,
2288 segments: Vec<EvidenceJournalSegment>,
2289 plugins: &[&dyn SimulationPlugin],
2290 ) -> Result<Self, CanwuError> {
2291 let mut simulation = Simulation::from_checkpoint_and_journal(checkpoint, segments)?;
2292 for plugin in plugins {
2293 simulation.register_plugin(*plugin)?;
2294 }
2295 simulation.ensure_runtime_ready()?;
2296 simulation.into_compacted()
2297 }
2298}
2299
2300impl Simulation {
2301 pub fn into_compacted(self) -> Result<CompactedSimulation, CanwuError> {
2307 self.ensure_runtime_ready()?;
2308 Ok(CompactedSimulation {
2309 simulation: self,
2310 committed_seal_tokens: BTreeSet::new(),
2311 })
2312 }
2313
2314 fn evidence_dependencies(&self) -> Result<Vec<EvidenceDependency>, CanwuError> {
2315 let mut dependencies = BTreeMap::new();
2316 let identity = EvidenceRequirement::IdentityOnly;
2317
2318 for records in self.state.current.knowledge.records.values() {
2319 for record in records.values() {
2320 for reference in &record.origin.evidence {
2321 promote_dependency(&mut dependencies, reference.clone(), identity);
2322 }
2323 }
2324 }
2325
2326 for record in self.state.current.domain_records.values() {
2327 let schema = self
2328 .plugins
2329 .record_schemas
2330 .get(&record.reference.kind)
2331 .map(|(_, schema)| schema)
2332 .ok_or_else(|| {
2333 CanwuError::new(
2334 ErrorCode::ArchiveNotReady,
2335 format!(
2336 "current domain record {} has no registered schema",
2337 record.reference
2338 ),
2339 )
2340 })?;
2341 add_identity_evidence_dependencies(&mut dependencies, record, schema)?;
2342 add_payload_required_continuation_dependencies(&mut dependencies, record, schema)?;
2343 let version = self
2344 .state
2345 .metadata
2346 .current_domain_record_versions
2347 .get(&record.reference)
2348 .filter(|version| version.version == record.version)
2349 .ok_or_else(|| {
2350 CanwuError::new(
2351 ErrorCode::ArchiveNotReady,
2352 format!(
2353 "current domain record {} has no exact version provenance",
2354 record.reference
2355 ),
2356 )
2357 })?;
2358 if !matches!(
2359 version.established_by,
2360 DomainRecordVersionSource::InitialScenario
2361 ) {
2362 promote_dependency(
2363 &mut dependencies,
2364 EvidenceRef::DomainRecordVersion(version.clone()),
2365 identity,
2366 );
2367 } else if !self.bound_initial_scenario().is_some_and(|scenario| {
2368 scenario.domain_records.iter().any(|initial| {
2369 initial.reference == record.reference && initial.version == record.version
2370 })
2371 }) {
2372 return Err(CanwuError::new(
2373 ErrorCode::ArchiveNotReady,
2374 format!(
2375 "current domain record {} has invalid initial-scenario provenance",
2376 record.reference
2377 ),
2378 ));
2379 }
2380 }
2381
2382 for reservation in &self.state.evidence.keyed_draw_reservations {
2383 promote_dependency(
2384 &mut dependencies,
2385 reservation.operation_evidence.clone(),
2386 identity,
2387 );
2388 promote_dependency(
2389 &mut dependencies,
2390 reservation.draw_receipt.evidence.clone(),
2391 identity,
2392 );
2393 }
2394 for draw in &self.state.evidence.random_draws {
2395 if matches!(draw.address, RandomDrawAddress::OperationV1(_)) {
2396 let reference = draw.operation_evidence.clone().ok_or_else(|| {
2397 CanwuError::new(
2398 ErrorCode::ArchiveNotReady,
2399 "operation-keyed retained draw is missing its operation evidence",
2400 )
2401 })?;
2402 promote_dependency(&mut dependencies, reference, identity);
2403 promote_dependency(
2404 &mut dependencies,
2405 EvidenceRef::RandomDraw(draw.id),
2406 identity,
2407 );
2408 }
2409 }
2410
2411 for archived in self.state.evidence.archived_command_requests.values() {
2412 add_command_outcome_dependencies(&mut dependencies, &archived.outcome);
2413 }
2414 for archived in self.state.evidence.archived_ingress_requests.values() {
2415 promote_dependency(
2416 &mut dependencies,
2417 EvidenceRef::Ingress(archived.receipt.ingress_id),
2418 identity,
2419 );
2420 }
2421 for archived in self.state.evidence.archived_decision_requests.values() {
2422 promote_dependency(
2423 &mut dependencies,
2424 EvidenceRef::Ingress(archived.receipt.ingress_id),
2425 identity,
2426 );
2427 }
2428 for attempt in &self.state.evidence.command_attempts {
2429 if attempt.request_id.is_none() {
2430 continue;
2431 }
2432 promote_dependency(
2433 &mut dependencies,
2434 EvidenceRef::CommandAttempt(attempt.id),
2435 identity,
2436 );
2437 if let CommandAttemptOutcome::Accepted { command_id } = attempt.outcome {
2438 promote_dependency(
2439 &mut dependencies,
2440 EvidenceRef::Command(command_id),
2441 identity,
2442 );
2443 if let Some(command) = self
2444 .state
2445 .evidence
2446 .commands
2447 .iter()
2448 .find(|command| command.id == command_id)
2449 {
2450 for event in &command.emitted_events {
2451 promote_dependency(&mut dependencies, EvidenceRef::Event(*event), identity);
2452 }
2453 }
2454 }
2455 }
2456 for ingress in &self.state.evidence.ingress {
2457 if matches!(
2458 ingress.payload,
2459 IngressPayload::Command { .. } | IngressPayload::Decision { .. }
2460 ) {
2461 promote_dependency(
2462 &mut dependencies,
2463 EvidenceRef::Ingress(ingress.id),
2464 identity,
2465 );
2466 }
2467 }
2468
2469 for action in self.state.scheduler.actions.values() {
2470 match action {
2471 ScheduledAction::ArmyArrival { order_event, .. }
2472 | ScheduledAction::PersonArrival { order_event, .. } => promote_dependency(
2473 &mut dependencies,
2474 EvidenceRef::Event(*order_event),
2475 identity,
2476 ),
2477 ScheduledAction::KnowledgeReport { dispatch_event, .. } => promote_dependency(
2478 &mut dependencies,
2479 EvidenceRef::Event(*dispatch_event),
2480 identity,
2481 ),
2482 ScheduledAction::PluginDirective { cause, .. } => {
2483 add_cause_dependency(&mut dependencies, cause);
2484 }
2485 }
2486 }
2487
2488 Ok(dependencies
2489 .into_iter()
2490 .map(|(reference, requirement)| EvidenceDependency {
2491 reference,
2492 requirement,
2493 })
2494 .collect())
2495 }
2496
2497 fn validate_payload_required_archive(
2498 &self,
2499 provider: &dyn ArchiveProvider,
2500 ) -> Result<(), CanwuError> {
2501 for dependency in self
2502 .evidence_dependencies()?
2503 .into_iter()
2504 .filter(|dependency| dependency.requirement == EvidenceRequirement::PayloadRequired)
2505 {
2506 let receipt = self
2507 .state
2508 .evidence
2509 .archived_evidence_receipts
2510 .get(&dependency.reference)
2511 .ok_or_else(|| {
2512 CanwuError::new(
2513 ErrorCode::ArchiveNotReady,
2514 "payload-required evidence has no committed archive receipt",
2515 )
2516 })?;
2517 load_verified_archived_evidence_segment(receipt, provider)?;
2518 }
2519 Ok(())
2520 }
2521
2522 fn ensure_retained_evidence_is_sealable(&self) -> Result<(), CanwuError> {
2523 if !self.state.scheduler.pending_ingress.is_empty() {
2524 return Err(CanwuError::new(
2525 ErrorCode::ArchiveNotReady,
2526 "live evidence can be sealed only when the canonical ingress queue is empty",
2527 ));
2528 }
2529 let admitted_attempts: std::collections::BTreeSet<_> = self
2530 .state
2531 .evidence
2532 .boundaries
2533 .iter()
2534 .flat_map(|record| record.admitted_attempts.iter().copied())
2535 .collect();
2536 if admitted_attempts.len() != self.state.evidence.command_attempts.len()
2537 || self
2538 .state
2539 .evidence
2540 .command_attempts
2541 .iter()
2542 .any(|attempt| !admitted_attempts.contains(&attempt.id))
2543 {
2544 return Err(CanwuError::new(
2545 ErrorCode::ArchiveNotReady,
2546 "live evidence sealing requires every retained command attempt to belong to a completed boundary",
2547 ));
2548 }
2549 let admitted_commands: std::collections::BTreeSet<_> = self
2550 .state
2551 .evidence
2552 .boundaries
2553 .iter()
2554 .flat_map(|record| record.admitted_commands.iter().copied())
2555 .collect();
2556 if admitted_commands.len() != self.state.evidence.commands.len()
2557 || self
2558 .state
2559 .evidence
2560 .commands
2561 .iter()
2562 .any(|command| !admitted_commands.contains(&command.id))
2563 {
2564 return Err(CanwuError::new(
2565 ErrorCode::ArchiveNotReady,
2566 "live evidence sealing requires every retained command to belong to a completed boundary",
2567 ));
2568 }
2569 let admitted_ingress: std::collections::BTreeSet<_> = self
2570 .state
2571 .evidence
2572 .boundaries
2573 .iter()
2574 .flat_map(|record| record.admitted_ingress.iter().copied())
2575 .collect();
2576 if admitted_ingress.len() != self.state.evidence.ingress.len()
2577 || self
2578 .state
2579 .evidence
2580 .ingress
2581 .iter()
2582 .any(|record| !admitted_ingress.contains(&record.id))
2583 {
2584 return Err(CanwuError::new(
2585 ErrorCode::ArchiveNotReady,
2586 "live evidence sealing requires every retained ingress record to belong to a completed boundary",
2587 ));
2588 }
2589 let admitted_events: std::collections::BTreeSet<_> = self
2590 .state
2591 .evidence
2592 .boundaries
2593 .iter()
2594 .flat_map(|record| record.admitted_events.iter().copied())
2595 .collect();
2596 if self.state.counters.admitted_event_count
2597 != self
2598 .state
2599 .evidence
2600 .archived
2601 .event_count
2602 .checked_add(
2603 u64::try_from(self.state.evidence.events.len()).map_err(|_| {
2604 CanwuError::new(
2605 ErrorCode::ArchiveNotReady,
2606 "retained event count exceeds the live archive cursor range",
2607 )
2608 })?,
2609 )
2610 .ok_or_else(|| {
2611 CanwuError::new(
2612 ErrorCode::ArchiveNotReady,
2613 "retained event cursor is exhausted",
2614 )
2615 })?
2616 || admitted_events.len() != self.state.evidence.events.len()
2617 || self
2618 .state
2619 .evidence
2620 .events
2621 .iter()
2622 .any(|event| !admitted_events.contains(&event.id))
2623 {
2624 return Err(CanwuError::new(
2625 ErrorCode::ArchiveNotReady,
2626 "live evidence sealing requires every retained event to be admitted by a later completed boundary",
2627 ));
2628 }
2629 Ok(())
2630 }
2631
2632 fn seal_retained_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
2633 let start = self.state.evidence.archived;
2634 let end = self.evidence_cursor()?;
2635 if start == end {
2636 return Ok(None);
2637 }
2638 self.ensure_retained_evidence_is_sealable()?;
2639 let dependencies = self.evidence_dependencies()?;
2640 let checkpoint_hash = self.state.metadata.checkpoint_hash.clone();
2641 let commitment_roots = self.state.metadata.commitment_roots.clone();
2642 let commitment_cache = self.state.metadata.commitment_cache.clone();
2643 let prepared = (|| {
2644 self.refresh_checkpoint_hash()?;
2645
2646 let mut archived_command_requests = Vec::new();
2647 for attempt in &self.state.evidence.command_attempts {
2648 let Some(request_id) = attempt.request_id else {
2649 continue;
2650 };
2651 let outcome = self.command_outcome_from_attempt(attempt)?;
2652 archived_command_requests.push((
2653 request_id,
2654 ArchivedCommandRequestOutcome {
2655 input_hash: super::canonical_hash(
2656 "canwu.archive.command.request.v1",
2657 &(attempt.expected_revision, &attempt.envelope),
2658 )?,
2659 outcome,
2660 },
2661 ));
2662 }
2663
2664 let mut archived_ingress_requests = Vec::new();
2665 let mut archived_decision_requests = Vec::new();
2666 let mut archived_decision_command_requests = Vec::new();
2667 for record in &self.state.evidence.ingress {
2668 let receipt = IngressReceipt {
2669 ingress_id: record.id,
2670 issued_at: record.issued_at,
2671 due_at: record.due_at,
2672 };
2673 match &record.payload {
2674 IngressPayload::Command { request } => {
2675 archived_ingress_requests.push((
2676 request.request_id,
2677 ArchivedIngressRequest {
2678 input_hash: super::canonical_hash(
2679 "canwu.archive.ingress.command.v1",
2680 &(record.due_at, record.priority, request.as_ref()),
2681 )?,
2682 receipt,
2683 },
2684 ));
2685 }
2686 IngressPayload::Decision { request } => {
2687 if let Some(command) = &request.command {
2688 archived_decision_command_requests.push(command.request_id);
2689 }
2690 archived_decision_requests.push((
2691 request.request_id,
2692 ArchivedIngressRequest {
2693 input_hash: super::canonical_hash(
2694 "canwu.ingress.decision-request.v1",
2695 &(record.due_at, record.priority, request.as_ref()),
2696 )?,
2697 receipt,
2698 },
2699 ));
2700 }
2701 IngressPayload::Plugin { .. }
2702 | IngressPayload::Calendar { .. }
2703 | IngressPayload::Maintenance { .. } => {}
2704 }
2705 }
2706 Ok::<_, CanwuError>((
2707 archived_command_requests,
2708 archived_ingress_requests,
2709 archived_decision_requests,
2710 archived_decision_command_requests,
2711 ))
2712 })();
2713 let (
2714 archived_command_requests,
2715 archived_ingress_requests,
2716 archived_decision_requests,
2717 archived_decision_command_requests,
2718 ) = match prepared {
2719 Ok(prepared) => prepared,
2720 Err(error) => {
2721 self.state.metadata.checkpoint_hash = checkpoint_hash;
2722 self.state.metadata.commitment_roots = commitment_roots;
2723 self.state.metadata.commitment_cache = commitment_cache;
2724 return Err(error);
2725 }
2726 };
2727
2728 let mut segment = EvidenceJournalSegment {
2729 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
2730 start,
2731 end,
2732 events: self.state.evidence.events.clone(),
2733 commands: self.state.evidence.commands.clone(),
2734 command_attempts: self.state.evidence.command_attempts.clone(),
2735 ingress: self.state.evidence.ingress.clone(),
2736 boundaries: self.state.evidence.boundaries.clone(),
2737 random_draws: self.state.evidence.random_draws.clone(),
2738 archive: None,
2739 };
2740 let (archive, receipts) = evidence_archive_index(&segment)?;
2741 segment.archive = Some(archive.clone());
2742 verify_archived_segment(&segment)?;
2743
2744 let receipt_map: BTreeMap<_, _> = receipts
2745 .iter()
2746 .cloned()
2747 .map(|receipt| (receipt.evidence.clone(), receipt))
2748 .collect();
2749 let mut reservations = Vec::new();
2750 for draw in &segment.random_draws {
2751 let RandomDrawAddress::OperationV1(address) = &draw.address else {
2752 continue;
2753 };
2754 let operation_evidence = draw.operation_evidence.clone().ok_or_else(|| {
2755 archive_error("operation-keyed archived draw is missing operation evidence")
2756 })?;
2757 let draw_receipt = receipt_map
2758 .get(&EvidenceRef::RandomDraw(draw.id))
2759 .cloned()
2760 .ok_or_else(|| {
2761 archive_error("operation-keyed archived draw is missing its receipt")
2762 })?;
2763 reservations.push(KeyedDrawReservation {
2764 stream: draw.stream.clone(),
2765 address: address.clone(),
2766 upper_exclusive: draw.upper_exclusive,
2767 purpose_hash: super::random::purpose_hash_hex_v1(&draw.purpose)?,
2768 result: draw.value,
2769 draw_id: draw.id,
2770 operation_evidence,
2771 draw_receipt,
2772 });
2773 }
2774 reservations.sort_by(|left, right| {
2775 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2776 });
2777
2778 self.state.evidence.archived_boundary_head = self
2779 .state
2780 .evidence
2781 .boundaries
2782 .last()
2783 .map(|record| record.hash.clone())
2784 .or_else(|| self.state.evidence.archived_boundary_head.clone());
2785 self.state.evidence.archived_legacy_commands |= self
2786 .state
2787 .evidence
2788 .commands
2789 .iter()
2790 .any(|record| record.attempt_id.is_none());
2791 self.state.evidence.archived_tracked_attempts |=
2792 !self.state.evidence.command_attempts.is_empty()
2793 || !self.state.evidence.ingress.is_empty();
2794 self.state.evidence.archived_unqueued_command_history |= has_unqueued_command_history(
2795 &self.state.evidence.commands,
2796 &self.state.evidence.command_attempts,
2797 &self.state.evidence.ingress,
2798 );
2799 self.state
2800 .evidence
2801 .archived_command_requests
2802 .extend(archived_command_requests);
2803 self.state
2804 .evidence
2805 .archived_ingress_requests
2806 .extend(archived_ingress_requests);
2807 self.state
2808 .evidence
2809 .archived_decision_requests
2810 .extend(archived_decision_requests);
2811 self.state
2812 .evidence
2813 .archived_decision_command_requests
2814 .extend(archived_decision_command_requests);
2815 self.state.evidence.archived = end;
2816 self.state.evidence.events.clear();
2817 self.state.evidence.commands.clear();
2818 self.state.evidence.command_attempts.clear();
2819 self.state.evidence.ingress.clear();
2820 self.state.evidence.boundaries.clear();
2821 self.state.evidence.random_draws.clear();
2822 self.state
2823 .evidence
2824 .archived_segment_headers
2825 .push(archive.header);
2826 for receipt in receipts {
2827 if self
2828 .state
2829 .evidence
2830 .archived_evidence_receipts
2831 .insert(receipt.evidence.clone(), receipt)
2832 .is_some()
2833 {
2834 return Err(archive_error("archived evidence receipt was duplicated"));
2835 }
2836 }
2837 self.state
2838 .evidence
2839 .keyed_draw_reservations
2840 .extend(reservations);
2841 self.state
2842 .evidence
2843 .keyed_draw_reservations
2844 .sort_by(|left, right| {
2845 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2846 });
2847 if self
2848 .state
2849 .evidence
2850 .keyed_draw_reservations
2851 .windows(2)
2852 .any(|reservations| {
2853 (&reservations[0].stream, &reservations[0].address)
2854 == (&reservations[1].stream, &reservations[1].address)
2855 })
2856 {
2857 return Err(archive_error("archived keyed reservation was duplicated"));
2858 }
2859 retain_reachable_archived_evidence_receipts(
2860 &mut self.state.evidence.archived_evidence_receipts,
2861 &dependencies,
2862 &self.state.evidence.keyed_draw_reservations,
2863 );
2864 for dependency in &dependencies {
2865 if matches!(
2866 &dependency.reference,
2867 EvidenceRef::DomainRecordVersion(version)
2868 if matches!(
2869 version.established_by,
2870 DomainRecordVersionSource::InitialScenario
2871 )
2872 ) {
2873 continue;
2874 }
2875 if !self
2876 .state
2877 .evidence
2878 .archived_evidence_receipts
2879 .contains_key(&dependency.reference)
2880 {
2881 return Err(CanwuError::new(
2882 ErrorCode::ArchiveNotReady,
2883 "declared evidence was not present in the sealed archive prefix",
2884 ));
2885 }
2886 }
2887 Ok(Some(segment))
2888 }
2889
2890 pub(super) fn checkpoint_state(&self) -> SimulationSnapshot {
2891 self.checkpoint_state_with_paged_payloads(true)
2892 }
2893
2894 fn checkpoint_state_with_paged_payloads(
2895 &self,
2896 include_paged_payloads: bool,
2897 ) -> SimulationSnapshot {
2898 SimulationSnapshot {
2899 engine_version: ENGINE_VERSION.to_owned(),
2900 snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
2901 run_manifest: Some(self.state.metadata.run_manifest.clone()),
2902 run_manifest_hash: self.state.metadata.run_manifest_hash.clone(),
2903 run_configuration: Some(self.state.metadata.run_configuration.clone()),
2904 checkpoint_hash: self.state.metadata.checkpoint_hash.clone(),
2905 commitment_format_version: self.state.metadata.commitment_format_version,
2906 commitment_roots: self.state.metadata.commitment_roots.clone(),
2907 revision_format_version: STATE_REVISION_FORMAT_VERSION,
2908 state_revision: self.state.counters.state_revision,
2909 replay_revision_format_version: self.state.metadata.replay_revision_format_version,
2910 admission_cursor_format_version: ADMISSION_CURSOR_FORMAT_VERSION,
2911 admitted_attempt_count: self.state.counters.admitted_attempt_count,
2912 admitted_command_count: self.state.counters.admitted_command_count,
2913 admitted_event_count: self.state.counters.admitted_event_count,
2914 initial_time: self.state.scheduler.initial_time,
2915 initial_scenario: self.bound_initial_scenario().cloned(),
2916 now: self.state.scheduler.now,
2917 plugin_registration_closed: self.state.metadata.plugin_registration_closed,
2918 entities: self.state.current.entities.iter().cloned().collect(),
2919 world: self.world(),
2920 knowledge: self.state.current.knowledge.clone(),
2921 events: Vec::new(),
2922 commands: Vec::new(),
2923 command_attempts: Vec::new(),
2924 ingress: Vec::new(),
2925 boundaries: Vec::new(),
2926 plugin_components: self
2927 .state
2928 .current
2929 .plugin_components
2930 .values()
2931 .cloned()
2932 .collect(),
2933 domain_records: if include_paged_payloads {
2934 self.state
2935 .current
2936 .domain_records
2937 .values()
2938 .cloned()
2939 .collect()
2940 } else {
2941 Vec::new()
2942 },
2943 decisions: if include_paged_payloads {
2944 self.state.current.decisions.clone()
2945 } else {
2946 DecisionState::default()
2947 },
2948 plugin_descriptors: self.plugins.descriptors().cloned().collect(),
2949 schema: self.schema.clone(),
2950 root_seed: self.state.current.root_seed,
2951 authority_root_seed: self.state.current.authority_root_seed,
2952 random_streams: self
2953 .state
2954 .current
2955 .random_streams
2956 .values()
2957 .cloned()
2958 .collect(),
2959 random_draws: Vec::new(),
2960 scheduled: self
2961 .state
2962 .scheduler
2963 .actions
2964 .iter()
2965 .map(|(key, action)| ScheduledRecord {
2966 key: key.clone(),
2967 action: action.clone(),
2968 })
2969 .collect(),
2970 legacy_rng: None,
2971 next_event_id: self.state.counters.next_event_id,
2972 next_command_id: self.state.counters.next_command_id,
2973 next_command_attempt_id: self.state.counters.next_command_attempt_id,
2974 next_ingress_id: self.state.counters.next_ingress_id,
2975 next_boundary_id: self.state.counters.next_boundary_id,
2976 next_random_draw_id: self.state.counters.next_random_draw_id,
2977 next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
2978 next_schedule_sequence: self.state.counters.next_schedule_sequence,
2979 next_correlation_id: self.state.counters.next_correlation_id,
2980 next_decision_trace_id: self.state.counters.next_decision_trace_id,
2981 }
2982 }
2983
2984 fn build_paged_checkpoint_pages(
2985 &self,
2986 provider: Option<&dyn StatePageProvider>,
2987 ) -> Result<(PagedSimulationCheckpoint, Vec<StatePageBlob>), CanwuError> {
2988 let (domain_records, domain_pages) = match provider {
2989 Some(provider) => self
2990 .state
2991 .current
2992 .domain_records
2993 .missing_state_pages(provider)?,
2994 None => self.state.current.domain_records.state_pages()?,
2995 };
2996 let decisions = &self.state.current.decisions;
2997 let hot_decisions = decisions.paged_checkpoint_hot_state();
2998 let hot_decision_page =
2999 StatePageBlob::new(serde_json::to_vec(&hot_decisions).map_err(|error| {
3000 invalid_snapshot_error(format!("cannot encode paged hot decision state: {error}"))
3001 })?)?;
3002 let decision_bucket_pages = decisions
3003 .decision_archive_bucket_page_ids()
3004 .iter()
3005 .map(|(bucket, page_id)| (*bucket, page_id.clone()))
3006 .collect::<BTreeMap<_, _>>();
3007 let mut decision_pages = Vec::with_capacity(
3008 decision_bucket_pages
3009 .len()
3010 .saturating_add(
3011 decision_bucket_pages
3012 .len()
3013 .div_ceil(MAX_PAGED_DECISION_DIRECTORY_ENTRIES),
3014 )
3015 .saturating_add(2),
3016 );
3017 for (bucket, expected_page_id) in &decision_bucket_pages {
3018 let Some(bucket_page) = decisions
3019 .decision_archive_bucket_page(*bucket)
3020 .map_err(|error| invalid_snapshot_error(error.to_string()))?
3021 else {
3022 if provider.is_none() {
3023 return Err(invalid_snapshot_error(
3024 "portable decision checkpoint requires every archive bucket to be resident",
3025 ));
3026 }
3027 continue;
3032 };
3033 let page = StatePageBlob::new(serde_json::to_vec(&bucket_page).map_err(|error| {
3034 invalid_snapshot_error(format!(
3035 "cannot encode paged decision archive bucket: {error}"
3036 ))
3037 })?)?;
3038 if page.page_id != *expected_page_id {
3039 return Err(invalid_snapshot_error(
3040 "decision archive bucket page disagrees with its cached commitment",
3041 ));
3042 }
3043 decision_pages.push(page);
3044 }
3045 let decision_bucket_directory = decision_bucket_pages.into_iter().collect::<Vec<_>>();
3046 let mut archive_directory_page_ids = Vec::with_capacity(
3047 decision_bucket_directory
3048 .len()
3049 .div_ceil(MAX_PAGED_DECISION_DIRECTORY_ENTRIES),
3050 );
3051 for chunk in decision_bucket_directory.chunks(MAX_PAGED_DECISION_DIRECTORY_ENTRIES) {
3052 let directory_page = PagedDecisionDirectoryPage {
3053 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3054 archive_bucket_pages: chunk.to_vec(),
3055 };
3056 directory_page.validate()?;
3057 let page =
3058 StatePageBlob::new(serde_json::to_vec(&directory_page).map_err(|error| {
3059 invalid_snapshot_error(format!(
3060 "cannot encode paged decision directory page: {error}"
3061 ))
3062 })?)?;
3063 archive_directory_page_ids.push(page.page_id.clone());
3064 decision_pages.push(page);
3065 }
3066 let archive_receipt_count = u64::try_from(decisions.archived_history_count())
3067 .map_err(|_| invalid_snapshot_error("decision archive receipt count exceeds u64"))?;
3068 let decision_manifest_page = StatePageBlob::new(
3069 serde_json::to_vec(&PagedDecisionManifest {
3070 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3071 hot_page_id: hot_decision_page.page_id.clone(),
3072 archive_receipt_root: self
3073 .state
3074 .current
3075 .decisions
3076 .archive_receipt_commitment()
3077 .map_err(|error| invalid_snapshot_error(error.to_string()))?,
3078 archive_receipt_count,
3079 archive_directory_page_ids,
3080 })
3081 .map_err(|error| {
3082 invalid_snapshot_error(format!("cannot encode paged decision manifest: {error}"))
3083 })?,
3084 )?;
3085 let checkpoint_without_paged_state = self.checkpoint_with_paged_payloads(false)?;
3086 let envelope = PagedCheckpointEnvelope {
3087 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3088 checkpoint_without_paged_state,
3089 domain_records,
3090 decision_manifest_page_id: decision_manifest_page.page_id.clone(),
3091 };
3092 let envelope_page =
3093 StatePageBlob::new(serde_json::to_vec(&envelope).map_err(|error| {
3094 invalid_snapshot_error(format!("cannot encode paged checkpoint envelope: {error}"))
3095 })?)?;
3096 let checkpoint = PagedSimulationCheckpoint {
3097 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3098 root_page_id: envelope_page.page_id.clone(),
3099 checkpoint_hash: envelope
3100 .checkpoint_without_paged_state
3101 .state
3102 .checkpoint_hash
3103 .clone(),
3104 };
3105 let mut pages = BTreeMap::new();
3106 for page in domain_pages.into_iter().chain(decision_pages).chain([
3107 hot_decision_page,
3108 decision_manifest_page,
3109 envelope_page,
3110 ]) {
3111 if let Some(existing) = pages.insert(page.page_id.clone(), page.clone())
3112 && existing != page
3113 {
3114 return Err(invalid_snapshot_error(
3115 "one state-page ID resolved to conflicting canonical bytes",
3116 ));
3117 }
3118 }
3119 Ok((checkpoint, pages.into_values().collect()))
3120 }
3121
3122 pub fn prepare_paged_checkpoint(
3126 &self,
3127 source: Option<&PagedSimulationCheckpoint>,
3128 provider: &dyn StatePageProvider,
3129 ) -> Result<PreparedPagedSimulationCheckpoint, CanwuError> {
3130 if let Some(source) = source {
3131 if source.format_version != PAGED_CHECKPOINT_FORMAT_VERSION {
3132 return Err(invalid_snapshot_error(
3133 "paged checkpoint source uses an unsupported format",
3134 ));
3135 }
3136 let page = provider
3137 .load_state_page(&source.root_page_id)?
3138 .ok_or_else(|| {
3139 CanwuError::new(
3140 ErrorCode::StatePageUnavailable,
3141 "paged checkpoint source root is unavailable",
3142 )
3143 })?;
3144 page.validate()?;
3145 if page.page_id != source.root_page_id {
3146 return Err(invalid_snapshot_error(
3147 "paged checkpoint source provider returned the wrong root",
3148 ));
3149 }
3150 }
3151 let (checkpoint, pages) = self.build_paged_checkpoint_pages(Some(provider))?;
3152 let mut new_pages = Vec::new();
3153 for page in pages {
3154 match provider.load_state_page(&page.page_id)? {
3155 Some(existing) => {
3156 existing.validate()?;
3157 if existing != page {
3158 return Err(invalid_snapshot_error(
3159 "state-page provider contains conflicting canonical bytes",
3160 ));
3161 }
3162 }
3163 None => new_pages.push(page),
3164 }
3165 }
3166 let source_root = source.map_or_else(
3167 || canonical_byte_hash("canwu.paged-checkpoint.genesis.v1", &[]),
3168 |source| source.root_page_id.clone(),
3169 );
3170 let delta = prepare_state_delta(&source_root, &checkpoint.root_page_id, new_pages)?;
3171 Ok(PreparedPagedSimulationCheckpoint { checkpoint, delta })
3172 }
3173
3174 pub fn portable_paged_checkpoint(
3177 &self,
3178 ) -> Result<PortablePagedSimulationCheckpoint, CanwuError> {
3179 let (checkpoint, pages) = self.build_paged_checkpoint_pages(None)?;
3180 Ok(PortablePagedSimulationCheckpoint { checkpoint, pages })
3181 }
3182
3183 pub fn from_paged_checkpoint(
3187 checkpoint: &PagedSimulationCheckpoint,
3188 provider: &dyn StatePageProvider,
3189 ) -> Result<Self, CanwuError> {
3190 Self::from_paged_checkpoint_and_journal(checkpoint, provider, Vec::new())
3191 }
3192
3193 pub fn from_paged_checkpoint_and_journal(
3196 checkpoint: &PagedSimulationCheckpoint,
3197 provider: &dyn StatePageProvider,
3198 segments: Vec<EvidenceJournalSegment>,
3199 ) -> Result<Self, CanwuError> {
3200 if checkpoint.format_version != PAGED_CHECKPOINT_FORMAT_VERSION {
3201 return Err(invalid_snapshot_error(
3202 "paged checkpoint uses an unsupported format",
3203 ));
3204 }
3205 let root_page = provider
3206 .load_state_page(&checkpoint.root_page_id)?
3207 .ok_or_else(|| {
3208 CanwuError::new(
3209 ErrorCode::StatePageUnavailable,
3210 "paged checkpoint root is unavailable",
3211 )
3212 })?;
3213 root_page.validate()?;
3214 if root_page.page_id != checkpoint.root_page_id {
3215 return Err(invalid_snapshot_error(
3216 "paged checkpoint provider returned the wrong root page",
3217 ));
3218 }
3219 let envelope: PagedCheckpointEnvelope =
3220 serde_json::from_slice(&root_page.bytes).map_err(|error| {
3221 invalid_snapshot_error(format!("invalid paged checkpoint envelope: {error}"))
3222 })?;
3223 if envelope.format_version != PAGED_CHECKPOINT_FORMAT_VERSION
3224 || envelope
3225 .checkpoint_without_paged_state
3226 .state
3227 .checkpoint_hash
3228 != checkpoint.checkpoint_hash
3229 || !envelope
3230 .checkpoint_without_paged_state
3231 .state
3232 .domain_records
3233 .is_empty()
3234 || !envelope
3235 .checkpoint_without_paged_state
3236 .state
3237 .decisions
3238 .is_empty()
3239 {
3240 return Err(invalid_snapshot_error(
3241 "paged checkpoint envelope metadata is inconsistent",
3242 ));
3243 }
3244 let decision_manifest_page = provider
3245 .load_state_page(&envelope.decision_manifest_page_id)?
3246 .ok_or_else(|| {
3247 CanwuError::new(
3248 ErrorCode::StatePageUnavailable,
3249 "paged decision manifest is unavailable",
3250 )
3251 })?;
3252 decision_manifest_page.validate()?;
3253 if decision_manifest_page.page_id != envelope.decision_manifest_page_id {
3254 return Err(invalid_snapshot_error(
3255 "paged decision provider returned the wrong manifest page",
3256 ));
3257 }
3258 let decision_manifest: PagedDecisionManifest =
3259 serde_json::from_slice(&decision_manifest_page.bytes).map_err(|error| {
3260 invalid_snapshot_error(format!("invalid paged decision manifest: {error}"))
3261 })?;
3262 validate_paged_decision_manifest(&decision_manifest)?;
3263 let hot_decision_page = provider
3264 .load_state_page(&decision_manifest.hot_page_id)?
3265 .ok_or_else(|| {
3266 CanwuError::new(
3267 ErrorCode::StatePageUnavailable,
3268 "paged hot decision state is unavailable",
3269 )
3270 })?;
3271 hot_decision_page.validate()?;
3272 if hot_decision_page.page_id != decision_manifest.hot_page_id {
3273 return Err(invalid_snapshot_error(
3274 "paged decision provider returned the wrong hot-state page",
3275 ));
3276 }
3277 let hot_decisions: DecisionState = serde_json::from_slice(&hot_decision_page.bytes)
3278 .map_err(|error| {
3279 invalid_snapshot_error(format!("invalid paged hot decision state: {error}"))
3280 })?;
3281 let mut directory_pages =
3282 Vec::with_capacity(decision_manifest.archive_directory_page_ids.len());
3283 for directory_page_id in &decision_manifest.archive_directory_page_ids {
3284 let directory_state_page =
3285 provider
3286 .load_state_page(directory_page_id)?
3287 .ok_or_else(|| {
3288 CanwuError::new(
3289 ErrorCode::StatePageUnavailable,
3290 "paged decision directory page is unavailable",
3291 )
3292 })?;
3293 directory_state_page.validate()?;
3294 if directory_state_page.page_id != *directory_page_id {
3295 return Err(invalid_snapshot_error(
3296 "paged decision provider returned the wrong directory page",
3297 ));
3298 }
3299 let directory_page: PagedDecisionDirectoryPage =
3300 serde_json::from_slice(&directory_state_page.bytes).map_err(|error| {
3301 invalid_snapshot_error(format!(
3302 "invalid paged decision directory page: {error}"
3303 ))
3304 })?;
3305 directory_pages.push(directory_page);
3306 }
3307 let archive_bucket_pages = assemble_paged_decision_directory(directory_pages)?;
3308 let required_dependency_pages = hot_decisions
3309 .required_archived_dependency_page_keys()
3310 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3311 let mut resident_dependency_pages = Vec::with_capacity(required_dependency_pages.len());
3312 for page_key in required_dependency_pages {
3313 let page_id = archive_bucket_pages.get(&page_key).ok_or_else(|| {
3314 invalid_snapshot_error(
3315 "hot decision state references history absent from the archive directory",
3316 )
3317 })?;
3318 let state_page = provider.load_state_page(page_id)?.ok_or_else(|| {
3319 CanwuError::new(
3320 ErrorCode::StatePageUnavailable,
3321 "decision dependency locator page is unavailable",
3322 )
3323 })?;
3324 state_page.validate()?;
3325 if state_page.page_id != *page_id {
3326 return Err(invalid_snapshot_error(
3327 "paged decision provider returned the wrong dependency locator page",
3328 ));
3329 }
3330 let page: DecisionArchiveBucketPage = serde_json::from_slice(&state_page.bytes)
3331 .map_err(|error| {
3332 invalid_snapshot_error(format!(
3333 "invalid decision dependency locator page: {error}"
3334 ))
3335 })?;
3336 page.validate()
3337 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3338 if page.bucket != page_key.bucket
3339 || page.segment != page_key.segment
3340 || page
3341 .state_page_id()
3342 .map_err(|error| invalid_snapshot_error(error.to_string()))?
3343 != *page_id
3344 {
3345 return Err(invalid_snapshot_error(
3346 "decision dependency locator page disagrees with the committed directory",
3347 ));
3348 }
3349 resident_dependency_pages.push(page);
3350 }
3351 let decisions = DecisionState::from_paged_checkpoint_root_with_resident_pages(
3352 hot_decisions,
3353 archive_bucket_pages.into(),
3354 decision_manifest.archive_receipt_count,
3355 &decision_manifest.archive_receipt_root,
3356 resident_dependency_pages,
3357 )
3358 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3359 let domain_records = super::PersistentDomainRecordStore::from_state_pages(
3360 &envelope.domain_records,
3361 provider,
3362 )?;
3363 let mut compact_checkpoint = envelope.checkpoint_without_paged_state;
3364 compact_checkpoint.state.domain_records = domain_records.values().cloned().collect();
3365 compact_checkpoint.state.decisions = decisions;
3366 let mut simulation = Self::from_checkpoint_and_journal(compact_checkpoint, segments)?;
3367 simulation.state.current.domain_records = domain_records;
3368 Ok(simulation)
3369 }
3370
3371 pub fn from_portable_paged_checkpoint(
3372 portable: PortablePagedSimulationCheckpoint,
3373 ) -> Result<Self, CanwuError> {
3374 let mut provider = EmbeddedPageProvider::default();
3375 for page in portable.pages {
3376 page.validate()?;
3377 if let Some(existing) = provider.pages.insert(page.page_id.clone(), page.clone())
3378 && existing != page
3379 {
3380 return Err(invalid_snapshot_error(
3381 "portable paged checkpoint contains conflicting duplicate pages",
3382 ));
3383 }
3384 }
3385 Self::from_paged_checkpoint(&portable.checkpoint, &provider)
3386 }
3387
3388 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
3390 EvidenceCursor::from_evidence(&self.state.evidence)
3391 }
3392
3393 fn checkpoint_with_paged_payloads(
3394 &self,
3395 include_paged_payloads: bool,
3396 ) -> Result<SimulationCheckpoint, CanwuError> {
3397 let archived_segment_headers = self.state.evidence.archived_segment_headers.clone();
3398 let evidence_dependencies = self.evidence_dependencies()?;
3399 let mut keyed_draw_reservations = self.state.evidence.keyed_draw_reservations.clone();
3400 keyed_draw_reservations.sort_by(|left, right| {
3401 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
3402 });
3403 let required_receipts =
3404 required_archived_receipt_references(&evidence_dependencies, &keyed_draw_reservations);
3405 let archived_evidence_receipts: Vec<_> = self
3406 .state
3407 .evidence
3408 .archived_evidence_receipts
3409 .iter()
3410 .filter(|(reference, _)| required_receipts.contains(*reference))
3411 .map(|(_, receipt)| receipt.clone())
3412 .collect();
3413 let checkpoint = SimulationCheckpoint {
3414 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
3415 journal_end: self.evidence_cursor()?,
3416 state: self.checkpoint_state_with_paged_payloads(include_paged_payloads),
3417 archived_segment_manifest_root: skipped_commitment_root(
3418 ARCHIVED_SEGMENT_MANIFEST_DOMAIN,
3419 &archived_segment_headers,
3420 )?,
3421 archived_segment_headers,
3422 archived_receipt_root: skipped_commitment_root(
3423 ARCHIVED_RECEIPT_DOMAIN,
3424 &archived_evidence_receipts,
3425 )?,
3426 archived_evidence_receipts,
3427 evidence_dependency_root: skipped_commitment_root(
3428 EVIDENCE_DEPENDENCY_DOMAIN,
3429 &evidence_dependencies,
3430 )?,
3431 evidence_dependencies,
3432 keyed_reservation_root: skipped_commitment_root(
3433 KEYED_RESERVATION_DOMAIN,
3434 &keyed_draw_reservations,
3435 )?,
3436 keyed_draw_reservations,
3437 };
3438 if include_paged_payloads {
3439 validate_compact_continuation(&checkpoint)?;
3440 }
3441 Ok(checkpoint)
3442 }
3443
3444 pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
3446 self.checkpoint_with_paged_payloads(true)
3447 }
3448
3449 pub fn archive_reachability_manifest(
3454 &self,
3455 retained_checkpoints: &[SimulationCheckpoint],
3456 page_retention: &StatePageRetentionLedger,
3457 decision_provider: &dyn DecisionArchiveProvider,
3458 plugin_provider: &dyn super::PluginArchiveObjectProvider,
3459 ) -> Result<ArchiveReachabilityManifest, CanwuError> {
3460 let mut manifest = ArchiveReachabilityManifest {
3461 state_page_ids: page_retention.reachable_page_ids(),
3462 evidence_segment_ids: SimulationCheckpoint::reachable_archive_segment_ids(
3463 retained_checkpoints,
3464 )?,
3465 ..ArchiveReachabilityManifest::default()
3466 };
3467 manifest.evidence_segment_ids.extend(
3468 self.state
3469 .evidence
3470 .archived_segment_headers
3471 .iter()
3472 .map(|header| header.segment_id.clone()),
3473 );
3474 let decisions = self
3475 .state
3476 .current
3477 .decisions
3478 .archive_reachability(decision_provider)
3479 .map_err(super::decision::decision_error)?;
3480 manifest.state_page_ids.extend(decisions.bucket_page_ids);
3481 manifest.decision_blob_ids.extend(decisions.blob_locators);
3482 let pending_ingress_ids = self
3483 .state
3484 .scheduler
3485 .pending_ingress
3486 .iter()
3487 .map(|key| key.id)
3488 .collect::<BTreeSet<_>>();
3489 for record in &self.state.evidence.ingress {
3490 if let IngressPayload::Maintenance { request } = &record.payload
3491 && let super::MaintenanceIngressRequest::DecisionArchive { commit } =
3492 request.as_ref()
3493 {
3494 manifest
3495 .decision_blob_ids
3496 .extend(commit.archive_locators().map(str::to_owned));
3497 }
3498 if pending_ingress_ids.contains(&record.id)
3499 && let IngressPayload::Plugin {
3500 archive_retention, ..
3501 } = &record.payload
3502 {
3503 for retention in archive_retention {
3504 manifest.insert_plugin_object(
3505 retention.namespace.clone(),
3506 retention.object_id.clone(),
3507 );
3508 }
3509 }
3510 }
3511 for (plugin, participant) in &self.plugins.archive_reachability_participants {
3512 let reads = self
3513 .plugins
3514 .state_owners
3515 .iter()
3516 .filter_map(|(state, owner)| (owner == plugin).then_some(state.clone()))
3517 .collect::<Vec<_>>();
3518 let reader = format!("{plugin}.archive_reachability");
3519 let view = self.plugin_view(&reader, &reads);
3520 participant(&view, plugin_provider, &mut manifest)?;
3521 }
3522 Ok(manifest)
3523 }
3524
3525 pub fn journal_segment_since(
3527 &self,
3528 start: EvidenceCursor,
3529 ) -> Result<EvidenceJournalSegment, CanwuError> {
3530 let end = self.evidence_cursor()?;
3531 let cut = |value: u64, archived: u64, len: usize, label: &str| {
3532 let value = value.checked_sub(archived).ok_or_else(|| {
3533 CanwuError::new(
3534 ErrorCode::InvalidSnapshot,
3535 format!("{label} journal cursor precedes the retained live evidence window"),
3536 )
3537 })?;
3538 let value = usize::try_from(value).map_err(|_| {
3539 CanwuError::new(
3540 ErrorCode::InvalidSnapshot,
3541 format!("{label} journal cursor is not representable on this platform"),
3542 )
3543 })?;
3544 if value > len {
3545 return Err(CanwuError::new(
3546 ErrorCode::InvalidSnapshot,
3547 format!("{label} journal cursor exceeds the current evidence tail"),
3548 ));
3549 }
3550 Ok(value)
3551 };
3552 let archived = self.state.evidence.archived;
3553 let event_start = cut(
3554 start.event_count,
3555 archived.event_count,
3556 self.state.evidence.events.len(),
3557 "event",
3558 )?;
3559 let command_start = cut(
3560 start.command_count,
3561 archived.command_count,
3562 self.state.evidence.commands.len(),
3563 "command",
3564 )?;
3565 let attempt_start = cut(
3566 start.command_attempt_count,
3567 archived.command_attempt_count,
3568 self.state.evidence.command_attempts.len(),
3569 "command-attempt",
3570 )?;
3571 let ingress_start = cut(
3572 start.ingress_count,
3573 archived.ingress_count,
3574 self.state.evidence.ingress.len(),
3575 "ingress",
3576 )?;
3577 let boundary_start = cut(
3578 start.boundary_count,
3579 archived.boundary_count,
3580 self.state.evidence.boundaries.len(),
3581 "boundary",
3582 )?;
3583 let draw_start = cut(
3584 start.random_draw_count,
3585 archived.random_draw_count,
3586 self.state.evidence.random_draws.len(),
3587 "random-draw",
3588 )?;
3589 Ok(EvidenceJournalSegment {
3590 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
3591 start,
3592 end,
3593 events: self.state.evidence.events[event_start..].to_vec(),
3594 commands: self.state.evidence.commands[command_start..].to_vec(),
3595 command_attempts: self.state.evidence.command_attempts[attempt_start..].to_vec(),
3596 ingress: self.state.evidence.ingress[ingress_start..].to_vec(),
3597 boundaries: self.state.evidence.boundaries[boundary_start..].to_vec(),
3598 random_draws: self.state.evidence.random_draws[draw_start..].to_vec(),
3599 archive: None,
3600 })
3601 }
3602
3603 pub fn checkpoint_journal(&self) -> Result<CheckpointJournal, CanwuError> {
3605 if self.state.evidence.archived != EvidenceCursor::default() {
3606 return Err(CanwuError::new(
3607 ErrorCode::InvalidSnapshot,
3608 "a compact live runtime requires its previously sealed evidence segments to build a portable save",
3609 ));
3610 }
3611 let segment = self.journal_segment_since(EvidenceCursor::default())?;
3612 Ok(CheckpointJournal {
3613 checkpoint: self.checkpoint()?,
3614 segments: (segment.start != segment.end)
3615 .then_some(segment)
3616 .into_iter()
3617 .collect(),
3618 })
3619 }
3620
3621 pub fn checkpoint_journal_json(&self) -> Result<String, CanwuError> {
3623 serde_json::to_string_pretty(&self.checkpoint_journal()?).map_err(|error| {
3624 CanwuError::new(
3625 ErrorCode::InvalidSnapshot,
3626 format!("could not serialize checkpoint journal: {error}"),
3627 )
3628 })
3629 }
3630
3631 fn snapshot_from_checkpoint_and_journal(
3632 checkpoint: SimulationCheckpoint,
3633 segments: Vec<EvidenceJournalSegment>,
3634 ) -> Result<SimulationSnapshot, CanwuError> {
3635 if checkpoint.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
3636 return Err(invalid_snapshot_error(format!(
3637 "checkpoint-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
3638 checkpoint.format_version
3639 )));
3640 }
3641 validate_compact_continuation(&checkpoint)?;
3642 let expected_headers = checkpoint.archived_segment_headers;
3643 let expected_receipts = checkpoint.archived_evidence_receipts;
3644 let expected_dependencies = checkpoint.evidence_dependencies;
3645 let expected_reservations = checkpoint.keyed_draw_reservations;
3646 let mut snapshot = checkpoint.state;
3647 if snapshot.snapshot_format_version != SNAPSHOT_FORMAT_VERSION {
3648 return Err(invalid_snapshot_error(format!(
3649 "checkpoint-journal format {CHECKPOINT_JOURNAL_FORMAT_VERSION} requires snapshot format {SNAPSHOT_FORMAT_VERSION}"
3650 )));
3651 }
3652 if !snapshot.events.is_empty()
3653 || !snapshot.commands.is_empty()
3654 || !snapshot.command_attempts.is_empty()
3655 || !snapshot.ingress.is_empty()
3656 || !snapshot.boundaries.is_empty()
3657 || !snapshot.random_draws.is_empty()
3658 {
3659 return Err(invalid_snapshot_error(
3660 "checkpoint current state must not duplicate append-only evidence",
3661 ));
3662 }
3663
3664 let mut cursor = EvidenceCursor::default();
3665 let mut rebuilt_headers = Vec::new();
3666 let mut rebuilt_receipts = BTreeMap::new();
3667 let mut rebuilt_reservations = Vec::new();
3668 for segment in segments {
3669 if segment.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
3670 return Err(invalid_snapshot_error(format!(
3671 "evidence-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
3672 segment.format_version
3673 )));
3674 }
3675 if segment.start != cursor {
3676 return Err(invalid_snapshot_error(
3677 "evidence-journal segments must form one contiguous global prefix",
3678 ));
3679 }
3680 let end = cursor.checked_advance(&segment)?;
3681 if end == cursor {
3682 return Err(invalid_snapshot_error(
3683 "evidence-journal segments must advance at least one journal cursor",
3684 ));
3685 }
3686 if segment.end != end {
3687 return Err(invalid_snapshot_error(
3688 "evidence-journal segment end does not match its encoded records",
3689 ));
3690 }
3691 if let Some(archive) = &segment.archive {
3692 let receipts = verify_archived_segment(&segment).map_err(|error| {
3693 invalid_snapshot_error(format!(
3694 "archived evidence segment is invalid: {}",
3695 error.message
3696 ))
3697 })?;
3698 rebuilt_headers.push(archive.header.clone());
3699 for receipt in receipts {
3700 if rebuilt_receipts
3701 .insert(receipt.evidence.clone(), receipt)
3702 .is_some()
3703 {
3704 return Err(invalid_snapshot_error(
3705 "archived evidence segments contain duplicate receipts",
3706 ));
3707 }
3708 }
3709 for draw in &segment.random_draws {
3710 let RandomDrawAddress::OperationV1(address) = &draw.address else {
3711 continue;
3712 };
3713 let operation_evidence = draw.operation_evidence.clone().ok_or_else(|| {
3714 invalid_snapshot_error("archived keyed draw is missing operation evidence")
3715 })?;
3716 let draw_receipt = rebuilt_receipts
3717 .get(&EvidenceRef::RandomDraw(draw.id))
3718 .cloned()
3719 .ok_or_else(|| {
3720 invalid_snapshot_error("archived keyed draw receipt is missing")
3721 })?;
3722 rebuilt_reservations.push(KeyedDrawReservation {
3723 stream: draw.stream.clone(),
3724 address: address.clone(),
3725 upper_exclusive: draw.upper_exclusive,
3726 purpose_hash: super::random::purpose_hash_hex_v1(&draw.purpose)?,
3727 result: draw.value,
3728 draw_id: draw.id,
3729 operation_evidence,
3730 draw_receipt,
3731 });
3732 }
3733 }
3734 snapshot.events.extend(segment.events);
3735 snapshot.commands.extend(segment.commands);
3736 snapshot.command_attempts.extend(segment.command_attempts);
3737 snapshot.ingress.extend(segment.ingress);
3738 snapshot.boundaries.extend(segment.boundaries);
3739 snapshot.random_draws.extend(segment.random_draws);
3740 cursor = end;
3741 }
3742 if cursor != checkpoint.journal_end {
3743 return Err(invalid_snapshot_error(
3744 "evidence-journal segments do not reach the checkpoint journal cut",
3745 ));
3746 }
3747 rebuilt_reservations.sort_by(|left, right| {
3748 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
3749 });
3750 let rebuilt_dependencies =
3751 Simulation::from_snapshot(snapshot.clone())?.evidence_dependencies()?;
3752 retain_reachable_archived_evidence_receipts(
3753 &mut rebuilt_receipts,
3754 &rebuilt_dependencies,
3755 &rebuilt_reservations,
3756 );
3757 let rebuilt_receipts: Vec<_> = rebuilt_receipts.into_values().collect();
3758 if rebuilt_headers != expected_headers
3759 || rebuilt_receipts != expected_receipts
3760 || rebuilt_dependencies != expected_dependencies
3761 || rebuilt_reservations != expected_reservations
3762 {
3763 return Err(invalid_snapshot_error(
3764 "checkpoint compact continuation does not match its reconstructed authoritative indexes",
3765 ));
3766 }
3767 Ok(snapshot)
3768 }
3769
3770 pub fn from_checkpoint_and_journal(
3772 checkpoint: SimulationCheckpoint,
3773 segments: Vec<EvidenceJournalSegment>,
3774 ) -> Result<Self, CanwuError> {
3775 Self::from_snapshot(Self::snapshot_from_checkpoint_and_journal(
3776 checkpoint, segments,
3777 )?)
3778 }
3779
3780 pub fn from_checkpoint_journal(bundle: CheckpointJournal) -> Result<Self, CanwuError> {
3782 Self::from_checkpoint_and_journal(bundle.checkpoint, bundle.segments)
3783 }
3784
3785 pub fn from_checkpoint_journal_with_plugins(
3787 bundle: CheckpointJournal,
3788 plugins: &[&dyn SimulationPlugin],
3789 ) -> Result<Self, CanwuError> {
3790 let mut simulation = Self::from_checkpoint_journal(bundle)?;
3791 for plugin in plugins {
3792 simulation.register_plugin(*plugin)?;
3793 }
3794 simulation.ensure_runtime_ready()?;
3795 Ok(simulation)
3796 }
3797
3798 pub fn from_checkpoint_journal_json(json: &str) -> Result<Self, CanwuError> {
3800 let bundle: CheckpointJournal =
3801 super::deserialize_current_json(json, "checkpoint journal")?;
3802 Self::from_checkpoint_journal(bundle)
3803 }
3804
3805 pub fn from_checkpoint_journal_json_with_plugins(
3807 json: &str,
3808 plugins: &[&dyn SimulationPlugin],
3809 ) -> Result<Self, CanwuError> {
3810 let bundle: CheckpointJournal =
3811 super::deserialize_current_json(json, "checkpoint journal")?;
3812 Self::from_checkpoint_journal_with_plugins(bundle, plugins)
3813 }
3814}
3815
3816pub fn format8_paged_checkpoint_scale_probe(
3820 decision_count: usize,
3821) -> Result<PagedCheckpointScaleMetrics, CanwuError> {
3822 let fixture = canwu_decision::format8_decision_locator_scale_fixture(decision_count)
3823 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3824 let decision_metrics = fixture.metrics.clone();
3825 let decision_blobs = fixture
3826 .archive_blobs
3827 .into_iter()
3828 .map(|blob| {
3829 blob.content_id()
3830 .map(|content_id| (content_id, blob))
3831 .map_err(|error| invalid_snapshot_error(error.to_string()))
3832 })
3833 .collect::<Result<BTreeMap<_, _>, _>>()?;
3834 let store = Format8ScalePageStore {
3835 pages: RefCell::new(BTreeMap::new()),
3836 decision_blobs: RefCell::new(decision_blobs),
3837 provider_calls: RefCell::new(0),
3838 };
3839 let mut simulation = Simulation::new(8, Scenario::new(SimTime::EPOCH, Vec::new()))?;
3840 simulation.state.current.decisions = fixture.state;
3841 simulation.state.metadata.plugin_registration_closed = true;
3842 simulation.state.metadata.commitment_cache = None;
3843 simulation.refresh_checkpoint_hash()?;
3844
3845 let prepared = simulation.prepare_paged_checkpoint(None, &store)?;
3846 let initial_provider_calls = *store.provider_calls.borrow();
3847 let initial_delta_pages = prepared.delta.new_pages.len() as u64;
3848 prepared.store_and_verify(&store)?;
3849 let root_page = store
3850 .load_state_page(&prepared.checkpoint.root_page_id)?
3851 .ok_or_else(|| invalid_snapshot_error("Format-8 scale root page is unavailable"))?;
3852 let envelope: PagedCheckpointEnvelope = serde_json::from_slice(&root_page.bytes)
3853 .map_err(|error| invalid_snapshot_error(format!("invalid scale envelope: {error}")))?;
3854 let manifest_page = store
3855 .load_state_page(&envelope.decision_manifest_page_id)?
3856 .ok_or_else(|| invalid_snapshot_error("Format-8 scale manifest page is unavailable"))?;
3857 let manifest: PagedDecisionManifest = serde_json::from_slice(&manifest_page.bytes)
3858 .map_err(|error| invalid_snapshot_error(format!("invalid scale manifest: {error}")))?;
3859
3860 let restored = Simulation::from_paged_checkpoint(&prepared.checkpoint, &store)?;
3861 let replayed =
3862 Simulation::from_paged_checkpoint_and_journal(&prepared.checkpoint, &store, Vec::new())?;
3863 let samples = [
3864 1_usize,
3865 decision_count.saturating_div(2).max(1),
3866 decision_count,
3867 ]
3868 .into_iter()
3869 .filter(|ordinal| *ordinal <= decision_count)
3870 .collect::<BTreeSet<_>>();
3871 for ordinal in &samples {
3872 let key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(
3873 u64::try_from(*ordinal)
3874 .map_err(|_| invalid_snapshot_error("decision scale sample exceeds u64"))?,
3875 ));
3876 if !matches!(
3877 restored.decision_history_location_with_provider(&key, &store)?,
3878 DecisionHistoryLocation::Archived { .. }
3879 ) || !matches!(
3880 replayed.decision_history_location_with_provider(&key, &store)?,
3881 DecisionHistoryLocation::Archived { .. }
3882 ) {
3883 return Err(invalid_snapshot_error(
3884 "paged checkpoint scale restore lost exact decision history",
3885 ));
3886 }
3887 }
3888 *store.provider_calls.borrow_mut() = 0;
3889 let repeat = restored.prepare_paged_checkpoint(Some(&prepared.checkpoint), &store)?;
3890 let repeat_provider_calls = *store.provider_calls.borrow();
3891 let replay_repeat = replayed.prepare_paged_checkpoint(Some(&prepared.checkpoint), &store)?;
3892 let mut changed = restored;
3893 let changed_request_id = canwu_core::DecisionRequestId::new(
3894 u64::try_from(decision_count)
3895 .map_err(|_| invalid_snapshot_error("decision scale count exceeds u64"))?
3896 .checked_add(1)
3897 .ok_or_else(|| invalid_snapshot_error("decision scale request ID overflowed"))?,
3898 );
3899 changed
3900 .state
3901 .current
3902 .decisions
3903 .append_attempt(super::DecisionAttemptRecord {
3904 request_id: changed_request_id,
3905 request_commitment: canonical_byte_hash(
3906 "canwu.format8.single-page-change.v1",
3907 &changed_request_id.get().to_be_bytes(),
3908 ),
3909 at: SimTime::from_minutes(
3910 i64::try_from(changed_request_id.get())
3911 .map_err(|_| invalid_snapshot_error("decision scale time exceeds i64"))?,
3912 ),
3913 revision_before: 0,
3914 expected_revision: 0,
3915 outcome: super::DecisionAttemptOutcome::Rejected {
3916 code: super::DecisionAttemptErrorCode::InvalidDecision,
3917 message: "Format-8 single locator-page change".to_owned(),
3918 },
3919 })
3920 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3921 let changed_key = DecisionHistoryKey::Attempt(changed_request_id);
3922 let changed_archive = changed
3923 .state
3924 .current
3925 .decisions
3926 .prepare_decision_archive(std::slice::from_ref(&changed_key))
3927 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3928 for blob in &changed_archive.blobs {
3929 store.decision_blobs.borrow_mut().insert(
3930 blob.content_id()
3931 .map_err(|error| invalid_snapshot_error(error.to_string()))?,
3932 blob.clone(),
3933 );
3934 }
3935 let verified_change = changed
3936 .state
3937 .current
3938 .decisions
3939 .verify_decision_archive(&changed_archive, &store)
3940 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3941 changed.state.current.decisions = changed
3942 .state
3943 .current
3944 .decisions
3945 .commit_verified_decision_archive(&verified_change)
3946 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3947 changed.state.metadata.commitment_cache = None;
3948 changed.refresh_checkpoint_hash()?;
3949 *store.provider_calls.borrow_mut() = 0;
3950 let single_page_change =
3951 changed.prepare_paged_checkpoint(Some(&prepared.checkpoint), &store)?;
3952 let single_page_change_provider_calls = *store.provider_calls.borrow();
3953 let pages = store.pages.borrow();
3954 Ok(PagedCheckpointScaleMetrics {
3955 decision_entries: decision_metrics.entries,
3956 decision_locator: decision_metrics,
3957 state_pages: pages.len() as u64,
3958 decision_directory_pages: manifest.archive_directory_page_ids.len() as u64,
3959 max_state_page_bytes: pages
3960 .values()
3961 .map(|page| page.decoded_bytes)
3962 .max()
3963 .unwrap_or(0),
3964 initial_delta_pages,
3965 repeat_delta_pages: repeat.delta.new_pages.len() as u64,
3966 single_page_change_delta_pages: single_page_change.delta.new_pages.len() as u64,
3967 initial_provider_calls,
3968 repeat_provider_calls,
3969 single_page_change_provider_calls,
3970 exact_restart_queries: samples.len() as u64,
3971 restored_root_matches: repeat.checkpoint.root_page_id == prepared.checkpoint.root_page_id,
3972 replayed_root_matches: replay_repeat.checkpoint.root_page_id
3973 == prepared.checkpoint.root_page_id,
3974 root_page_id: prepared.checkpoint.root_page_id,
3975 })
3976}
3977
3978#[cfg(test)]
3979mod tests {
3980 #![allow(clippy::unnecessary_literal_bound, clippy::unnecessary_wraps)]
3981 use super::super::{
3982 BoundaryContext, BoundaryDirective, BoundaryPhase, BoundaryProposal,
3983 BoundarySystemContract, Command, CommandContext, CommandRequestId, DomainRecordClass,
3984 DomainRecordDraft, DomainRecordMutation, Issuer, KnowledgeHolderRef, KnowledgeOrigin,
3985 KnowledgeRecordDraft, KnowledgeRecordKind, KnowledgeSchemaId, KnowledgeWriteGrant,
3986 PluginActionDescriptor, PluginKnowledgeSchema, PluginRegistrar, SimulationView, StateKey,
3987 StateVisibility, SystemDirective, demo_scenario,
3988 };
3989 use super::*;
3990 use canwu_core::{BoundaryId, CommandId, DomainRecordKind, PersonId};
3991 use serde_json::{Map, Value, json};
3992 use std::cell::RefCell;
3993
3994 fn decision_directory_page(start: u32, len: usize) -> PagedDecisionDirectoryPage {
3995 PagedDecisionDirectoryPage {
3996 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3997 archive_bucket_pages: (0..len)
3998 .map(|offset| {
3999 let ordinal = start + u32::try_from(offset).expect("test offset fits u32");
4000 (
4001 super::super::DecisionArchivePageKey {
4002 bucket: u16::try_from(ordinal / 256).expect("bucket fits"),
4003 segment: u8::try_from(ordinal % 256).expect("segment fits"),
4004 },
4005 canonical_byte_hash(
4006 "canwu.test.decision-directory-page.v1",
4007 &ordinal.to_be_bytes(),
4008 ),
4009 )
4010 })
4011 .collect(),
4012 }
4013 }
4014
4015 #[test]
4016 fn paged_decision_directory_rejects_noncanonical_cross_page_encodings() {
4017 let first = decision_directory_page(0, MAX_PAGED_DECISION_DIRECTORY_ENTRIES);
4018 let second = decision_directory_page(
4019 u32::try_from(MAX_PAGED_DECISION_DIRECTORY_ENTRIES).expect("directory bound fits u32"),
4020 1,
4021 );
4022 assemble_paged_decision_directory(vec![second.clone(), first.clone()])
4023 .expect_err("swapped directory pages must fail");
4024 assemble_paged_decision_directory(vec![decision_directory_page(0, 1), second])
4025 .expect_err("a short middle directory page must fail");
4026
4027 let duplicate_id = canonical_byte_hash("canwu.test.duplicate-directory.v1", b"same");
4028 let manifest = PagedDecisionManifest {
4029 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
4030 hot_page_id: canonical_byte_hash("canwu.test.hot.v1", b"hot"),
4031 archive_receipt_root: canonical_byte_hash("canwu.test.archive.v1", b"archive"),
4032 archive_receipt_count: 1,
4033 archive_directory_page_ids: vec![duplicate_id.clone(), duplicate_id],
4034 };
4035 validate_paged_decision_manifest(&manifest)
4036 .expect_err("duplicate directory page IDs must fail");
4037 }
4038
4039 #[test]
4040 fn archived_plugin_ingress_provenance_is_merkle_bound() {
4041 let provenance = ArchivedPluginIngressProvenance {
4042 plugin: "fixture-provider".to_owned(),
4043 packet_type: "recognized-practice".to_owned(),
4044 producer_boundary: BoundaryId::new(7),
4045 };
4046 let entry = EvidenceIndexEntry {
4047 reference: EvidenceRef::Ingress(super::super::IngressId::new(9)),
4048 item: EvidenceItemLocator {
4049 journal: EvidenceJournalKind::Ingress,
4050 absolute_index: 9,
4051 nested: EvidenceNestedLocator::None,
4052 },
4053 item_commitment: "0101010101010101010101010101010101010101010101010101010101010101"
4054 .to_owned(),
4055 plugin_ingress_provenance: Some(provenance.clone()),
4056 };
4057 let (root, proofs) = archive_merkle(std::slice::from_ref(&entry)).expect("build proof");
4058 let header = ArchivedSegmentHeader {
4059 segment_id: "0202020202020202020202020202020202020202020202020202020202020202"
4060 .to_owned(),
4061 start: EvidenceCursor {
4062 ingress_count: 8,
4063 ..EvidenceCursor::default()
4064 },
4065 end: EvidenceCursor {
4066 ingress_count: 9,
4067 ..EvidenceCursor::default()
4068 },
4069 journal_roots: EvidenceJournalRoots {
4070 events: archive_empty_root(),
4071 commands: archive_empty_root(),
4072 command_attempts: archive_empty_root(),
4073 ingress: archive_empty_root(),
4074 boundaries: archive_empty_root(),
4075 random_draws: archive_empty_root(),
4076 },
4077 evidence_index_root: root,
4078 evidence_index_entry_count: 1,
4079 };
4080 let mut receipt = ArchivedEvidenceReceipt {
4081 evidence: entry.reference,
4082 locator: ArchivedEvidenceLocator {
4083 segment_id: header.segment_id.clone(),
4084 item: entry.item,
4085 },
4086 evidence_index_leaf: 0,
4087 item_commitment: entry.item_commitment,
4088 plugin_ingress_provenance: Some(provenance),
4089 merkle_path: proofs.into_iter().next().expect("one proof"),
4090 };
4091 verify_archive_receipt(&receipt, &header).expect("canonical provenance proof");
4092
4093 receipt
4094 .plugin_ingress_provenance
4095 .as_mut()
4096 .expect("provenance")
4097 .plugin = "forged-provider".to_owned();
4098 assert_eq!(
4099 verify_archive_receipt(&receipt, &header)
4100 .expect_err("tampered provenance must break the proof")
4101 .code,
4102 ErrorCode::InvalidArchive
4103 );
4104 }
4105
4106 #[derive(Default)]
4107 struct TestArchive {
4108 segments: RefCell<BTreeMap<String, EvidenceJournalSegment>>,
4109 }
4110
4111 impl TestArchive {
4112 fn segment_ids(&self) -> Vec<String> {
4113 self.segments.borrow().keys().cloned().collect()
4114 }
4115 }
4116
4117 impl ArchiveProvider for TestArchive {
4118 fn load_evidence_segment(
4119 &self,
4120 segment_id: &str,
4121 ) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
4122 Ok(self.segments.borrow().get(segment_id).cloned())
4123 }
4124 }
4125
4126 impl ArchiveStore for TestArchive {
4127 fn store_evidence_segment(
4128 &self,
4129 segment: &EvidenceJournalSegment,
4130 ) -> Result<ArchiveStoreOutcome, CanwuError> {
4131 let segment_id = segment
4132 .archive
4133 .as_ref()
4134 .ok_or_else(|| archive_error("test archive segment has no index"))?
4135 .header
4136 .segment_id
4137 .clone();
4138 let mut segments = self.segments.borrow_mut();
4139 if let Some(existing) = segments.get(&segment_id) {
4140 return if existing == segment {
4141 Ok(ArchiveStoreOutcome::AlreadyPresent)
4142 } else {
4143 Err(archive_error(
4144 "content-addressed test segment ID has conflicting bytes",
4145 ))
4146 };
4147 }
4148 segments.insert(segment_id, segment.clone());
4149 Ok(ArchiveStoreOutcome::Stored)
4150 }
4151 }
4152
4153 fn archived_identity_schema() -> KnowledgeSchemaId {
4154 KnowledgeSchemaId::new(
4155 KnowledgeRecordKind::new("fixture.archive", "archived_command_notice"),
4156 1,
4157 )
4158 }
4159
4160 #[allow(clippy::unnecessary_wraps)]
4161 fn retain_archive_source_command(
4162 _view: &SimulationView<'_>,
4163 _context: &CommandContext,
4164 _payload: &Value,
4165 ) -> Result<Vec<SystemDirective>, CanwuError> {
4166 Ok(Vec::new())
4167 }
4168
4169 #[allow(clippy::unnecessary_wraps)]
4170 fn publish_archived_command_identity(
4171 _view: &SimulationView<'_>,
4172 context: &BoundaryContext,
4173 ) -> Result<BoundaryProposal, CanwuError> {
4174 let record = KnowledgeRecordDraft {
4175 schema: archived_identity_schema(),
4176 subjects: Vec::new(),
4177 payload: json!({ "boundary": context.boundary_id.get() }),
4178 as_of: None,
4179 confidence_per_mille: 1_000,
4180 origin: KnowledgeOrigin {
4181 method: "archived_command_identity_v1".to_owned(),
4182 evidence: vec![EvidenceRef::Command(CommandId::new(1))],
4183 },
4184 supersedes: Vec::new(),
4185 contradicts: Vec::new(),
4186 };
4187 Ok(BoundaryProposal {
4188 directives: vec![BoundaryDirective::PublishKnowledge {
4189 holder: KnowledgeHolderRef::Person(PersonId::new(1)),
4190 visibility: StateVisibility::SameBoundary,
4191 producer_correlation: Some(format!(
4192 "archived-command-boundary-{}",
4193 context.boundary_id.get()
4194 )),
4195 records: vec![record],
4196 summary: "Publish knowledge from a retained or archived command identity"
4197 .to_owned(),
4198 }],
4199 ..BoundaryProposal::default()
4200 })
4201 }
4202
4203 struct ArchivedIdentityPublicationPlugin;
4204
4205 impl SimulationPlugin for ArchivedIdentityPublicationPlugin {
4206 fn name(&self) -> &str {
4207 "fixture-archived-identity-publication"
4208 }
4209
4210 fn version(&self) -> &str {
4211 "test-v1"
4212 }
4213
4214 fn semantic_hash(&self) -> &str {
4215 "7100000000000000000000000000000000000000000000000000000000000000"
4216 }
4217
4218 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
4219 registrar.register_knowledge_schema(PluginKnowledgeSchema {
4220 id: archived_identity_schema(),
4221 schema_hash: "7200000000000000000000000000000000000000000000000000000000000000"
4222 .to_owned(),
4223 writable: true,
4224 payload_schema: PayloadSchema::Any,
4225 subjects: Vec::new(),
4226 })?;
4227 registrar.register_command(
4228 PluginActionDescriptor {
4229 name: "retain_archive_source_v1".to_owned(),
4230 description: "Persist one neutral command identity for archive testing"
4231 .to_owned(),
4232 payload_schema: PayloadSchema::Any,
4233 reads: Vec::new(),
4234 writes: Vec::new(),
4235 },
4236 retain_archive_source_command,
4237 )?;
4238 let mut publisher = BoundarySystemContract::new(
4239 "publish-archived-command-identity",
4240 BoundaryPhase::PerspectiveAndReportMaterialization,
4241 SystemCadence::Daily,
4242 );
4243 publisher.knowledge_writes = vec![KnowledgeWriteGrant {
4244 schema: archived_identity_schema(),
4245 visibilities: vec![StateVisibility::SameBoundary],
4246 }];
4247 registrar.register_boundary_system(publisher, publish_archived_command_identity)
4248 }
4249 }
4250
4251 fn archived_identity_two_segment_fixture() -> (SimulationCheckpoint, Vec<EvidenceJournalSegment>)
4252 {
4253 let (scenario, _) = demo_scenario();
4254 let plugin = ArchivedIdentityPublicationPlugin;
4255 let mut simulation = Simulation::new(711, scenario).expect("fixture scenario should load");
4256 simulation
4257 .register_plugin(&plugin)
4258 .expect("archived-identity plugin should register");
4259 simulation
4260 .enqueue_command(
4261 SimTime::EPOCH,
4262 0,
4263 CommandRequest::new(
4264 CommandRequestId::new(1),
4265 simulation.revision(),
4266 CommandEnvelope::new(
4267 Issuer::System("archive-identity-fixture".to_owned()),
4268 Command::Plugin {
4269 plugin: plugin.name().to_owned(),
4270 command: "retain_archive_source_v1".to_owned(),
4271 payload: json!({ "format_version": 1 }),
4272 },
4273 )
4274 .at_time(SimTime::EPOCH),
4275 ),
4276 )
4277 .expect("archive source command should queue");
4278 simulation
4279 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
4280 .expect("boundary one should retain the command-backed publication");
4281 simulation
4282 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4283 .expect("the following cut should admit boundary-one events before sealing");
4284
4285 let mut compact = simulation
4286 .into_compacted()
4287 .expect("the fixture should enter compact mode");
4288 let first_segment = compact
4289 .seal_evidence()
4290 .expect("boundary-one evidence should seal")
4291 .expect("boundary one should produce an archive segment");
4292 let first_checkpoint = compact.checkpoint().expect("checkpoint one should build");
4293 assert!(
4294 first_checkpoint
4295 .archived_evidence_receipts
4296 .iter()
4297 .any(|receipt| receipt.evidence == EvidenceRef::Command(CommandId::new(1)))
4298 );
4299
4300 let mut restored = CompactedSimulation::from_checkpoint_and_journal_with_plugins(
4301 first_checkpoint,
4302 vec![first_segment.clone()],
4303 &[&plugin],
4304 )
4305 .expect("checkpoint one should restore with its exact archive prefix");
4306 let restored_prefix = restored
4307 .seal_evidence()
4308 .expect("restored boundary-one evidence should reseal")
4309 .expect("restored boundary-one evidence should remain non-empty");
4310 assert_eq!(restored_prefix, first_segment);
4311 assert!(
4312 restored
4313 .archived_evidence_receipt(&EvidenceRef::Command(CommandId::new(1)))
4314 .is_some(),
4315 "the second publication must consume an archived identity, not retained payload"
4316 );
4317 let archived_reads = [StateKey::core_commands()];
4318 let archived_view = restored
4319 .simulation
4320 .plugin_view("archive-identity-probe", &archived_reads);
4321 let error = archived_view
4322 .command(CommandId::new(1))
4323 .expect_err("archived identity must not expose retained command payload");
4324 assert_eq!(error.code, ErrorCode::EvidenceContentUnavailable);
4325 restored
4326 .settle_boundary(
4327 BoundaryRequest::at(SimTime::EPOCH + SimDuration::days(1))
4328 .with_cadence(SystemCadence::Daily),
4329 )
4330 .expect("boundary two should accept the archived command identity");
4331 let holder = KnowledgeHolderRef::Person(PersonId::new(1));
4332 let records = restored
4333 .knowledge()
4334 .for_holder(&holder)
4335 .expect("the holder should have both publications");
4336 assert_eq!(records.len(), 2);
4337 assert!(records.values().all(|record| {
4338 record.origin.evidence == vec![EvidenceRef::Command(CommandId::new(1))]
4339 }));
4340 restored
4341 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH + SimDuration::days(1)))
4342 .expect("the following cut should admit boundary-two events before sealing");
4343
4344 let second_segment = restored
4345 .seal_evidence()
4346 .expect("boundary-two evidence should seal")
4347 .expect("boundary two should produce an archive segment");
4348 let second_checkpoint = restored.checkpoint().expect("checkpoint two should build");
4349 (second_checkpoint, vec![restored_prefix, second_segment])
4350 }
4351
4352 #[test]
4353 fn compact_restore_preserves_archived_identity_and_rejects_noncontiguous_segments() {
4354 let plugin = ArchivedIdentityPublicationPlugin;
4355 let (checkpoint, segments) = archived_identity_two_segment_fixture();
4356 let restored = CompactedSimulation::from_checkpoint_and_journal_with_plugins(
4357 checkpoint.clone(),
4358 segments.clone(),
4359 &[&plugin],
4360 )
4361 .expect("the complete two-segment archive should restore");
4362 assert_eq!(
4363 restored
4364 .knowledge()
4365 .for_holder(&KnowledgeHolderRef::Person(PersonId::new(1)))
4366 .expect("published holder ledger")
4367 .len(),
4368 2
4369 );
4370
4371 let first = segments[0].clone();
4372 let second = segments[1].clone();
4373 for (label, tampered) in [
4374 ("omission", vec![first.clone()]),
4375 (
4376 "overlap",
4377 vec![first.clone(), first.clone(), second.clone()],
4378 ),
4379 ("reorder", vec![second, first]),
4380 ] {
4381 let error = Simulation::from_checkpoint_and_journal(checkpoint.clone(), tampered)
4382 .err()
4383 .expect(label);
4384 assert_eq!(error.code, ErrorCode::InvalidSnapshot, "{label}");
4385 }
4386 }
4387
4388 struct PayloadContinuationPlugin;
4389
4390 fn continuation_record_ref() -> DomainRecordRef {
4391 DomainRecordRef {
4392 kind: DomainRecordKind::new("fixture.archive", "payload_continuation"),
4393 id: "primary".to_owned(),
4394 }
4395 }
4396
4397 fn continuation_payload(continuation: PayloadRequiredEvidenceContinuationV1) -> Value {
4398 Value::Object(Map::from_iter([(
4399 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD.to_owned(),
4400 serde_json::to_value(continuation).expect("fixture continuation should encode"),
4401 )]))
4402 }
4403
4404 fn mutate_payload_continuation(
4405 _view: &SimulationView<'_>,
4406 context: &BoundaryContext,
4407 ) -> Result<BoundaryProposal, CanwuError> {
4408 let directive = match context.boundary_id.get() {
4409 1 => Some(BoundaryDirective::MutateRecord {
4410 mutation: DomainRecordMutation::Create {
4411 record: DomainRecordDraft::new(
4412 continuation_record_ref(),
4413 continuation_payload(PayloadRequiredEvidenceContinuationV1::active(vec![
4414 EvidenceRef::Boundary(BoundaryId::new(1)),
4415 ])),
4416 ),
4417 },
4418 summary: "Create an active payload continuation".to_owned(),
4419 }),
4420 3 => Some(BoundaryDirective::MutateRecord {
4421 mutation: DomainRecordMutation::Update {
4422 record: DomainRecordDraft::new(
4423 continuation_record_ref(),
4424 continuation_payload(PayloadRequiredEvidenceContinuationV1::completed()),
4425 ),
4426 expected_version: 1,
4427 },
4428 summary: "Complete the payload continuation".to_owned(),
4429 }),
4430 _ => None,
4431 };
4432 Ok(BoundaryProposal {
4433 directives: directive.into_iter().collect(),
4434 ..BoundaryProposal::default()
4435 })
4436 }
4437
4438 impl SimulationPlugin for PayloadContinuationPlugin {
4439 fn name(&self) -> &str {
4440 "fixture-payload-continuation"
4441 }
4442
4443 fn version(&self) -> &str {
4444 "test-v1"
4445 }
4446
4447 fn semantic_hash(&self) -> &str {
4448 "7000000000000000000000000000000000000000000000000000000000000000"
4449 }
4450
4451 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
4452 let mut properties = BTreeMap::new();
4453 properties.insert(
4454 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD.to_owned(),
4455 payload_required_evidence_continuation_property_v1(),
4456 );
4457 let mut schema =
4458 DomainRecordSchema::new(continuation_record_ref().kind, DomainRecordClass::Record);
4459 schema.payload_schema = PayloadSchema::Object {
4460 properties,
4461 allow_additional: false,
4462 };
4463 let state = schema.state_key();
4464 registrar.register_record_schema(schema)?;
4465 let mut contract = BoundarySystemContract::new(
4466 "payload-continuation",
4467 BoundaryPhase::DomainDeltaProposal,
4468 SystemCadence::Daily,
4469 );
4470 contract.writes = vec![state];
4471 contract.visibility = StateVisibility::SameBoundary;
4472 registrar.register_boundary_system(contract, mutate_payload_continuation)
4473 }
4474 }
4475
4476 fn payload_continuation_runtime() -> CompactedSimulation {
4477 let (scenario, _) = demo_scenario();
4478 let mut simulation = Simulation::new(701, scenario).expect("fixture scenario should load");
4479 simulation
4480 .register_plugin(&PayloadContinuationPlugin)
4481 .expect("payload-continuation plugin should register");
4482 simulation
4483 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
4484 .expect("the active continuation should be created");
4485 simulation
4486 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4487 .expect("a later boundary should admit the record-change event");
4488 simulation
4489 .into_compacted()
4490 .expect("the fixture should enter compact mode")
4491 }
4492
4493 fn store_prepared(
4494 compact: &CompactedSimulation,
4495 archive: &TestArchive,
4496 ) -> PreparedEvidenceSeal {
4497 let prepared = compact
4498 .prepare_evidence_seal()
4499 .expect("the fixture should prepare a seal")
4500 .expect("the retained tail should be non-empty");
4501 assert_eq!(
4502 archive
4503 .store_evidence_segment(&prepared.segment)
4504 .expect("the prepared segment should store"),
4505 ArchiveStoreOutcome::Stored
4506 );
4507 prepared
4508 }
4509
4510 #[test]
4511 fn payload_required_receipts_are_exactly_reachable_and_prune_after_completion() {
4512 let mut compact = payload_continuation_runtime();
4513 assert_eq!(
4514 compact
4515 .seal_evidence()
4516 .expect_err("direct sealing must reject payload continuations")
4517 .code,
4518 ErrorCode::ArchiveNotReady
4519 );
4520
4521 let archive = TestArchive::default();
4522 let first = store_prepared(&compact, &archive);
4523 compact
4524 .commit_evidence_seal(&first.token, &archive)
4525 .expect("provider-backed sealing should commit");
4526 let first_checkpoint = compact.checkpoint().expect("checkpoint should build");
4527 let first_references = first_checkpoint
4528 .archived_evidence_receipts
4529 .iter()
4530 .map(|receipt| receipt.evidence.clone())
4531 .collect::<BTreeSet<_>>();
4532 let first_dependencies = first_checkpoint
4533 .evidence_dependencies
4534 .iter()
4535 .map(|dependency| (dependency.reference.clone(), dependency.requirement))
4536 .collect::<BTreeMap<_, _>>();
4537 assert_eq!(
4538 first_dependencies.get(&EvidenceRef::Boundary(BoundaryId::new(1))),
4539 Some(&EvidenceRequirement::PayloadRequired)
4540 );
4541 assert_eq!(first_references.len(), first_dependencies.len());
4542 assert_eq!(
4543 first_references,
4544 first_dependencies.keys().cloned().collect::<BTreeSet<_>>()
4545 );
4546 assert!(
4547 first
4548 .segment
4549 .archive
4550 .as_ref()
4551 .expect("prepared segment should have an archive index")
4552 .entries
4553 .len()
4554 > first_references.len(),
4555 "the full segment index must outlive the reachable-only receipt set"
4556 );
4557
4558 compact
4559 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
4560 .expect("the continuation should complete");
4561 compact
4562 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4563 .expect("a later boundary should admit the completion event");
4564 let second = store_prepared(&compact, &archive);
4565 compact
4566 .commit_evidence_seal(&second.token, &archive)
4567 .expect("the completed continuation should seal");
4568 let second_checkpoint = compact.checkpoint().expect("checkpoint should build");
4569 assert_eq!(second_checkpoint.archived_segment_headers.len(), 2);
4570 assert_eq!(second_checkpoint.archived_evidence_receipts.len(), 1);
4571 assert_eq!(second_checkpoint.evidence_dependencies.len(), 1);
4572 assert_eq!(
4573 second_checkpoint.archived_evidence_receipts[0].evidence,
4574 second_checkpoint.evidence_dependencies[0].reference
4575 );
4576 assert_eq!(
4577 second_checkpoint.evidence_dependencies[0].requirement,
4578 EvidenceRequirement::IdentityOnly
4579 );
4580 assert!(
4581 !second_checkpoint
4582 .archived_evidence_receipts
4583 .iter()
4584 .any(|receipt| receipt.evidence == EvidenceRef::Boundary(BoundaryId::new(1)))
4585 );
4586
4587 Simulation::from_checkpoint_and_journal(
4588 second_checkpoint,
4589 vec![first.segment, second.segment],
4590 )
4591 .expect("reconstruction must filter full segment indexes to the compact receipt frontier");
4592 }
4593
4594 #[test]
4595 fn payload_required_commit_fails_closed_when_an_older_segment_is_missing() {
4596 let mut compact = payload_continuation_runtime();
4597 let complete_archive = TestArchive::default();
4598 let first = store_prepared(&compact, &complete_archive);
4599 compact
4600 .commit_evidence_seal(&first.token, &complete_archive)
4601 .expect("the initial provider-backed seal should commit");
4602
4603 compact
4604 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4605 .expect("a new tail should preserve the active continuation");
4606 let incomplete_archive = TestArchive::default();
4607 let second = store_prepared(&compact, &incomplete_archive);
4608 let before = compact.checkpoint().expect("checkpoint should build");
4609 let error = compact
4610 .commit_evidence_seal(&second.token, &incomplete_archive)
4611 .expect_err("the provider must retain every payload-required segment");
4612 assert_eq!(error.code, ErrorCode::EvidenceContentUnavailable);
4613 assert_eq!(
4614 compact.checkpoint().expect("checkpoint should build"),
4615 before
4616 );
4617
4618 incomplete_archive
4619 .store_evidence_segment(&first.segment)
4620 .expect("restoring the older required segment should succeed");
4621 compact
4622 .commit_evidence_seal(&second.token, &incomplete_archive)
4623 .expect("the exact provider set should permit commit");
4624 }
4625
4626 #[test]
4627 fn host_orphan_candidates_follow_all_retained_manifests_without_deleting() {
4628 let mut compact = payload_continuation_runtime();
4629 let archive = TestArchive::default();
4630 let prepared = store_prepared(&compact, &archive);
4631 let stored_ids = archive.segment_ids();
4632 let before = compact.checkpoint().expect("checkpoint should build");
4633 assert_eq!(
4634 SimulationCheckpoint::orphaned_archive_segment_ids(
4635 std::slice::from_ref(&before),
4636 &stored_ids,
4637 )
4638 .expect("manifest reachability should validate"),
4639 vec![prepared.token.segment_id.clone()]
4640 );
4641 assert!(
4642 archive
4643 .load_evidence_segment(&prepared.token.segment_id)
4644 .expect("the store should remain readable")
4645 .is_some(),
4646 "the conformance API must not delete host content"
4647 );
4648
4649 compact
4650 .commit_evidence_seal(&prepared.token, &archive)
4651 .expect("the stored segment should commit");
4652 let after = compact.checkpoint().expect("checkpoint should build");
4653 assert!(
4654 SimulationCheckpoint::orphaned_archive_segment_ids(
4655 std::slice::from_ref(&after),
4656 &stored_ids,
4657 )
4658 .expect("committed manifest reachability should validate")
4659 .is_empty()
4660 );
4661 assert_eq!(
4662 SimulationCheckpoint::reachable_archive_segment_ids(&[before, after])
4663 .expect("all retained manifests should be scanned"),
4664 BTreeSet::from([prepared.token.segment_id.clone()])
4665 );
4666 }
4667}