Skip to main content

canwu_sim/runtime/
persistence.rs

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