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