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