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 .map(|current| ¤t.version)
2415 .filter(|version| version.version == record.version)
2416 .ok_or_else(|| {
2417 CanwuError::new(
2418 ErrorCode::ArchiveNotReady,
2419 format!(
2420 "current domain record {} has no exact version provenance",
2421 record.reference
2422 ),
2423 )
2424 })?;
2425 if !matches!(
2426 version.established_by,
2427 DomainRecordVersionSource::InitialScenario
2428 ) {
2429 promote_dependency(
2430 &mut dependencies,
2431 EvidenceRef::DomainRecordVersion(version.clone()),
2432 identity,
2433 );
2434 } else if !self.bound_initial_scenario().is_some_and(|scenario| {
2435 scenario.domain_records.iter().any(|initial| {
2436 initial.reference == record.reference && initial.version == record.version
2437 })
2438 }) {
2439 return Err(CanwuError::new(
2440 ErrorCode::ArchiveNotReady,
2441 format!(
2442 "current domain record {} has invalid initial-scenario provenance",
2443 record.reference
2444 ),
2445 ));
2446 }
2447 }
2448
2449 for reservation in &self.state.evidence.keyed_draw_reservations {
2450 promote_dependency(
2451 &mut dependencies,
2452 reservation.operation_evidence.clone(),
2453 identity,
2454 );
2455 promote_dependency(
2456 &mut dependencies,
2457 reservation.draw_receipt.evidence.clone(),
2458 identity,
2459 );
2460 }
2461 for draw in &self.state.evidence.random_draws {
2462 if matches!(draw.address, RandomDrawAddress::OperationV1(_)) {
2463 let reference = draw.operation_evidence.clone().ok_or_else(|| {
2464 CanwuError::new(
2465 ErrorCode::ArchiveNotReady,
2466 "operation-keyed retained draw is missing its operation evidence",
2467 )
2468 })?;
2469 promote_dependency(&mut dependencies, reference, identity);
2470 promote_dependency(
2471 &mut dependencies,
2472 EvidenceRef::RandomDraw(draw.id),
2473 identity,
2474 );
2475 }
2476 }
2477
2478 for archived in self.state.evidence.archived_command_requests.values() {
2479 add_command_outcome_dependencies(&mut dependencies, &archived.outcome);
2480 }
2481 for archived in self.state.evidence.archived_ingress_requests.values() {
2482 promote_dependency(
2483 &mut dependencies,
2484 EvidenceRef::Ingress(archived.receipt.ingress_id),
2485 identity,
2486 );
2487 }
2488 for archived in self.state.evidence.archived_decision_requests.values() {
2489 promote_dependency(
2490 &mut dependencies,
2491 EvidenceRef::Ingress(archived.receipt.ingress_id),
2492 identity,
2493 );
2494 }
2495 for attempt in &self.state.evidence.command_attempts {
2496 if attempt.request_id.is_none() {
2497 continue;
2498 }
2499 promote_dependency(
2500 &mut dependencies,
2501 EvidenceRef::CommandAttempt(attempt.id),
2502 identity,
2503 );
2504 if let CommandAttemptOutcome::Accepted { command_id } = attempt.outcome {
2505 promote_dependency(
2506 &mut dependencies,
2507 EvidenceRef::Command(command_id),
2508 identity,
2509 );
2510 if let Some(command) = self
2511 .state
2512 .evidence
2513 .commands
2514 .iter()
2515 .find(|command| command.id == command_id)
2516 {
2517 for event in &command.emitted_events {
2518 promote_dependency(&mut dependencies, EvidenceRef::Event(*event), identity);
2519 }
2520 }
2521 }
2522 }
2523 for ingress in &self.state.evidence.ingress {
2524 if matches!(
2525 ingress.payload,
2526 IngressPayload::Command { .. } | IngressPayload::Decision { .. }
2527 ) {
2528 promote_dependency(
2529 &mut dependencies,
2530 EvidenceRef::Ingress(ingress.id),
2531 identity,
2532 );
2533 }
2534 }
2535
2536 for action in self.state.scheduler.actions.values() {
2537 match action {
2538 ScheduledAction::ArmyArrival { order_event, .. }
2539 | ScheduledAction::PersonArrival { order_event, .. } => promote_dependency(
2540 &mut dependencies,
2541 EvidenceRef::Event(*order_event),
2542 identity,
2543 ),
2544 ScheduledAction::KnowledgeReport { dispatch_event, .. } => promote_dependency(
2545 &mut dependencies,
2546 EvidenceRef::Event(*dispatch_event),
2547 identity,
2548 ),
2549 ScheduledAction::PluginDirective { cause, .. } => {
2550 add_cause_dependency(&mut dependencies, cause);
2551 }
2552 }
2553 }
2554
2555 Ok(dependencies
2556 .into_iter()
2557 .map(|(reference, requirement)| EvidenceDependency {
2558 reference,
2559 requirement,
2560 })
2561 .collect())
2562 }
2563
2564 fn validate_payload_required_archive(
2565 &self,
2566 provider: &dyn ArchiveProvider,
2567 ) -> Result<(), CanwuError> {
2568 for dependency in self
2569 .evidence_dependencies()?
2570 .into_iter()
2571 .filter(|dependency| dependency.requirement == EvidenceRequirement::PayloadRequired)
2572 {
2573 let receipt = self
2574 .state
2575 .evidence
2576 .archived_evidence_receipts
2577 .get(&dependency.reference)
2578 .ok_or_else(|| {
2579 CanwuError::new(
2580 ErrorCode::ArchiveNotReady,
2581 "payload-required evidence has no committed archive receipt",
2582 )
2583 })?;
2584 load_verified_archived_evidence_segment(receipt, provider)?;
2585 }
2586 Ok(())
2587 }
2588
2589 fn ensure_retained_evidence_is_sealable(&self) -> Result<(), CanwuError> {
2590 if !self.state.scheduler.pending_ingress.is_empty() {
2591 return Err(CanwuError::new(
2592 ErrorCode::ArchiveNotReady,
2593 "live evidence can be sealed only when the canonical ingress queue is empty",
2594 ));
2595 }
2596 let admitted_attempts: std::collections::BTreeSet<_> = self
2597 .state
2598 .evidence
2599 .boundaries
2600 .iter()
2601 .flat_map(|record| record.admitted_attempts.iter().copied())
2602 .collect();
2603 if admitted_attempts.len() != self.state.evidence.command_attempts.len()
2604 || self
2605 .state
2606 .evidence
2607 .command_attempts
2608 .iter()
2609 .any(|attempt| !admitted_attempts.contains(&attempt.id))
2610 {
2611 return Err(CanwuError::new(
2612 ErrorCode::ArchiveNotReady,
2613 "live evidence sealing requires every retained command attempt to belong to a completed boundary",
2614 ));
2615 }
2616 let admitted_commands: std::collections::BTreeSet<_> = self
2617 .state
2618 .evidence
2619 .boundaries
2620 .iter()
2621 .flat_map(|record| record.admitted_commands.iter().copied())
2622 .collect();
2623 if admitted_commands.len() != self.state.evidence.commands.len()
2624 || self
2625 .state
2626 .evidence
2627 .commands
2628 .iter()
2629 .any(|command| !admitted_commands.contains(&command.id))
2630 {
2631 return Err(CanwuError::new(
2632 ErrorCode::ArchiveNotReady,
2633 "live evidence sealing requires every retained command to belong to a completed boundary",
2634 ));
2635 }
2636 let mut settled_ingress: std::collections::BTreeSet<_> = self
2637 .state
2638 .evidence
2639 .boundaries
2640 .iter()
2641 .flat_map(|record| record.admitted_ingress.iter().copied())
2642 .collect();
2643 for record in &self.state.evidence.ingress {
2644 if let IngressPayload::PluginCancellation { cancelled, .. } = &record.payload {
2645 settled_ingress.insert(record.id);
2646 settled_ingress.insert(*cancelled);
2647 }
2648 }
2649 if settled_ingress.len() != self.state.evidence.ingress.len()
2650 || self
2651 .state
2652 .evidence
2653 .ingress
2654 .iter()
2655 .any(|record| !settled_ingress.contains(&record.id))
2656 {
2657 return Err(CanwuError::new(
2658 ErrorCode::ArchiveNotReady,
2659 "live evidence sealing requires every retained ingress record to be admitted by a completed boundary or terminally cancelled",
2660 ));
2661 }
2662 let admitted_events: std::collections::BTreeSet<_> = self
2663 .state
2664 .evidence
2665 .boundaries
2666 .iter()
2667 .flat_map(|record| record.admitted_events.iter().copied())
2668 .collect();
2669 if self.state.counters.admitted_event_count
2670 != self
2671 .state
2672 .evidence
2673 .archived
2674 .event_count
2675 .checked_add(
2676 u64::try_from(self.state.evidence.events.len()).map_err(|_| {
2677 CanwuError::new(
2678 ErrorCode::ArchiveNotReady,
2679 "retained event count exceeds the live archive cursor range",
2680 )
2681 })?,
2682 )
2683 .ok_or_else(|| {
2684 CanwuError::new(
2685 ErrorCode::ArchiveNotReady,
2686 "retained event cursor is exhausted",
2687 )
2688 })?
2689 || admitted_events.len() != self.state.evidence.events.len()
2690 || self
2691 .state
2692 .evidence
2693 .events
2694 .iter()
2695 .any(|event| !admitted_events.contains(&event.id))
2696 {
2697 return Err(CanwuError::new(
2698 ErrorCode::ArchiveNotReady,
2699 "live evidence sealing requires every retained event to be admitted by a later completed boundary",
2700 ));
2701 }
2702 Ok(())
2703 }
2704
2705 fn seal_retained_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
2706 let start = self.state.evidence.archived;
2707 let end = self.evidence_cursor()?;
2708 if start == end {
2709 return Ok(None);
2710 }
2711 self.ensure_retained_evidence_is_sealable()?;
2712 let dependencies = self.evidence_dependencies()?;
2713 let checkpoint_hash = self.state.metadata.checkpoint_hash.clone();
2714 let commitment_roots = self.state.metadata.commitment_roots.clone();
2715 let commitment_cache = self.state.metadata.commitment_cache.clone();
2716 let prepared = (|| {
2717 self.refresh_checkpoint_hash()?;
2718
2719 let mut archived_command_requests = Vec::new();
2720 for attempt in &self.state.evidence.command_attempts {
2721 let Some(request_id) = attempt.request_id else {
2722 continue;
2723 };
2724 let outcome = self.command_outcome_from_attempt(attempt)?;
2725 archived_command_requests.push((
2726 request_id,
2727 ArchivedCommandRequestOutcome {
2728 input_hash: super::canonical_hash(
2729 "canwu.archive.command.request.v1",
2730 &(attempt.expected_revision, &attempt.envelope),
2731 )?,
2732 outcome,
2733 },
2734 ));
2735 }
2736
2737 let mut archived_ingress_requests = Vec::new();
2738 let mut archived_decision_requests = Vec::new();
2739 let mut archived_decision_command_requests = Vec::new();
2740 for record in &self.state.evidence.ingress {
2741 let receipt = IngressReceipt {
2742 ingress_id: record.id,
2743 issued_at: record.issued_at,
2744 due_at: record.due_at,
2745 };
2746 match &record.payload {
2747 IngressPayload::Command { request } => {
2748 archived_ingress_requests.push((
2749 request.request_id,
2750 ArchivedIngressRequest {
2751 input_hash: super::canonical_hash(
2752 "canwu.archive.ingress.command.v1",
2753 &(record.due_at, record.priority, request.as_ref()),
2754 )?,
2755 receipt,
2756 },
2757 ));
2758 }
2759 IngressPayload::Decision { request } => {
2760 if let Some(command) = &request.command {
2761 archived_decision_command_requests.push(command.request_id);
2762 }
2763 archived_decision_requests.push((
2764 request.request_id,
2765 ArchivedIngressRequest {
2766 input_hash: super::canonical_hash(
2767 "canwu.ingress.decision-request.v1",
2768 &(record.due_at, record.priority, request.as_ref()),
2769 )?,
2770 receipt,
2771 },
2772 ));
2773 }
2774 IngressPayload::Plugin { .. }
2775 | IngressPayload::Calendar { .. }
2776 | IngressPayload::Maintenance { .. }
2777 | IngressPayload::PluginCancellation { .. } => {}
2778 }
2779 }
2780 Ok::<_, CanwuError>((
2781 archived_command_requests,
2782 archived_ingress_requests,
2783 archived_decision_requests,
2784 archived_decision_command_requests,
2785 ))
2786 })();
2787 let (
2788 archived_command_requests,
2789 archived_ingress_requests,
2790 archived_decision_requests,
2791 archived_decision_command_requests,
2792 ) = match prepared {
2793 Ok(prepared) => prepared,
2794 Err(error) => {
2795 self.state.metadata.checkpoint_hash = checkpoint_hash;
2796 self.state.metadata.commitment_roots = commitment_roots;
2797 self.state.metadata.commitment_cache = commitment_cache;
2798 return Err(error);
2799 }
2800 };
2801
2802 let mut segment = EvidenceJournalSegment {
2803 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
2804 start,
2805 end,
2806 events: self.state.evidence.events.clone(),
2807 commands: self.state.evidence.commands.clone(),
2808 command_attempts: self.state.evidence.command_attempts.clone(),
2809 ingress: self.state.evidence.ingress.clone(),
2810 boundaries: self.state.evidence.boundaries.clone(),
2811 random_draws: self.state.evidence.random_draws.clone(),
2812 archive: None,
2813 };
2814 let (archive, receipts) = evidence_archive_index(&segment)?;
2815 segment.archive = Some(archive.clone());
2816 verify_archived_segment(&segment)?;
2817
2818 let receipt_map: BTreeMap<_, _> = receipts
2819 .iter()
2820 .cloned()
2821 .map(|receipt| (receipt.evidence.clone(), receipt))
2822 .collect();
2823 let mut reservations = Vec::new();
2824 for draw in &segment.random_draws {
2825 let RandomDrawAddress::OperationV1(address) = &draw.address else {
2826 continue;
2827 };
2828 let operation_evidence = draw.operation_evidence.clone().ok_or_else(|| {
2829 archive_error("operation-keyed archived draw is missing operation evidence")
2830 })?;
2831 let draw_receipt = receipt_map
2832 .get(&EvidenceRef::RandomDraw(draw.id))
2833 .cloned()
2834 .ok_or_else(|| {
2835 archive_error("operation-keyed archived draw is missing its receipt")
2836 })?;
2837 reservations.push(KeyedDrawReservation {
2838 stream: draw.stream.clone(),
2839 address: address.clone(),
2840 upper_exclusive: draw.upper_exclusive,
2841 purpose_hash: super::random::purpose_hash_hex_v1(&draw.purpose)?,
2842 result: draw.value,
2843 draw_id: draw.id,
2844 operation_evidence,
2845 draw_receipt,
2846 });
2847 }
2848 reservations.sort_by(|left, right| {
2849 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2850 });
2851
2852 self.state.evidence.archived_boundary_head = self
2853 .state
2854 .evidence
2855 .boundaries
2856 .last()
2857 .map(|record| record.hash.clone())
2858 .or_else(|| self.state.evidence.archived_boundary_head.clone());
2859 self.state.evidence.archived_legacy_commands |= self
2860 .state
2861 .evidence
2862 .commands
2863 .iter()
2864 .any(|record| record.attempt_id.is_none());
2865 self.state.evidence.archived_tracked_attempts |=
2866 !self.state.evidence.command_attempts.is_empty()
2867 || !self.state.evidence.ingress.is_empty();
2868 self.state.evidence.archived_unqueued_command_history |= has_unqueued_command_history(
2869 &self.state.evidence.commands,
2870 &self.state.evidence.command_attempts,
2871 &self.state.evidence.ingress,
2872 );
2873 self.state
2874 .evidence
2875 .archived_command_requests
2876 .extend(archived_command_requests);
2877 self.state
2878 .evidence
2879 .archived_ingress_requests
2880 .extend(archived_ingress_requests);
2881 self.state
2882 .evidence
2883 .archived_decision_requests
2884 .extend(archived_decision_requests);
2885 self.state
2886 .evidence
2887 .archived_decision_command_requests
2888 .extend(archived_decision_command_requests);
2889 self.state.evidence.archived = end;
2890 self.state.evidence.events.clear();
2891 self.state.evidence.commands.clear();
2892 self.state.evidence.command_attempts.clear();
2893 self.state.evidence.ingress.clear();
2894 self.state.scheduler.cancelled_ingress.clear();
2895 self.state.evidence.boundaries.clear();
2896 self.state.evidence.random_draws.clear();
2897 self.state
2898 .evidence
2899 .archived_segment_headers
2900 .push(archive.header);
2901 for receipt in receipts {
2902 if self
2903 .state
2904 .evidence
2905 .archived_evidence_receipts
2906 .insert(receipt.evidence.clone(), receipt)
2907 .is_some()
2908 {
2909 return Err(archive_error("archived evidence receipt was duplicated"));
2910 }
2911 }
2912 self.state
2913 .evidence
2914 .keyed_draw_reservations
2915 .extend(reservations);
2916 self.state
2917 .evidence
2918 .keyed_draw_reservations
2919 .sort_by(|left, right| {
2920 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
2921 });
2922 if self
2923 .state
2924 .evidence
2925 .keyed_draw_reservations
2926 .windows(2)
2927 .any(|reservations| {
2928 (&reservations[0].stream, &reservations[0].address)
2929 == (&reservations[1].stream, &reservations[1].address)
2930 })
2931 {
2932 return Err(archive_error("archived keyed reservation was duplicated"));
2933 }
2934 retain_reachable_archived_evidence_receipts(
2935 &mut self.state.evidence.archived_evidence_receipts,
2936 &dependencies,
2937 &self.state.evidence.keyed_draw_reservations,
2938 );
2939 for dependency in &dependencies {
2940 if matches!(
2941 &dependency.reference,
2942 EvidenceRef::DomainRecordVersion(version)
2943 if matches!(
2944 version.established_by,
2945 DomainRecordVersionSource::InitialScenario
2946 )
2947 ) {
2948 continue;
2949 }
2950 if !self
2951 .state
2952 .evidence
2953 .archived_evidence_receipts
2954 .contains_key(&dependency.reference)
2955 {
2956 return Err(CanwuError::new(
2957 ErrorCode::ArchiveNotReady,
2958 "declared evidence was not present in the sealed archive prefix",
2959 ));
2960 }
2961 }
2962 Ok(Some(segment))
2963 }
2964
2965 pub(super) fn checkpoint_state(&self) -> SimulationSnapshot {
2966 self.checkpoint_state_with_paged_payloads(true)
2967 }
2968
2969 fn checkpoint_state_with_paged_payloads(
2970 &self,
2971 include_paged_payloads: bool,
2972 ) -> SimulationSnapshot {
2973 SimulationSnapshot {
2974 engine_version: ENGINE_VERSION.to_owned(),
2975 snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
2976 run_manifest: Some(self.state.metadata.run_manifest.clone()),
2977 run_manifest_hash: self.state.metadata.run_manifest_hash.clone(),
2978 run_configuration: Some(self.state.metadata.run_configuration.clone()),
2979 checkpoint_hash: self.state.metadata.checkpoint_hash.clone(),
2980 commitment_format_version: self.state.metadata.commitment_format_version,
2981 commitment_roots: self.state.metadata.commitment_roots.clone(),
2982 revision_format_version: STATE_REVISION_FORMAT_VERSION,
2983 state_revision: self.state.counters.state_revision,
2984 replay_revision_format_version: self.state.metadata.replay_revision_format_version,
2985 admission_cursor_format_version: ADMISSION_CURSOR_FORMAT_VERSION,
2986 admitted_attempt_count: self.state.counters.admitted_attempt_count,
2987 admitted_command_count: self.state.counters.admitted_command_count,
2988 admitted_event_count: self.state.counters.admitted_event_count,
2989 initial_time: self.state.scheduler.initial_time,
2990 initial_scenario: self.bound_initial_scenario().cloned(),
2991 now: self.state.scheduler.now,
2992 plugin_registration_closed: self.state.metadata.plugin_registration_closed,
2993 entities: self.state.current.entities.iter().cloned().collect(),
2994 world: self.world(),
2995 person_availability: self.state.current.person_availability.clone(),
2996 created_persons: self.state.current.created_persons.clone(),
2997 knowledge: self.state.current.knowledge.clone(),
2998 events: Vec::new(),
2999 commands: Vec::new(),
3000 command_attempts: Vec::new(),
3001 ingress: Vec::new(),
3002 boundaries: Vec::new(),
3003 plugin_components: self
3004 .state
3005 .current
3006 .plugin_components
3007 .values()
3008 .cloned()
3009 .collect(),
3010 domain_records: if include_paged_payloads {
3011 self.state
3012 .current
3013 .domain_records
3014 .values()
3015 .cloned()
3016 .collect()
3017 } else {
3018 Vec::new()
3019 },
3020 decisions: if include_paged_payloads {
3021 self.state.current.decisions.clone()
3022 } else {
3023 DecisionState::default()
3024 },
3025 plugin_descriptors: self.plugins.descriptors().cloned().collect(),
3026 schema: self.schema.clone(),
3027 root_seed: self.state.current.root_seed,
3028 authority_root_seed: self.state.current.authority_root_seed,
3029 random_streams: self
3030 .state
3031 .current
3032 .random_streams
3033 .values()
3034 .cloned()
3035 .collect(),
3036 random_draws: Vec::new(),
3037 scheduled: self
3038 .state
3039 .scheduler
3040 .actions
3041 .iter()
3042 .map(|(key, action)| ScheduledRecord {
3043 key: key.clone(),
3044 action: action.clone(),
3045 })
3046 .collect(),
3047 pending_transition_manifests: self
3048 .state
3049 .scheduler
3050 .transition_manifests
3051 .values()
3052 .cloned()
3053 .collect(),
3054 legacy_rng: None,
3055 next_event_id: self.state.counters.next_event_id,
3056 next_command_id: self.state.counters.next_command_id,
3057 next_command_attempt_id: self.state.counters.next_command_attempt_id,
3058 next_ingress_id: self.state.counters.next_ingress_id,
3059 next_boundary_id: self.state.counters.next_boundary_id,
3060 next_random_draw_id: self.state.counters.next_random_draw_id,
3061 next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
3062 next_schedule_sequence: self.state.counters.next_schedule_sequence,
3063 next_correlation_id: self.state.counters.next_correlation_id,
3064 next_decision_trace_id: self.state.counters.next_decision_trace_id,
3065 next_person_id: self.state.counters.next_person_id,
3066 }
3067 }
3068
3069 fn build_paged_checkpoint_pages(
3070 &self,
3071 provider: Option<&dyn StatePageProvider>,
3072 ) -> Result<(PagedSimulationCheckpoint, Vec<StatePageBlob>), CanwuError> {
3073 let (domain_records, domain_pages) = match provider {
3074 Some(provider) => self
3075 .state
3076 .current
3077 .domain_records
3078 .missing_state_pages(provider)?,
3079 None => self.state.current.domain_records.state_pages()?,
3080 };
3081 let decisions = &self.state.current.decisions;
3082 let hot_decisions = decisions.paged_checkpoint_hot_state();
3083 let hot_decision_page =
3084 StatePageBlob::new(serde_json::to_vec(&hot_decisions).map_err(|error| {
3085 invalid_snapshot_error(format!("cannot encode paged hot decision state: {error}"))
3086 })?)?;
3087 let decision_bucket_pages = decisions
3088 .decision_archive_bucket_page_ids()
3089 .iter()
3090 .map(|(bucket, page_id)| (*bucket, page_id.clone()))
3091 .collect::<BTreeMap<_, _>>();
3092 let mut decision_pages = Vec::with_capacity(
3093 decision_bucket_pages
3094 .len()
3095 .saturating_add(
3096 decision_bucket_pages
3097 .len()
3098 .div_ceil(MAX_PAGED_DECISION_DIRECTORY_ENTRIES),
3099 )
3100 .saturating_add(2),
3101 );
3102 for (bucket, expected_page_id) in &decision_bucket_pages {
3103 let Some(bucket_page) = decisions
3104 .decision_archive_bucket_page(*bucket)
3105 .map_err(|error| invalid_snapshot_error(error.to_string()))?
3106 else {
3107 if provider.is_none() {
3108 return Err(invalid_snapshot_error(
3109 "portable decision checkpoint requires every archive bucket to be resident",
3110 ));
3111 }
3112 continue;
3117 };
3118 let page = StatePageBlob::new(serde_json::to_vec(&bucket_page).map_err(|error| {
3119 invalid_snapshot_error(format!(
3120 "cannot encode paged decision archive bucket: {error}"
3121 ))
3122 })?)?;
3123 if page.page_id != *expected_page_id {
3124 return Err(invalid_snapshot_error(
3125 "decision archive bucket page disagrees with its cached commitment",
3126 ));
3127 }
3128 decision_pages.push(page);
3129 }
3130 let decision_bucket_directory = decision_bucket_pages.into_iter().collect::<Vec<_>>();
3131 let mut archive_directory_page_ids = Vec::with_capacity(
3132 decision_bucket_directory
3133 .len()
3134 .div_ceil(MAX_PAGED_DECISION_DIRECTORY_ENTRIES),
3135 );
3136 for chunk in decision_bucket_directory.chunks(MAX_PAGED_DECISION_DIRECTORY_ENTRIES) {
3137 let directory_page = PagedDecisionDirectoryPage {
3138 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3139 archive_bucket_pages: chunk.to_vec(),
3140 };
3141 directory_page.validate()?;
3142 let page =
3143 StatePageBlob::new(serde_json::to_vec(&directory_page).map_err(|error| {
3144 invalid_snapshot_error(format!(
3145 "cannot encode paged decision directory page: {error}"
3146 ))
3147 })?)?;
3148 archive_directory_page_ids.push(page.page_id.clone());
3149 decision_pages.push(page);
3150 }
3151 let archive_receipt_count = u64::try_from(decisions.archived_history_count())
3152 .map_err(|_| invalid_snapshot_error("decision archive receipt count exceeds u64"))?;
3153 let decision_manifest_page = StatePageBlob::new(
3154 serde_json::to_vec(&PagedDecisionManifest {
3155 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3156 hot_page_id: hot_decision_page.page_id.clone(),
3157 archive_receipt_root: self
3158 .state
3159 .current
3160 .decisions
3161 .archive_receipt_commitment()
3162 .map_err(|error| invalid_snapshot_error(error.to_string()))?,
3163 archive_receipt_count,
3164 archive_directory_page_ids,
3165 })
3166 .map_err(|error| {
3167 invalid_snapshot_error(format!("cannot encode paged decision manifest: {error}"))
3168 })?,
3169 )?;
3170 let checkpoint_without_paged_state = self.checkpoint_with_paged_payloads(false)?;
3171 let envelope = PagedCheckpointEnvelope {
3172 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3173 checkpoint_without_paged_state,
3174 domain_records,
3175 decision_manifest_page_id: decision_manifest_page.page_id.clone(),
3176 };
3177 let envelope_page =
3178 StatePageBlob::new(serde_json::to_vec(&envelope).map_err(|error| {
3179 invalid_snapshot_error(format!("cannot encode paged checkpoint envelope: {error}"))
3180 })?)?;
3181 let checkpoint = PagedSimulationCheckpoint {
3182 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
3183 root_page_id: envelope_page.page_id.clone(),
3184 checkpoint_hash: envelope
3185 .checkpoint_without_paged_state
3186 .state
3187 .checkpoint_hash
3188 .clone(),
3189 };
3190 let mut pages = BTreeMap::new();
3191 for page in domain_pages.into_iter().chain(decision_pages).chain([
3192 hot_decision_page,
3193 decision_manifest_page,
3194 envelope_page,
3195 ]) {
3196 if let Some(existing) = pages.insert(page.page_id.clone(), page.clone())
3197 && existing != page
3198 {
3199 return Err(invalid_snapshot_error(
3200 "one state-page ID resolved to conflicting canonical bytes",
3201 ));
3202 }
3203 }
3204 Ok((checkpoint, pages.into_values().collect()))
3205 }
3206
3207 pub fn prepare_paged_checkpoint(
3211 &self,
3212 source: Option<&PagedSimulationCheckpoint>,
3213 provider: &dyn StatePageProvider,
3214 ) -> Result<PreparedPagedSimulationCheckpoint, CanwuError> {
3215 if let Some(source) = source {
3216 if source.format_version != PAGED_CHECKPOINT_FORMAT_VERSION {
3217 return Err(invalid_snapshot_error(
3218 "paged checkpoint source uses an unsupported format",
3219 ));
3220 }
3221 let page = provider
3222 .load_state_page(&source.root_page_id)?
3223 .ok_or_else(|| {
3224 CanwuError::new(
3225 ErrorCode::StatePageUnavailable,
3226 "paged checkpoint source root is unavailable",
3227 )
3228 })?;
3229 page.validate()?;
3230 if page.page_id != source.root_page_id {
3231 return Err(invalid_snapshot_error(
3232 "paged checkpoint source provider returned the wrong root",
3233 ));
3234 }
3235 }
3236 let (checkpoint, pages) = self.build_paged_checkpoint_pages(Some(provider))?;
3237 let mut new_pages = Vec::new();
3238 for page in pages {
3239 match provider.load_state_page(&page.page_id)? {
3240 Some(existing) => {
3241 existing.validate()?;
3242 if existing != page {
3243 return Err(invalid_snapshot_error(
3244 "state-page provider contains conflicting canonical bytes",
3245 ));
3246 }
3247 }
3248 None => new_pages.push(page),
3249 }
3250 }
3251 let source_root = source.map_or_else(
3252 || canonical_byte_hash("canwu.paged-checkpoint.genesis.v1", &[]),
3253 |source| source.root_page_id.clone(),
3254 );
3255 let delta = prepare_state_delta(&source_root, &checkpoint.root_page_id, new_pages)?;
3256 Ok(PreparedPagedSimulationCheckpoint { checkpoint, delta })
3257 }
3258
3259 pub fn portable_paged_checkpoint(
3262 &self,
3263 ) -> Result<PortablePagedSimulationCheckpoint, CanwuError> {
3264 let (checkpoint, pages) = self.build_paged_checkpoint_pages(None)?;
3265 Ok(PortablePagedSimulationCheckpoint { checkpoint, pages })
3266 }
3267
3268 pub fn from_paged_checkpoint(
3272 checkpoint: &PagedSimulationCheckpoint,
3273 provider: &dyn StatePageProvider,
3274 ) -> Result<Self, CanwuError> {
3275 Self::from_paged_checkpoint_and_journal(checkpoint, provider, Vec::new())
3276 }
3277
3278 pub fn from_paged_checkpoint_and_journal(
3281 checkpoint: &PagedSimulationCheckpoint,
3282 provider: &dyn StatePageProvider,
3283 segments: Vec<EvidenceJournalSegment>,
3284 ) -> Result<Self, CanwuError> {
3285 if checkpoint.format_version != PAGED_CHECKPOINT_FORMAT_VERSION {
3286 return Err(invalid_snapshot_error(
3287 "paged checkpoint uses an unsupported format",
3288 ));
3289 }
3290 let root_page = provider
3291 .load_state_page(&checkpoint.root_page_id)?
3292 .ok_or_else(|| {
3293 CanwuError::new(
3294 ErrorCode::StatePageUnavailable,
3295 "paged checkpoint root is unavailable",
3296 )
3297 })?;
3298 root_page.validate()?;
3299 if root_page.page_id != checkpoint.root_page_id {
3300 return Err(invalid_snapshot_error(
3301 "paged checkpoint provider returned the wrong root page",
3302 ));
3303 }
3304 let envelope: PagedCheckpointEnvelope =
3305 serde_json::from_slice(&root_page.bytes).map_err(|error| {
3306 invalid_snapshot_error(format!("invalid paged checkpoint envelope: {error}"))
3307 })?;
3308 if envelope.format_version != PAGED_CHECKPOINT_FORMAT_VERSION
3309 || envelope
3310 .checkpoint_without_paged_state
3311 .state
3312 .checkpoint_hash
3313 != checkpoint.checkpoint_hash
3314 || !envelope
3315 .checkpoint_without_paged_state
3316 .state
3317 .domain_records
3318 .is_empty()
3319 || !envelope
3320 .checkpoint_without_paged_state
3321 .state
3322 .decisions
3323 .is_empty()
3324 {
3325 return Err(invalid_snapshot_error(
3326 "paged checkpoint envelope metadata is inconsistent",
3327 ));
3328 }
3329 let decision_manifest_page = provider
3330 .load_state_page(&envelope.decision_manifest_page_id)?
3331 .ok_or_else(|| {
3332 CanwuError::new(
3333 ErrorCode::StatePageUnavailable,
3334 "paged decision manifest is unavailable",
3335 )
3336 })?;
3337 decision_manifest_page.validate()?;
3338 if decision_manifest_page.page_id != envelope.decision_manifest_page_id {
3339 return Err(invalid_snapshot_error(
3340 "paged decision provider returned the wrong manifest page",
3341 ));
3342 }
3343 let decision_manifest: PagedDecisionManifest =
3344 serde_json::from_slice(&decision_manifest_page.bytes).map_err(|error| {
3345 invalid_snapshot_error(format!("invalid paged decision manifest: {error}"))
3346 })?;
3347 validate_paged_decision_manifest(&decision_manifest)?;
3348 let hot_decision_page = provider
3349 .load_state_page(&decision_manifest.hot_page_id)?
3350 .ok_or_else(|| {
3351 CanwuError::new(
3352 ErrorCode::StatePageUnavailable,
3353 "paged hot decision state is unavailable",
3354 )
3355 })?;
3356 hot_decision_page.validate()?;
3357 if hot_decision_page.page_id != decision_manifest.hot_page_id {
3358 return Err(invalid_snapshot_error(
3359 "paged decision provider returned the wrong hot-state page",
3360 ));
3361 }
3362 let hot_decisions: DecisionState = serde_json::from_slice(&hot_decision_page.bytes)
3363 .map_err(|error| {
3364 invalid_snapshot_error(format!("invalid paged hot decision state: {error}"))
3365 })?;
3366 let mut directory_pages =
3367 Vec::with_capacity(decision_manifest.archive_directory_page_ids.len());
3368 for directory_page_id in &decision_manifest.archive_directory_page_ids {
3369 let directory_state_page =
3370 provider
3371 .load_state_page(directory_page_id)?
3372 .ok_or_else(|| {
3373 CanwuError::new(
3374 ErrorCode::StatePageUnavailable,
3375 "paged decision directory page is unavailable",
3376 )
3377 })?;
3378 directory_state_page.validate()?;
3379 if directory_state_page.page_id != *directory_page_id {
3380 return Err(invalid_snapshot_error(
3381 "paged decision provider returned the wrong directory page",
3382 ));
3383 }
3384 let directory_page: PagedDecisionDirectoryPage =
3385 serde_json::from_slice(&directory_state_page.bytes).map_err(|error| {
3386 invalid_snapshot_error(format!(
3387 "invalid paged decision directory page: {error}"
3388 ))
3389 })?;
3390 directory_pages.push(directory_page);
3391 }
3392 let archive_bucket_pages = assemble_paged_decision_directory(directory_pages)?;
3393 let required_dependency_pages = hot_decisions
3394 .required_archived_dependency_page_keys()
3395 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3396 let mut resident_dependency_pages = Vec::with_capacity(required_dependency_pages.len());
3397 for page_key in required_dependency_pages {
3398 let page_id = archive_bucket_pages.get(&page_key).ok_or_else(|| {
3399 invalid_snapshot_error(
3400 "hot decision state references history absent from the archive directory",
3401 )
3402 })?;
3403 let state_page = provider.load_state_page(page_id)?.ok_or_else(|| {
3404 CanwuError::new(
3405 ErrorCode::StatePageUnavailable,
3406 "decision dependency locator page is unavailable",
3407 )
3408 })?;
3409 state_page.validate()?;
3410 if state_page.page_id != *page_id {
3411 return Err(invalid_snapshot_error(
3412 "paged decision provider returned the wrong dependency locator page",
3413 ));
3414 }
3415 let page: DecisionArchiveBucketPage = serde_json::from_slice(&state_page.bytes)
3416 .map_err(|error| {
3417 invalid_snapshot_error(format!(
3418 "invalid decision dependency locator page: {error}"
3419 ))
3420 })?;
3421 page.validate()
3422 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3423 if page.bucket != page_key.bucket
3424 || page.segment != page_key.segment
3425 || page
3426 .state_page_id()
3427 .map_err(|error| invalid_snapshot_error(error.to_string()))?
3428 != *page_id
3429 {
3430 return Err(invalid_snapshot_error(
3431 "decision dependency locator page disagrees with the committed directory",
3432 ));
3433 }
3434 resident_dependency_pages.push(page);
3435 }
3436 let decisions = DecisionState::from_paged_checkpoint_root_with_resident_pages(
3437 hot_decisions,
3438 archive_bucket_pages.into(),
3439 decision_manifest.archive_receipt_count,
3440 &decision_manifest.archive_receipt_root,
3441 resident_dependency_pages,
3442 )
3443 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3444 let domain_records = super::PersistentDomainRecordStore::from_state_pages(
3445 &envelope.domain_records,
3446 provider,
3447 )?;
3448 let mut compact_checkpoint = envelope.checkpoint_without_paged_state;
3449 compact_checkpoint.state.domain_records = domain_records.values().cloned().collect();
3450 compact_checkpoint.state.decisions = decisions;
3451 let mut simulation = Self::from_checkpoint_and_journal(compact_checkpoint, segments)?;
3452 simulation.state.current.domain_records = domain_records;
3453 Ok(simulation)
3454 }
3455
3456 pub fn from_portable_paged_checkpoint(
3457 portable: PortablePagedSimulationCheckpoint,
3458 ) -> Result<Self, CanwuError> {
3459 let mut provider = EmbeddedPageProvider::default();
3460 for page in portable.pages {
3461 page.validate()?;
3462 if let Some(existing) = provider.pages.insert(page.page_id.clone(), page.clone())
3463 && existing != page
3464 {
3465 return Err(invalid_snapshot_error(
3466 "portable paged checkpoint contains conflicting duplicate pages",
3467 ));
3468 }
3469 }
3470 Self::from_paged_checkpoint(&portable.checkpoint, &provider)
3471 }
3472
3473 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
3475 EvidenceCursor::from_evidence(&self.state.evidence)
3476 }
3477
3478 fn checkpoint_with_paged_payloads(
3479 &self,
3480 include_paged_payloads: bool,
3481 ) -> Result<SimulationCheckpoint, CanwuError> {
3482 let archived_segment_headers = self.state.evidence.archived_segment_headers.clone();
3483 let evidence_dependencies = self.evidence_dependencies()?;
3484 let mut keyed_draw_reservations = self.state.evidence.keyed_draw_reservations.clone();
3485 keyed_draw_reservations.sort_by(|left, right| {
3486 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
3487 });
3488 let required_receipts =
3489 required_archived_receipt_references(&evidence_dependencies, &keyed_draw_reservations);
3490 let archived_evidence_receipts: Vec<_> = self
3491 .state
3492 .evidence
3493 .archived_evidence_receipts
3494 .iter()
3495 .filter(|(reference, _)| required_receipts.contains(*reference))
3496 .map(|(_, receipt)| receipt.clone())
3497 .collect();
3498 let checkpoint = SimulationCheckpoint {
3499 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
3500 journal_end: self.evidence_cursor()?,
3501 state: self.checkpoint_state_with_paged_payloads(include_paged_payloads),
3502 archived_segment_manifest_root: skipped_commitment_root(
3503 ARCHIVED_SEGMENT_MANIFEST_DOMAIN,
3504 &archived_segment_headers,
3505 )?,
3506 archived_segment_headers,
3507 archived_receipt_root: skipped_commitment_root(
3508 ARCHIVED_RECEIPT_DOMAIN,
3509 &archived_evidence_receipts,
3510 )?,
3511 archived_evidence_receipts,
3512 evidence_dependency_root: skipped_commitment_root(
3513 EVIDENCE_DEPENDENCY_DOMAIN,
3514 &evidence_dependencies,
3515 )?,
3516 evidence_dependencies,
3517 keyed_reservation_root: skipped_commitment_root(
3518 KEYED_RESERVATION_DOMAIN,
3519 &keyed_draw_reservations,
3520 )?,
3521 keyed_draw_reservations,
3522 };
3523 if include_paged_payloads {
3524 validate_compact_continuation(&checkpoint)?;
3525 }
3526 Ok(checkpoint)
3527 }
3528
3529 pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
3531 self.checkpoint_with_paged_payloads(true)
3532 }
3533
3534 pub fn archive_reachability_manifest(
3539 &self,
3540 retained_checkpoints: &[SimulationCheckpoint],
3541 page_retention: &StatePageRetentionLedger,
3542 decision_provider: &dyn DecisionArchiveProvider,
3543 plugin_provider: &dyn super::PluginArchiveObjectProvider,
3544 ) -> Result<ArchiveReachabilityManifest, CanwuError> {
3545 let mut manifest = ArchiveReachabilityManifest {
3546 state_page_ids: page_retention.reachable_page_ids(),
3547 evidence_segment_ids: SimulationCheckpoint::reachable_archive_segment_ids(
3548 retained_checkpoints,
3549 )?,
3550 ..ArchiveReachabilityManifest::default()
3551 };
3552 manifest.evidence_segment_ids.extend(
3553 self.state
3554 .evidence
3555 .archived_segment_headers
3556 .iter()
3557 .map(|header| header.segment_id.clone()),
3558 );
3559 let decisions = self
3560 .state
3561 .current
3562 .decisions
3563 .archive_reachability(decision_provider)
3564 .map_err(super::decision::decision_error)?;
3565 manifest.state_page_ids.extend(decisions.bucket_page_ids);
3566 manifest.decision_blob_ids.extend(decisions.blob_locators);
3567 let pending_ingress_ids = self
3568 .state
3569 .scheduler
3570 .pending_ingress
3571 .iter()
3572 .map(|key| key.id)
3573 .collect::<BTreeSet<_>>();
3574 for record in &self.state.evidence.ingress {
3575 if let IngressPayload::Maintenance { request } = &record.payload
3576 && let super::MaintenanceIngressRequest::DecisionArchive { commit } =
3577 request.as_ref()
3578 {
3579 manifest
3580 .decision_blob_ids
3581 .extend(commit.archive_locators().map(str::to_owned));
3582 }
3583 if pending_ingress_ids.contains(&record.id)
3584 && let IngressPayload::Plugin {
3585 archive_retention, ..
3586 } = &record.payload
3587 {
3588 for retention in archive_retention {
3589 manifest.insert_plugin_object(
3590 retention.namespace.clone(),
3591 retention.object_id.clone(),
3592 );
3593 }
3594 }
3595 }
3596 for (plugin, participant) in &self.plugins.archive_reachability_participants {
3597 let reads = self
3598 .plugins
3599 .state_owners
3600 .iter()
3601 .filter_map(|(state, owner)| (owner == plugin).then_some(state.clone()))
3602 .collect::<Vec<_>>();
3603 let reader = format!("{plugin}.archive_reachability");
3604 let view = self.plugin_view(&reader, &reads);
3605 participant(&view, plugin_provider, &mut manifest)?;
3606 }
3607 Ok(manifest)
3608 }
3609
3610 pub fn journal_segment_since(
3612 &self,
3613 start: EvidenceCursor,
3614 ) -> Result<EvidenceJournalSegment, CanwuError> {
3615 let end = self.evidence_cursor()?;
3616 let cut = |value: u64, archived: u64, len: usize, label: &str| {
3617 let value = value.checked_sub(archived).ok_or_else(|| {
3618 CanwuError::new(
3619 ErrorCode::InvalidSnapshot,
3620 format!("{label} journal cursor precedes the retained live evidence window"),
3621 )
3622 })?;
3623 let value = usize::try_from(value).map_err(|_| {
3624 CanwuError::new(
3625 ErrorCode::InvalidSnapshot,
3626 format!("{label} journal cursor is not representable on this platform"),
3627 )
3628 })?;
3629 if value > len {
3630 return Err(CanwuError::new(
3631 ErrorCode::InvalidSnapshot,
3632 format!("{label} journal cursor exceeds the current evidence tail"),
3633 ));
3634 }
3635 Ok(value)
3636 };
3637 let archived = self.state.evidence.archived;
3638 let event_start = cut(
3639 start.event_count,
3640 archived.event_count,
3641 self.state.evidence.events.len(),
3642 "event",
3643 )?;
3644 let command_start = cut(
3645 start.command_count,
3646 archived.command_count,
3647 self.state.evidence.commands.len(),
3648 "command",
3649 )?;
3650 let attempt_start = cut(
3651 start.command_attempt_count,
3652 archived.command_attempt_count,
3653 self.state.evidence.command_attempts.len(),
3654 "command-attempt",
3655 )?;
3656 let ingress_start = cut(
3657 start.ingress_count,
3658 archived.ingress_count,
3659 self.state.evidence.ingress.len(),
3660 "ingress",
3661 )?;
3662 let boundary_start = cut(
3663 start.boundary_count,
3664 archived.boundary_count,
3665 self.state.evidence.boundaries.len(),
3666 "boundary",
3667 )?;
3668 let draw_start = cut(
3669 start.random_draw_count,
3670 archived.random_draw_count,
3671 self.state.evidence.random_draws.len(),
3672 "random-draw",
3673 )?;
3674 Ok(EvidenceJournalSegment {
3675 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
3676 start,
3677 end,
3678 events: self.state.evidence.events[event_start..].to_vec(),
3679 commands: self.state.evidence.commands[command_start..].to_vec(),
3680 command_attempts: self.state.evidence.command_attempts[attempt_start..].to_vec(),
3681 ingress: self.state.evidence.ingress[ingress_start..].to_vec(),
3682 boundaries: self.state.evidence.boundaries[boundary_start..].to_vec(),
3683 random_draws: self.state.evidence.random_draws[draw_start..].to_vec(),
3684 archive: None,
3685 })
3686 }
3687
3688 pub fn checkpoint_journal(&self) -> Result<CheckpointJournal, CanwuError> {
3690 if self.state.evidence.archived != EvidenceCursor::default() {
3691 return Err(CanwuError::new(
3692 ErrorCode::InvalidSnapshot,
3693 "a compact live runtime requires its previously sealed evidence segments to build a portable save",
3694 ));
3695 }
3696 let segment = self.journal_segment_since(EvidenceCursor::default())?;
3697 Ok(CheckpointJournal {
3698 checkpoint: self.checkpoint()?,
3699 segments: (segment.start != segment.end)
3700 .then_some(segment)
3701 .into_iter()
3702 .collect(),
3703 })
3704 }
3705
3706 pub fn checkpoint_journal_json(&self) -> Result<String, CanwuError> {
3708 serde_json::to_string_pretty(&self.checkpoint_journal()?).map_err(|error| {
3709 CanwuError::new(
3710 ErrorCode::InvalidSnapshot,
3711 format!("could not serialize checkpoint journal: {error}"),
3712 )
3713 })
3714 }
3715
3716 fn snapshot_from_checkpoint_and_journal(
3717 checkpoint: SimulationCheckpoint,
3718 segments: Vec<EvidenceJournalSegment>,
3719 ) -> Result<SimulationSnapshot, CanwuError> {
3720 if checkpoint.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
3721 return Err(invalid_snapshot_error(format!(
3722 "checkpoint-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
3723 checkpoint.format_version
3724 )));
3725 }
3726 validate_compact_continuation(&checkpoint)?;
3727 let expected_headers = checkpoint.archived_segment_headers;
3728 let expected_receipts = checkpoint.archived_evidence_receipts;
3729 let expected_dependencies = checkpoint.evidence_dependencies;
3730 let expected_reservations = checkpoint.keyed_draw_reservations;
3731 let mut snapshot = checkpoint.state;
3732 if snapshot.snapshot_format_version != SNAPSHOT_FORMAT_VERSION {
3733 return Err(invalid_snapshot_error(format!(
3734 "checkpoint-journal format {CHECKPOINT_JOURNAL_FORMAT_VERSION} requires snapshot format {SNAPSHOT_FORMAT_VERSION}"
3735 )));
3736 }
3737 if !snapshot.events.is_empty()
3738 || !snapshot.commands.is_empty()
3739 || !snapshot.command_attempts.is_empty()
3740 || !snapshot.ingress.is_empty()
3741 || !snapshot.boundaries.is_empty()
3742 || !snapshot.random_draws.is_empty()
3743 {
3744 return Err(invalid_snapshot_error(
3745 "checkpoint current state must not duplicate append-only evidence",
3746 ));
3747 }
3748
3749 let mut cursor = EvidenceCursor::default();
3750 let mut rebuilt_headers = Vec::new();
3751 let mut rebuilt_receipts = BTreeMap::new();
3752 let mut rebuilt_reservations = Vec::new();
3753 for segment in segments {
3754 if segment.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
3755 return Err(invalid_snapshot_error(format!(
3756 "evidence-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
3757 segment.format_version
3758 )));
3759 }
3760 if segment.start != cursor {
3761 return Err(invalid_snapshot_error(
3762 "evidence-journal segments must form one contiguous global prefix",
3763 ));
3764 }
3765 let end = cursor.checked_advance(&segment)?;
3766 if end == cursor {
3767 return Err(invalid_snapshot_error(
3768 "evidence-journal segments must advance at least one journal cursor",
3769 ));
3770 }
3771 if segment.end != end {
3772 return Err(invalid_snapshot_error(
3773 "evidence-journal segment end does not match its encoded records",
3774 ));
3775 }
3776 if let Some(archive) = &segment.archive {
3777 let receipts = verify_archived_segment(&segment).map_err(|error| {
3778 invalid_snapshot_error(format!(
3779 "archived evidence segment is invalid: {}",
3780 error.message
3781 ))
3782 })?;
3783 rebuilt_headers.push(archive.header.clone());
3784 for receipt in receipts {
3785 if rebuilt_receipts
3786 .insert(receipt.evidence.clone(), receipt)
3787 .is_some()
3788 {
3789 return Err(invalid_snapshot_error(
3790 "archived evidence segments contain duplicate receipts",
3791 ));
3792 }
3793 }
3794 for draw in &segment.random_draws {
3795 let RandomDrawAddress::OperationV1(address) = &draw.address else {
3796 continue;
3797 };
3798 let operation_evidence = draw.operation_evidence.clone().ok_or_else(|| {
3799 invalid_snapshot_error("archived keyed draw is missing operation evidence")
3800 })?;
3801 let draw_receipt = rebuilt_receipts
3802 .get(&EvidenceRef::RandomDraw(draw.id))
3803 .cloned()
3804 .ok_or_else(|| {
3805 invalid_snapshot_error("archived keyed draw receipt is missing")
3806 })?;
3807 rebuilt_reservations.push(KeyedDrawReservation {
3808 stream: draw.stream.clone(),
3809 address: address.clone(),
3810 upper_exclusive: draw.upper_exclusive,
3811 purpose_hash: super::random::purpose_hash_hex_v1(&draw.purpose)?,
3812 result: draw.value,
3813 draw_id: draw.id,
3814 operation_evidence,
3815 draw_receipt,
3816 });
3817 }
3818 }
3819 snapshot.events.extend(segment.events);
3820 snapshot.commands.extend(segment.commands);
3821 snapshot.command_attempts.extend(segment.command_attempts);
3822 snapshot.ingress.extend(segment.ingress);
3823 snapshot.boundaries.extend(segment.boundaries);
3824 snapshot.random_draws.extend(segment.random_draws);
3825 cursor = end;
3826 }
3827 if cursor != checkpoint.journal_end {
3828 return Err(invalid_snapshot_error(
3829 "evidence-journal segments do not reach the checkpoint journal cut",
3830 ));
3831 }
3832 rebuilt_reservations.sort_by(|left, right| {
3833 (&left.stream, &left.address).cmp(&(&right.stream, &right.address))
3834 });
3835 let rebuilt_dependencies =
3836 Simulation::from_snapshot(snapshot.clone())?.evidence_dependencies()?;
3837 retain_reachable_archived_evidence_receipts(
3838 &mut rebuilt_receipts,
3839 &rebuilt_dependencies,
3840 &rebuilt_reservations,
3841 );
3842 let rebuilt_receipts: Vec<_> = rebuilt_receipts.into_values().collect();
3843 if rebuilt_headers != expected_headers
3844 || rebuilt_receipts != expected_receipts
3845 || rebuilt_dependencies != expected_dependencies
3846 || rebuilt_reservations != expected_reservations
3847 {
3848 return Err(invalid_snapshot_error(
3849 "checkpoint compact continuation does not match its reconstructed authoritative indexes",
3850 ));
3851 }
3852 Ok(snapshot)
3853 }
3854
3855 pub fn from_checkpoint_and_journal(
3857 checkpoint: SimulationCheckpoint,
3858 segments: Vec<EvidenceJournalSegment>,
3859 ) -> Result<Self, CanwuError> {
3860 Self::from_snapshot(Self::snapshot_from_checkpoint_and_journal(
3861 checkpoint, segments,
3862 )?)
3863 }
3864
3865 pub fn from_checkpoint_journal(bundle: CheckpointJournal) -> Result<Self, CanwuError> {
3867 Self::from_checkpoint_and_journal(bundle.checkpoint, bundle.segments)
3868 }
3869
3870 pub fn from_checkpoint_journal_with_plugins(
3872 bundle: CheckpointJournal,
3873 plugins: &[&dyn SimulationPlugin],
3874 ) -> Result<Self, CanwuError> {
3875 let mut simulation = Self::from_checkpoint_journal(bundle)?;
3876 for plugin in plugins {
3877 simulation.register_plugin(*plugin)?;
3878 }
3879 simulation.ensure_runtime_ready()?;
3880 Ok(simulation)
3881 }
3882
3883 pub fn from_checkpoint_journal_json(json: &str) -> Result<Self, CanwuError> {
3885 let bundle: CheckpointJournal =
3886 super::deserialize_current_json(json, "checkpoint journal")?;
3887 Self::from_checkpoint_journal(bundle)
3888 }
3889
3890 pub fn from_checkpoint_journal_json_with_plugins(
3892 json: &str,
3893 plugins: &[&dyn SimulationPlugin],
3894 ) -> Result<Self, CanwuError> {
3895 let bundle: CheckpointJournal =
3896 super::deserialize_current_json(json, "checkpoint journal")?;
3897 Self::from_checkpoint_journal_with_plugins(bundle, plugins)
3898 }
3899}
3900
3901pub fn format8_paged_checkpoint_scale_probe(
3905 decision_count: usize,
3906) -> Result<PagedCheckpointScaleMetrics, CanwuError> {
3907 let fixture = canwu_decision::format8_decision_locator_scale_fixture(decision_count)
3908 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
3909 let decision_metrics = fixture.metrics.clone();
3910 let decision_blobs = fixture
3911 .archive_blobs
3912 .into_iter()
3913 .map(|blob| {
3914 blob.content_id()
3915 .map(|content_id| (content_id, blob))
3916 .map_err(|error| invalid_snapshot_error(error.to_string()))
3917 })
3918 .collect::<Result<BTreeMap<_, _>, _>>()?;
3919 let store = Format8ScalePageStore {
3920 pages: RefCell::new(BTreeMap::new()),
3921 decision_blobs: RefCell::new(decision_blobs),
3922 provider_calls: RefCell::new(0),
3923 };
3924 let mut simulation = Simulation::new(8, Scenario::new(SimTime::EPOCH, Vec::new()))?;
3925 simulation.state.current.decisions = fixture.state;
3926 simulation.state.metadata.plugin_registration_closed = true;
3927 simulation.state.metadata.commitment_cache = None;
3928 simulation.refresh_checkpoint_hash()?;
3929
3930 let prepared = simulation.prepare_paged_checkpoint(None, &store)?;
3931 let initial_provider_calls = *store.provider_calls.borrow();
3932 let initial_delta_pages = prepared.delta.new_pages.len() as u64;
3933 prepared.store_and_verify(&store)?;
3934 let root_page = store
3935 .load_state_page(&prepared.checkpoint.root_page_id)?
3936 .ok_or_else(|| invalid_snapshot_error("Format-8 scale root page is unavailable"))?;
3937 let envelope: PagedCheckpointEnvelope = serde_json::from_slice(&root_page.bytes)
3938 .map_err(|error| invalid_snapshot_error(format!("invalid scale envelope: {error}")))?;
3939 let manifest_page = store
3940 .load_state_page(&envelope.decision_manifest_page_id)?
3941 .ok_or_else(|| invalid_snapshot_error("Format-8 scale manifest page is unavailable"))?;
3942 let manifest: PagedDecisionManifest = serde_json::from_slice(&manifest_page.bytes)
3943 .map_err(|error| invalid_snapshot_error(format!("invalid scale manifest: {error}")))?;
3944
3945 let restored = Simulation::from_paged_checkpoint(&prepared.checkpoint, &store)?;
3946 let replayed =
3947 Simulation::from_paged_checkpoint_and_journal(&prepared.checkpoint, &store, Vec::new())?;
3948 let samples = [
3949 1_usize,
3950 decision_count.saturating_div(2).max(1),
3951 decision_count,
3952 ]
3953 .into_iter()
3954 .filter(|ordinal| *ordinal <= decision_count)
3955 .collect::<BTreeSet<_>>();
3956 for ordinal in &samples {
3957 let key = DecisionHistoryKey::Attempt(canwu_core::DecisionRequestId::new(
3958 u64::try_from(*ordinal)
3959 .map_err(|_| invalid_snapshot_error("decision scale sample exceeds u64"))?,
3960 ));
3961 if !matches!(
3962 restored.decision_history_location_with_provider(&key, &store)?,
3963 DecisionHistoryLocation::Archived { .. }
3964 ) || !matches!(
3965 replayed.decision_history_location_with_provider(&key, &store)?,
3966 DecisionHistoryLocation::Archived { .. }
3967 ) {
3968 return Err(invalid_snapshot_error(
3969 "paged checkpoint scale restore lost exact decision history",
3970 ));
3971 }
3972 }
3973 *store.provider_calls.borrow_mut() = 0;
3974 let repeat = restored.prepare_paged_checkpoint(Some(&prepared.checkpoint), &store)?;
3975 let repeat_provider_calls = *store.provider_calls.borrow();
3976 let replay_repeat = replayed.prepare_paged_checkpoint(Some(&prepared.checkpoint), &store)?;
3977 let mut changed = restored;
3978 let changed_request_id = canwu_core::DecisionRequestId::new(
3979 u64::try_from(decision_count)
3980 .map_err(|_| invalid_snapshot_error("decision scale count exceeds u64"))?
3981 .checked_add(1)
3982 .ok_or_else(|| invalid_snapshot_error("decision scale request ID overflowed"))?,
3983 );
3984 changed
3985 .state
3986 .current
3987 .decisions
3988 .append_attempt(super::DecisionAttemptRecord {
3989 request_id: changed_request_id,
3990 request_commitment: canonical_byte_hash(
3991 "canwu.format8.single-page-change.v1",
3992 &changed_request_id.get().to_be_bytes(),
3993 ),
3994 at: SimTime::from_minutes(
3995 i64::try_from(changed_request_id.get())
3996 .map_err(|_| invalid_snapshot_error("decision scale time exceeds i64"))?,
3997 ),
3998 revision_before: 0,
3999 expected_revision: 0,
4000 outcome: super::DecisionAttemptOutcome::Rejected {
4001 code: super::DecisionAttemptErrorCode::InvalidDecision,
4002 message: "Format-8 single locator-page change".to_owned(),
4003 },
4004 })
4005 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
4006 let changed_key = DecisionHistoryKey::Attempt(changed_request_id);
4007 let changed_archive = changed
4008 .state
4009 .current
4010 .decisions
4011 .prepare_decision_archive(std::slice::from_ref(&changed_key))
4012 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
4013 for blob in &changed_archive.blobs {
4014 store.decision_blobs.borrow_mut().insert(
4015 blob.content_id()
4016 .map_err(|error| invalid_snapshot_error(error.to_string()))?,
4017 blob.clone(),
4018 );
4019 }
4020 let verified_change = changed
4021 .state
4022 .current
4023 .decisions
4024 .verify_decision_archive(&changed_archive, &store)
4025 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
4026 changed.state.current.decisions = changed
4027 .state
4028 .current
4029 .decisions
4030 .commit_verified_decision_archive(&verified_change)
4031 .map_err(|error| invalid_snapshot_error(error.to_string()))?;
4032 changed.state.metadata.commitment_cache = None;
4033 changed.refresh_checkpoint_hash()?;
4034 *store.provider_calls.borrow_mut() = 0;
4035 let single_page_change =
4036 changed.prepare_paged_checkpoint(Some(&prepared.checkpoint), &store)?;
4037 let single_page_change_provider_calls = *store.provider_calls.borrow();
4038 let pages = store.pages.borrow();
4039 Ok(PagedCheckpointScaleMetrics {
4040 decision_entries: decision_metrics.entries,
4041 decision_locator: decision_metrics,
4042 state_pages: pages.len() as u64,
4043 decision_directory_pages: manifest.archive_directory_page_ids.len() as u64,
4044 max_state_page_bytes: pages
4045 .values()
4046 .map(|page| page.decoded_bytes)
4047 .max()
4048 .unwrap_or(0),
4049 initial_delta_pages,
4050 repeat_delta_pages: repeat.delta.new_pages.len() as u64,
4051 single_page_change_delta_pages: single_page_change.delta.new_pages.len() as u64,
4052 initial_provider_calls,
4053 repeat_provider_calls,
4054 single_page_change_provider_calls,
4055 exact_restart_queries: samples.len() as u64,
4056 restored_root_matches: repeat.checkpoint.root_page_id == prepared.checkpoint.root_page_id,
4057 replayed_root_matches: replay_repeat.checkpoint.root_page_id
4058 == prepared.checkpoint.root_page_id,
4059 root_page_id: prepared.checkpoint.root_page_id,
4060 })
4061}
4062
4063#[cfg(test)]
4064mod tests {
4065 #![allow(clippy::unnecessary_literal_bound, clippy::unnecessary_wraps)]
4066 use super::super::{
4067 BoundaryContext, BoundaryDirective, BoundaryPhase, BoundaryProposal,
4068 BoundarySystemContract, Command, CommandContext, CommandRequestId, DomainRecordClass,
4069 DomainRecordDraft, DomainRecordMutation, Issuer, KnowledgeHolderRef, KnowledgeOrigin,
4070 KnowledgeRecordDraft, KnowledgeRecordKind, KnowledgeSchemaId, KnowledgeWriteGrant,
4071 PluginActionDescriptor, PluginKnowledgeSchema, PluginRegistrar, SimulationView, StateKey,
4072 StateVisibility, SystemDirective, demo_scenario,
4073 };
4074 use super::*;
4075 use canwu_core::{BoundaryId, CommandId, DomainRecordKind, PersonId};
4076 use serde_json::{Map, Value, json};
4077 use std::cell::RefCell;
4078
4079 fn decision_directory_page(start: u32, len: usize) -> PagedDecisionDirectoryPage {
4080 PagedDecisionDirectoryPage {
4081 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
4082 archive_bucket_pages: (0..len)
4083 .map(|offset| {
4084 let ordinal = start + u32::try_from(offset).expect("test offset fits u32");
4085 (
4086 super::super::DecisionArchivePageKey {
4087 bucket: u16::try_from(ordinal / 256).expect("bucket fits"),
4088 segment: u8::try_from(ordinal % 256).expect("segment fits"),
4089 },
4090 canonical_byte_hash(
4091 "canwu.test.decision-directory-page.v1",
4092 &ordinal.to_be_bytes(),
4093 ),
4094 )
4095 })
4096 .collect(),
4097 }
4098 }
4099
4100 #[test]
4101 fn paged_decision_directory_rejects_noncanonical_cross_page_encodings() {
4102 let first = decision_directory_page(0, MAX_PAGED_DECISION_DIRECTORY_ENTRIES);
4103 let second = decision_directory_page(
4104 u32::try_from(MAX_PAGED_DECISION_DIRECTORY_ENTRIES).expect("directory bound fits u32"),
4105 1,
4106 );
4107 assemble_paged_decision_directory(vec![second.clone(), first.clone()])
4108 .expect_err("swapped directory pages must fail");
4109 assemble_paged_decision_directory(vec![decision_directory_page(0, 1), second])
4110 .expect_err("a short middle directory page must fail");
4111
4112 let duplicate_id = canonical_byte_hash("canwu.test.duplicate-directory.v1", b"same");
4113 let manifest = PagedDecisionManifest {
4114 format_version: PAGED_CHECKPOINT_FORMAT_VERSION,
4115 hot_page_id: canonical_byte_hash("canwu.test.hot.v1", b"hot"),
4116 archive_receipt_root: canonical_byte_hash("canwu.test.archive.v1", b"archive"),
4117 archive_receipt_count: 1,
4118 archive_directory_page_ids: vec![duplicate_id.clone(), duplicate_id],
4119 };
4120 validate_paged_decision_manifest(&manifest)
4121 .expect_err("duplicate directory page IDs must fail");
4122 }
4123
4124 #[test]
4125 fn archived_plugin_ingress_provenance_is_merkle_bound() {
4126 let provenance = ArchivedPluginIngressProvenance {
4127 plugin: "fixture-provider".to_owned(),
4128 packet_type: "recognized-practice".to_owned(),
4129 producer_boundary: BoundaryId::new(7),
4130 };
4131 let entry = EvidenceIndexEntry {
4132 reference: EvidenceRef::Ingress(super::super::IngressId::new(9)),
4133 item: EvidenceItemLocator {
4134 journal: EvidenceJournalKind::Ingress,
4135 absolute_index: 9,
4136 nested: EvidenceNestedLocator::None,
4137 },
4138 item_commitment: "0101010101010101010101010101010101010101010101010101010101010101"
4139 .to_owned(),
4140 plugin_ingress_provenance: Some(provenance.clone()),
4141 };
4142 let (root, proofs) = archive_merkle(std::slice::from_ref(&entry)).expect("build proof");
4143 let header = ArchivedSegmentHeader {
4144 segment_id: "0202020202020202020202020202020202020202020202020202020202020202"
4145 .to_owned(),
4146 start: EvidenceCursor {
4147 ingress_count: 8,
4148 ..EvidenceCursor::default()
4149 },
4150 end: EvidenceCursor {
4151 ingress_count: 9,
4152 ..EvidenceCursor::default()
4153 },
4154 journal_roots: EvidenceJournalRoots {
4155 events: archive_empty_root(),
4156 commands: archive_empty_root(),
4157 command_attempts: archive_empty_root(),
4158 ingress: archive_empty_root(),
4159 boundaries: archive_empty_root(),
4160 random_draws: archive_empty_root(),
4161 },
4162 evidence_index_root: root,
4163 evidence_index_entry_count: 1,
4164 };
4165 let mut receipt = ArchivedEvidenceReceipt {
4166 evidence: entry.reference,
4167 locator: ArchivedEvidenceLocator {
4168 segment_id: header.segment_id.clone(),
4169 item: entry.item,
4170 },
4171 evidence_index_leaf: 0,
4172 item_commitment: entry.item_commitment,
4173 plugin_ingress_provenance: Some(provenance),
4174 merkle_path: proofs.into_iter().next().expect("one proof"),
4175 };
4176 verify_archive_receipt(&receipt, &header).expect("canonical provenance proof");
4177
4178 receipt
4179 .plugin_ingress_provenance
4180 .as_mut()
4181 .expect("provenance")
4182 .plugin = "forged-provider".to_owned();
4183 assert_eq!(
4184 verify_archive_receipt(&receipt, &header)
4185 .expect_err("tampered provenance must break the proof")
4186 .code,
4187 ErrorCode::InvalidArchive
4188 );
4189 }
4190
4191 #[derive(Default)]
4192 struct TestArchive {
4193 segments: RefCell<BTreeMap<String, EvidenceJournalSegment>>,
4194 }
4195
4196 impl TestArchive {
4197 fn segment_ids(&self) -> Vec<String> {
4198 self.segments.borrow().keys().cloned().collect()
4199 }
4200 }
4201
4202 impl ArchiveProvider for TestArchive {
4203 fn load_evidence_segment(
4204 &self,
4205 segment_id: &str,
4206 ) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
4207 Ok(self.segments.borrow().get(segment_id).cloned())
4208 }
4209 }
4210
4211 impl ArchiveStore for TestArchive {
4212 fn store_evidence_segment(
4213 &self,
4214 segment: &EvidenceJournalSegment,
4215 ) -> Result<ArchiveStoreOutcome, CanwuError> {
4216 let segment_id = segment
4217 .archive
4218 .as_ref()
4219 .ok_or_else(|| archive_error("test archive segment has no index"))?
4220 .header
4221 .segment_id
4222 .clone();
4223 let mut segments = self.segments.borrow_mut();
4224 if let Some(existing) = segments.get(&segment_id) {
4225 return if existing == segment {
4226 Ok(ArchiveStoreOutcome::AlreadyPresent)
4227 } else {
4228 Err(archive_error(
4229 "content-addressed test segment ID has conflicting bytes",
4230 ))
4231 };
4232 }
4233 segments.insert(segment_id, segment.clone());
4234 Ok(ArchiveStoreOutcome::Stored)
4235 }
4236 }
4237
4238 fn archived_identity_schema() -> KnowledgeSchemaId {
4239 KnowledgeSchemaId::new(
4240 KnowledgeRecordKind::new("fixture.archive", "archived_command_notice"),
4241 1,
4242 )
4243 }
4244
4245 #[allow(clippy::unnecessary_wraps)]
4246 fn retain_archive_source_command(
4247 _view: &SimulationView<'_>,
4248 _context: &CommandContext,
4249 _payload: &Value,
4250 ) -> Result<Vec<SystemDirective>, CanwuError> {
4251 Ok(Vec::new())
4252 }
4253
4254 #[allow(clippy::unnecessary_wraps)]
4255 fn publish_archived_command_identity(
4256 _view: &SimulationView<'_>,
4257 context: &BoundaryContext,
4258 ) -> Result<BoundaryProposal, CanwuError> {
4259 let record = KnowledgeRecordDraft {
4260 schema: archived_identity_schema(),
4261 subjects: Vec::new(),
4262 payload: json!({ "boundary": context.boundary_id.get() }),
4263 as_of: None,
4264 confidence_per_mille: 1_000,
4265 origin: KnowledgeOrigin {
4266 method: "archived_command_identity_v1".to_owned(),
4267 evidence: vec![EvidenceRef::Command(CommandId::new(1))],
4268 },
4269 supersedes: Vec::new(),
4270 contradicts: Vec::new(),
4271 };
4272 Ok(BoundaryProposal {
4273 directives: vec![BoundaryDirective::PublishKnowledge {
4274 holder: KnowledgeHolderRef::Person(PersonId::new(1)),
4275 visibility: StateVisibility::SameBoundary,
4276 producer_correlation: Some(format!(
4277 "archived-command-boundary-{}",
4278 context.boundary_id.get()
4279 )),
4280 records: vec![record],
4281 summary: "Publish knowledge from a retained or archived command identity"
4282 .to_owned(),
4283 }],
4284 ..BoundaryProposal::default()
4285 })
4286 }
4287
4288 struct ArchivedIdentityPublicationPlugin;
4289
4290 impl SimulationPlugin for ArchivedIdentityPublicationPlugin {
4291 fn name(&self) -> &str {
4292 "fixture-archived-identity-publication"
4293 }
4294
4295 fn version(&self) -> &str {
4296 "test-v1"
4297 }
4298
4299 fn semantic_hash(&self) -> &str {
4300 "7100000000000000000000000000000000000000000000000000000000000000"
4301 }
4302
4303 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
4304 registrar.register_knowledge_schema(PluginKnowledgeSchema {
4305 id: archived_identity_schema(),
4306 schema_hash: "7200000000000000000000000000000000000000000000000000000000000000"
4307 .to_owned(),
4308 writable: true,
4309 payload_schema: PayloadSchema::Any,
4310 subjects: Vec::new(),
4311 })?;
4312 registrar.register_command(
4313 PluginActionDescriptor {
4314 name: "retain_archive_source_v1".to_owned(),
4315 description: "Persist one neutral command identity for archive testing"
4316 .to_owned(),
4317 payload_schema: PayloadSchema::Any,
4318 reads: Vec::new(),
4319 writes: Vec::new(),
4320 },
4321 retain_archive_source_command,
4322 )?;
4323 let mut publisher = BoundarySystemContract::new(
4324 "publish-archived-command-identity",
4325 BoundaryPhase::PerspectiveAndReportMaterialization,
4326 SystemCadence::Daily,
4327 );
4328 publisher.knowledge_writes = vec![KnowledgeWriteGrant {
4329 schema: archived_identity_schema(),
4330 visibilities: vec![StateVisibility::SameBoundary],
4331 }];
4332 registrar.register_boundary_system(publisher, publish_archived_command_identity)
4333 }
4334 }
4335
4336 fn archived_identity_two_segment_fixture() -> (SimulationCheckpoint, Vec<EvidenceJournalSegment>)
4337 {
4338 let (scenario, _) = demo_scenario();
4339 let plugin = ArchivedIdentityPublicationPlugin;
4340 let mut simulation = Simulation::new(711, scenario).expect("fixture scenario should load");
4341 simulation
4342 .register_plugin(&plugin)
4343 .expect("archived-identity plugin should register");
4344 simulation
4345 .enqueue_command(
4346 SimTime::EPOCH,
4347 0,
4348 CommandRequest::new(
4349 CommandRequestId::new(1),
4350 simulation.revision(),
4351 CommandEnvelope::new(
4352 Issuer::System("archive-identity-fixture".to_owned()),
4353 Command::Plugin {
4354 plugin: plugin.name().to_owned(),
4355 command: "retain_archive_source_v1".to_owned(),
4356 payload: json!({ "format_version": 1 }),
4357 },
4358 )
4359 .at_time(SimTime::EPOCH),
4360 ),
4361 )
4362 .expect("archive source command should queue");
4363 simulation
4364 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
4365 .expect("boundary one should retain the command-backed publication");
4366 simulation
4367 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4368 .expect("the following cut should admit boundary-one events before sealing");
4369
4370 let mut compact = simulation
4371 .into_compacted()
4372 .expect("the fixture should enter compact mode");
4373 let first_segment = compact
4374 .seal_evidence()
4375 .expect("boundary-one evidence should seal")
4376 .expect("boundary one should produce an archive segment");
4377 let first_checkpoint = compact.checkpoint().expect("checkpoint one should build");
4378 assert!(
4379 first_checkpoint
4380 .archived_evidence_receipts
4381 .iter()
4382 .any(|receipt| receipt.evidence == EvidenceRef::Command(CommandId::new(1)))
4383 );
4384
4385 let mut restored = CompactedSimulation::from_checkpoint_and_journal_with_plugins(
4386 first_checkpoint,
4387 vec![first_segment.clone()],
4388 &[&plugin],
4389 )
4390 .expect("checkpoint one should restore with its exact archive prefix");
4391 let restored_prefix = restored
4392 .seal_evidence()
4393 .expect("restored boundary-one evidence should reseal")
4394 .expect("restored boundary-one evidence should remain non-empty");
4395 assert_eq!(restored_prefix, first_segment);
4396 assert!(
4397 restored
4398 .archived_evidence_receipt(&EvidenceRef::Command(CommandId::new(1)))
4399 .is_some(),
4400 "the second publication must consume an archived identity, not retained payload"
4401 );
4402 let archived_reads = [StateKey::core_commands()];
4403 let archived_view = restored
4404 .simulation
4405 .plugin_view("archive-identity-probe", &archived_reads);
4406 let error = archived_view
4407 .command(CommandId::new(1))
4408 .expect_err("archived identity must not expose retained command payload");
4409 assert_eq!(error.code, ErrorCode::EvidenceContentUnavailable);
4410 restored
4411 .settle_boundary(
4412 BoundaryRequest::at(SimTime::EPOCH + SimDuration::days(1))
4413 .with_cadence(SystemCadence::Daily),
4414 )
4415 .expect("boundary two should accept the archived command identity");
4416 let holder = KnowledgeHolderRef::Person(PersonId::new(1));
4417 let records = restored
4418 .knowledge()
4419 .for_holder(&holder)
4420 .expect("the holder should have both publications");
4421 assert_eq!(records.len(), 2);
4422 assert!(records.values().all(|record| {
4423 record.origin.evidence == vec![EvidenceRef::Command(CommandId::new(1))]
4424 }));
4425 restored
4426 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH + SimDuration::days(1)))
4427 .expect("the following cut should admit boundary-two events before sealing");
4428
4429 let second_segment = restored
4430 .seal_evidence()
4431 .expect("boundary-two evidence should seal")
4432 .expect("boundary two should produce an archive segment");
4433 let second_checkpoint = restored.checkpoint().expect("checkpoint two should build");
4434 (second_checkpoint, vec![restored_prefix, second_segment])
4435 }
4436
4437 #[test]
4438 fn compact_restore_preserves_archived_identity_and_rejects_noncontiguous_segments() {
4439 let plugin = ArchivedIdentityPublicationPlugin;
4440 let (checkpoint, segments) = archived_identity_two_segment_fixture();
4441 let restored = CompactedSimulation::from_checkpoint_and_journal_with_plugins(
4442 checkpoint.clone(),
4443 segments.clone(),
4444 &[&plugin],
4445 )
4446 .expect("the complete two-segment archive should restore");
4447 assert_eq!(
4448 restored
4449 .knowledge()
4450 .for_holder(&KnowledgeHolderRef::Person(PersonId::new(1)))
4451 .expect("published holder ledger")
4452 .len(),
4453 2
4454 );
4455
4456 let first = segments[0].clone();
4457 let second = segments[1].clone();
4458 for (label, tampered) in [
4459 ("omission", vec![first.clone()]),
4460 (
4461 "overlap",
4462 vec![first.clone(), first.clone(), second.clone()],
4463 ),
4464 ("reorder", vec![second, first]),
4465 ] {
4466 let error = Simulation::from_checkpoint_and_journal(checkpoint.clone(), tampered)
4467 .err()
4468 .expect(label);
4469 assert_eq!(error.code, ErrorCode::InvalidSnapshot, "{label}");
4470 }
4471 }
4472
4473 struct PayloadContinuationPlugin;
4474
4475 fn continuation_record_ref() -> DomainRecordRef {
4476 DomainRecordRef {
4477 kind: DomainRecordKind::new("fixture.archive", "payload_continuation"),
4478 id: "primary".to_owned(),
4479 }
4480 }
4481
4482 fn continuation_payload(continuation: PayloadRequiredEvidenceContinuationV1) -> Value {
4483 Value::Object(Map::from_iter([(
4484 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD.to_owned(),
4485 serde_json::to_value(continuation).expect("fixture continuation should encode"),
4486 )]))
4487 }
4488
4489 fn mutate_payload_continuation(
4490 _view: &SimulationView<'_>,
4491 context: &BoundaryContext,
4492 ) -> Result<BoundaryProposal, CanwuError> {
4493 let directive = match context.boundary_id.get() {
4494 1 => Some(BoundaryDirective::MutateRecord {
4495 mutation: DomainRecordMutation::Create {
4496 record: DomainRecordDraft::new(
4497 continuation_record_ref(),
4498 continuation_payload(PayloadRequiredEvidenceContinuationV1::active(vec![
4499 EvidenceRef::Boundary(BoundaryId::new(1)),
4500 ])),
4501 ),
4502 },
4503 summary: "Create an active payload continuation".to_owned(),
4504 }),
4505 3 => Some(BoundaryDirective::MutateRecord {
4506 mutation: DomainRecordMutation::Update {
4507 record: DomainRecordDraft::new(
4508 continuation_record_ref(),
4509 continuation_payload(PayloadRequiredEvidenceContinuationV1::completed()),
4510 ),
4511 expected_version: 1,
4512 },
4513 summary: "Complete the payload continuation".to_owned(),
4514 }),
4515 _ => None,
4516 };
4517 Ok(BoundaryProposal {
4518 directives: directive.into_iter().collect(),
4519 ..BoundaryProposal::default()
4520 })
4521 }
4522
4523 impl SimulationPlugin for PayloadContinuationPlugin {
4524 fn name(&self) -> &str {
4525 "fixture-payload-continuation"
4526 }
4527
4528 fn version(&self) -> &str {
4529 "test-v1"
4530 }
4531
4532 fn semantic_hash(&self) -> &str {
4533 "7000000000000000000000000000000000000000000000000000000000000000"
4534 }
4535
4536 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
4537 let mut properties = BTreeMap::new();
4538 properties.insert(
4539 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD.to_owned(),
4540 payload_required_evidence_continuation_property_v1(),
4541 );
4542 let mut schema =
4543 DomainRecordSchema::new(continuation_record_ref().kind, DomainRecordClass::Record);
4544 schema.payload_schema = PayloadSchema::Object {
4545 properties,
4546 allow_additional: false,
4547 };
4548 let state = schema.state_key();
4549 registrar.register_record_schema(schema)?;
4550 let mut contract = BoundarySystemContract::new(
4551 "payload-continuation",
4552 BoundaryPhase::DomainDeltaProposal,
4553 SystemCadence::Daily,
4554 );
4555 contract.writes = vec![state];
4556 contract.visibility = StateVisibility::SameBoundary;
4557 registrar.register_boundary_system(contract, mutate_payload_continuation)
4558 }
4559 }
4560
4561 fn payload_continuation_runtime() -> CompactedSimulation {
4562 let (scenario, _) = demo_scenario();
4563 let mut simulation = Simulation::new(701, scenario).expect("fixture scenario should load");
4564 simulation
4565 .register_plugin(&PayloadContinuationPlugin)
4566 .expect("payload-continuation plugin should register");
4567 simulation
4568 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
4569 .expect("the active continuation should be created");
4570 simulation
4571 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4572 .expect("a later boundary should admit the record-change event");
4573 simulation
4574 .into_compacted()
4575 .expect("the fixture should enter compact mode")
4576 }
4577
4578 fn store_prepared(
4579 compact: &CompactedSimulation,
4580 archive: &TestArchive,
4581 ) -> PreparedEvidenceSeal {
4582 let prepared = compact
4583 .prepare_evidence_seal()
4584 .expect("the fixture should prepare a seal")
4585 .expect("the retained tail should be non-empty");
4586 assert_eq!(
4587 archive
4588 .store_evidence_segment(&prepared.segment)
4589 .expect("the prepared segment should store"),
4590 ArchiveStoreOutcome::Stored
4591 );
4592 prepared
4593 }
4594
4595 #[test]
4596 fn payload_required_receipts_are_exactly_reachable_and_prune_after_completion() {
4597 let mut compact = payload_continuation_runtime();
4598 assert_eq!(
4599 compact
4600 .seal_evidence()
4601 .expect_err("direct sealing must reject payload continuations")
4602 .code,
4603 ErrorCode::ArchiveNotReady
4604 );
4605
4606 let archive = TestArchive::default();
4607 let first = store_prepared(&compact, &archive);
4608 compact
4609 .commit_evidence_seal(&first.token, &archive)
4610 .expect("provider-backed sealing should commit");
4611 let first_checkpoint = compact.checkpoint().expect("checkpoint should build");
4612 let first_references = first_checkpoint
4613 .archived_evidence_receipts
4614 .iter()
4615 .map(|receipt| receipt.evidence.clone())
4616 .collect::<BTreeSet<_>>();
4617 let first_dependencies = first_checkpoint
4618 .evidence_dependencies
4619 .iter()
4620 .map(|dependency| (dependency.reference.clone(), dependency.requirement))
4621 .collect::<BTreeMap<_, _>>();
4622 assert_eq!(
4623 first_dependencies.get(&EvidenceRef::Boundary(BoundaryId::new(1))),
4624 Some(&EvidenceRequirement::PayloadRequired)
4625 );
4626 assert_eq!(first_references.len(), first_dependencies.len());
4627 assert_eq!(
4628 first_references,
4629 first_dependencies.keys().cloned().collect::<BTreeSet<_>>()
4630 );
4631 assert!(
4632 first
4633 .segment
4634 .archive
4635 .as_ref()
4636 .expect("prepared segment should have an archive index")
4637 .entries
4638 .len()
4639 > first_references.len(),
4640 "the full segment index must outlive the reachable-only receipt set"
4641 );
4642
4643 compact
4644 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH).with_cadence(SystemCadence::Daily))
4645 .expect("the continuation should complete");
4646 compact
4647 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4648 .expect("a later boundary should admit the completion event");
4649 let second = store_prepared(&compact, &archive);
4650 compact
4651 .commit_evidence_seal(&second.token, &archive)
4652 .expect("the completed continuation should seal");
4653 let second_checkpoint = compact.checkpoint().expect("checkpoint should build");
4654 assert_eq!(second_checkpoint.archived_segment_headers.len(), 2);
4655 assert_eq!(second_checkpoint.archived_evidence_receipts.len(), 1);
4656 assert_eq!(second_checkpoint.evidence_dependencies.len(), 1);
4657 assert_eq!(
4658 second_checkpoint.archived_evidence_receipts[0].evidence,
4659 second_checkpoint.evidence_dependencies[0].reference
4660 );
4661 assert_eq!(
4662 second_checkpoint.evidence_dependencies[0].requirement,
4663 EvidenceRequirement::IdentityOnly
4664 );
4665 assert!(
4666 !second_checkpoint
4667 .archived_evidence_receipts
4668 .iter()
4669 .any(|receipt| receipt.evidence == EvidenceRef::Boundary(BoundaryId::new(1)))
4670 );
4671
4672 Simulation::from_checkpoint_and_journal(
4673 second_checkpoint,
4674 vec![first.segment, second.segment],
4675 )
4676 .expect("reconstruction must filter full segment indexes to the compact receipt frontier");
4677 }
4678
4679 #[test]
4680 fn payload_required_commit_fails_closed_when_an_older_segment_is_missing() {
4681 let mut compact = payload_continuation_runtime();
4682 let complete_archive = TestArchive::default();
4683 let first = store_prepared(&compact, &complete_archive);
4684 compact
4685 .commit_evidence_seal(&first.token, &complete_archive)
4686 .expect("the initial provider-backed seal should commit");
4687
4688 compact
4689 .settle_boundary(BoundaryRequest::at(SimTime::EPOCH))
4690 .expect("a new tail should preserve the active continuation");
4691 let incomplete_archive = TestArchive::default();
4692 let second = store_prepared(&compact, &incomplete_archive);
4693 let before = compact.checkpoint().expect("checkpoint should build");
4694 let error = compact
4695 .commit_evidence_seal(&second.token, &incomplete_archive)
4696 .expect_err("the provider must retain every payload-required segment");
4697 assert_eq!(error.code, ErrorCode::EvidenceContentUnavailable);
4698 assert_eq!(
4699 compact.checkpoint().expect("checkpoint should build"),
4700 before
4701 );
4702
4703 incomplete_archive
4704 .store_evidence_segment(&first.segment)
4705 .expect("restoring the older required segment should succeed");
4706 compact
4707 .commit_evidence_seal(&second.token, &incomplete_archive)
4708 .expect("the exact provider set should permit commit");
4709 }
4710
4711 #[test]
4712 fn host_orphan_candidates_follow_all_retained_manifests_without_deleting() {
4713 let mut compact = payload_continuation_runtime();
4714 let archive = TestArchive::default();
4715 let prepared = store_prepared(&compact, &archive);
4716 let stored_ids = archive.segment_ids();
4717 let before = compact.checkpoint().expect("checkpoint should build");
4718 assert_eq!(
4719 SimulationCheckpoint::orphaned_archive_segment_ids(
4720 std::slice::from_ref(&before),
4721 &stored_ids,
4722 )
4723 .expect("manifest reachability should validate"),
4724 vec![prepared.token.segment_id.clone()]
4725 );
4726 assert!(
4727 archive
4728 .load_evidence_segment(&prepared.token.segment_id)
4729 .expect("the store should remain readable")
4730 .is_some(),
4731 "the conformance API must not delete host content"
4732 );
4733
4734 compact
4735 .commit_evidence_seal(&prepared.token, &archive)
4736 .expect("the stored segment should commit");
4737 let after = compact.checkpoint().expect("checkpoint should build");
4738 assert!(
4739 SimulationCheckpoint::orphaned_archive_segment_ids(
4740 std::slice::from_ref(&after),
4741 &stored_ids,
4742 )
4743 .expect("committed manifest reachability should validate")
4744 .is_empty()
4745 );
4746 assert_eq!(
4747 SimulationCheckpoint::reachable_archive_segment_ids(&[before, after])
4748 .expect("all retained manifests should be scanned"),
4749 BTreeSet::from([prepared.token.segment_id.clone()])
4750 );
4751 }
4752}