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