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