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