Skip to main content

canwu_sim/
persistence.rs

1use super::{
2    ADMISSION_CURSOR_FORMAT_VERSION, BoundaryReceipt, BoundaryRecord, BoundaryRequest, CanwuError,
3    CommandAttemptRecord, CommandEnvelope, CommandOutcome, CommandReceipt, CommandRecord,
4    CommandRequest, DomainRecord, DomainRecordRef, DomainRecordType, ENGINE_VERSION, ErrorCode,
5    IngressPayload, IngressReceipt, IngressRecord, KnowledgeSnapshot, PluginIngressRequest,
6    RandomDrawRecord, ReplayJournal, RuntimeEvidence, SNAPSHOT_FORMAT_VERSION,
7    STATE_REVISION_FORMAT_VERSION, ScheduledRecord, SimDuration, SimEvent, SimTime, Simulation,
8    SimulationPlugin, SimulationSnapshot, SystemCadence, TypedDomainRecordRef, WorldSnapshot,
9    has_unqueued_command_history, invalid_snapshot_error,
10};
11use crate::state::{ArchivedCommandRequestOutcome, ArchivedIngressRequest};
12use serde::{Deserialize, Serialize};
13
14/// Version of current-state checkpoints plus append-only evidence segments.
15pub const CHECKPOINT_JOURNAL_FORMAT_VERSION: u32 = 1;
16
17#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
18/// Monotonic cuts through every append-only evidence journal.
19pub struct EvidenceCursor {
20    pub event_count: u64,
21    pub command_count: u64,
22    pub command_attempt_count: u64,
23    pub ingress_count: u64,
24    pub boundary_count: u64,
25    pub random_draw_count: u64,
26}
27
28impl EvidenceCursor {
29    fn from_evidence(evidence: &RuntimeEvidence) -> Result<Self, CanwuError> {
30        let count = |len: usize, label: &str| {
31            u64::try_from(len).map_err(|_| {
32                CanwuError::new(
33                    ErrorCode::IdentifierExhausted,
34                    format!("{label} journal length exceeds the persistent cursor space"),
35                )
36            })
37        };
38        Ok(Self {
39            event_count: evidence
40                .archived
41                .event_count
42                .checked_add(count(evidence.events.len(), "event")?)
43                .ok_or_else(|| invalid_snapshot_error("event journal cursor is exhausted"))?,
44            command_count: evidence
45                .archived
46                .command_count
47                .checked_add(count(evidence.commands.len(), "command")?)
48                .ok_or_else(|| invalid_snapshot_error("command journal cursor is exhausted"))?,
49            command_attempt_count: evidence
50                .archived
51                .command_attempt_count
52                .checked_add(count(evidence.command_attempts.len(), "command-attempt")?)
53                .ok_or_else(|| {
54                    invalid_snapshot_error("command-attempt journal cursor is exhausted")
55                })?,
56            ingress_count: evidence
57                .archived
58                .ingress_count
59                .checked_add(count(evidence.ingress.len(), "ingress")?)
60                .ok_or_else(|| invalid_snapshot_error("ingress journal cursor is exhausted"))?,
61            boundary_count: evidence
62                .archived
63                .boundary_count
64                .checked_add(count(evidence.boundaries.len(), "boundary")?)
65                .ok_or_else(|| invalid_snapshot_error("boundary journal cursor is exhausted"))?,
66            random_draw_count: evidence
67                .archived
68                .random_draw_count
69                .checked_add(count(evidence.random_draws.len(), "random-draw")?)
70                .ok_or_else(|| invalid_snapshot_error("random-draw journal cursor is exhausted"))?,
71        })
72    }
73
74    pub(super) fn checked_advance(
75        self,
76        segment: &EvidenceJournalSegment,
77    ) -> Result<Self, CanwuError> {
78        let advance = |value: u64, len: usize, label: &str| {
79            value
80                .checked_add(u64::try_from(len).map_err(|_| {
81                    invalid_snapshot_error(format!(
82                        "{label} journal segment exceeds the persistent cursor space"
83                    ))
84                })?)
85                .ok_or_else(|| {
86                    invalid_snapshot_error(format!(
87                        "{label} journal cursor exceeds the persistent cursor space"
88                    ))
89                })
90        };
91        Ok(Self {
92            event_count: advance(self.event_count, segment.events.len(), "event")?,
93            command_count: advance(self.command_count, segment.commands.len(), "command")?,
94            command_attempt_count: advance(
95                self.command_attempt_count,
96                segment.command_attempts.len(),
97                "command-attempt",
98            )?,
99            ingress_count: advance(self.ingress_count, segment.ingress.len(), "ingress")?,
100            boundary_count: advance(self.boundary_count, segment.boundaries.len(), "boundary")?,
101            random_draw_count: advance(
102                self.random_draw_count,
103                segment.random_draws.len(),
104                "random-draw",
105            )?,
106        })
107    }
108}
109
110#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
111/// Current authoritative state plus the journal cut required to validate it.
112///
113/// `state` deliberately contains empty append-only evidence arrays. It is not a
114/// standalone `SimulationSnapshot`; load it only with the contiguous evidence
115/// segments ending at `journal_end`.
116pub struct SimulationCheckpoint {
117    pub format_version: u32,
118    pub journal_end: EvidenceCursor,
119    pub state: SimulationSnapshot,
120}
121
122#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
123/// One contiguous append-only evidence range for incremental archival.
124pub struct EvidenceJournalSegment {
125    pub format_version: u32,
126    pub start: EvidenceCursor,
127    pub end: EvidenceCursor,
128    #[serde(default, skip_serializing_if = "Vec::is_empty")]
129    pub events: Vec<SimEvent>,
130    #[serde(default, skip_serializing_if = "Vec::is_empty")]
131    pub commands: Vec<CommandRecord>,
132    #[serde(default, skip_serializing_if = "Vec::is_empty")]
133    pub command_attempts: Vec<CommandAttemptRecord>,
134    #[serde(default, skip_serializing_if = "Vec::is_empty")]
135    pub ingress: Vec<IngressRecord>,
136    #[serde(default, skip_serializing_if = "Vec::is_empty")]
137    pub boundaries: Vec<BoundaryRecord>,
138    #[serde(default, skip_serializing_if = "Vec::is_empty")]
139    pub random_draws: Vec<RandomDrawRecord>,
140}
141
142#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
143/// Portable full-save bundle built from a current-state checkpoint and journal segments.
144pub struct CheckpointJournal {
145    pub checkpoint: SimulationCheckpoint,
146    pub segments: Vec<EvidenceJournalSegment>,
147}
148
149/// A live simulation whose sealed evidence prefixes are owned by the caller.
150///
151/// This opt-in runtime preserves current authoritative state, deterministic
152/// commitments, idempotency, and continuation behavior while retaining only
153/// the evidence appended since the most recent seal. Every returned segment is
154/// part of the permanent replay record and must be stored contiguously by the
155/// caller.
156pub struct CompactedSimulation {
157    simulation: Simulation,
158}
159
160impl CompactedSimulation {
161    /// Returns the monotonic cut through sealed and retained evidence.
162    pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
163        self.simulation.evidence_cursor()
164    }
165
166    /// Captures current state and the total journal cut without cloning sealed evidence.
167    pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
168        self.simulation.checkpoint()
169    }
170
171    /// Seals and releases the current retained evidence tail.
172    ///
173    /// The runtime changes only after the segment is fully constructed and its
174    /// continuation indexes are prepared. An empty retained tail returns
175    /// `None`. The caller owns persistence and must keep all non-empty segments
176    /// in exact cursor order for save restoration or replay.
177    pub fn seal_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
178        self.simulation.seal_retained_evidence()
179    }
180
181    #[must_use]
182    pub const fn time(&self) -> SimTime {
183        self.simulation.time()
184    }
185
186    #[must_use]
187    pub const fn revision(&self) -> u64 {
188        self.simulation.revision()
189    }
190
191    #[must_use]
192    pub fn checkpoint_hash(&self) -> &str {
193        self.simulation.checkpoint_hash()
194    }
195
196    #[must_use]
197    pub fn boundary_head_hash(&self) -> Option<&str> {
198        self.simulation.boundary_head_hash()
199    }
200
201    #[must_use]
202    pub fn world(&self) -> WorldSnapshot {
203        self.simulation.world()
204    }
205
206    #[must_use]
207    pub fn knowledge(&self) -> &KnowledgeSnapshot {
208        self.simulation.knowledge()
209    }
210
211    #[must_use]
212    pub fn domain_record(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
213        self.simulation.domain_record(reference)
214    }
215
216    #[must_use]
217    pub fn typed_domain_record<T: DomainRecordType>(
218        &self,
219        reference: &TypedDomainRecordRef<T>,
220    ) -> Option<&DomainRecord> {
221        self.simulation.typed_domain_record(reference)
222    }
223
224    pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
225        self.simulation.submit(envelope)
226    }
227
228    pub fn process_command(
229        &mut self,
230        request: CommandRequest,
231    ) -> Result<CommandOutcome, CanwuError> {
232        self.simulation.process_command(request)
233    }
234
235    pub fn enqueue_command(
236        &mut self,
237        due_at: SimTime,
238        priority: i32,
239        request: CommandRequest,
240    ) -> Result<IngressReceipt, CanwuError> {
241        self.simulation.enqueue_command(due_at, priority, request)
242    }
243
244    pub fn enqueue_plugin_ingress(
245        &mut self,
246        request: PluginIngressRequest,
247    ) -> Result<IngressReceipt, CanwuError> {
248        self.simulation.enqueue_plugin_ingress(request)
249    }
250
251    pub fn schedule_calendar_boundary(
252        &mut self,
253        due_at: SimTime,
254        cadences: Vec<SystemCadence>,
255    ) -> Result<IngressReceipt, CanwuError> {
256        self.simulation.schedule_calendar_boundary(due_at, cadences)
257    }
258
259    pub fn advance(&mut self, duration: SimDuration) -> Result<Vec<SimEvent>, CanwuError> {
260        self.simulation.advance(duration)
261    }
262
263    pub fn advance_canonical(
264        &mut self,
265        duration: SimDuration,
266    ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
267        self.simulation.advance_canonical(duration)
268    }
269
270    pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
271        self.simulation.step_canonical()
272    }
273
274    pub fn settle_boundary(
275        &mut self,
276        request: BoundaryRequest,
277    ) -> Result<BoundaryReceipt, CanwuError> {
278        self.simulation.settle_boundary(request)
279    }
280
281    /// Reconstructs a validated full snapshot from the supplied sealed prefix
282    /// plus the currently retained tail.
283    pub fn snapshot_with_segments(
284        &self,
285        mut segments: Vec<EvidenceJournalSegment>,
286    ) -> Result<SimulationSnapshot, CanwuError> {
287        let tail = self
288            .simulation
289            .journal_segment_since(self.simulation.state.evidence.archived)?;
290        if tail.start != tail.end {
291            segments.push(tail);
292        }
293        let snapshot =
294            Simulation::snapshot_from_checkpoint_and_journal(self.checkpoint()?, segments)?;
295        Simulation::from_snapshot(snapshot.clone())?;
296        Ok(snapshot)
297    }
298
299    /// Produces the ordinary exact-replay journal after validating the supplied archive.
300    pub fn replay_journal_with_segments(
301        &self,
302        segments: Vec<EvidenceJournalSegment>,
303    ) -> Result<ReplayJournal, CanwuError> {
304        let snapshot = self.snapshot_with_segments(segments)?;
305        let simulation = Simulation::from_snapshot(snapshot)?;
306        Ok(simulation.replay_journal())
307    }
308
309    /// Restores and validates a checkpoint plus its archive, then enters the
310    /// compact interface with that evidence retained until the caller seals it.
311    pub fn from_checkpoint_and_journal(
312        checkpoint: SimulationCheckpoint,
313        segments: Vec<EvidenceJournalSegment>,
314    ) -> Result<Self, CanwuError> {
315        Simulation::from_checkpoint_and_journal(checkpoint, segments)?.into_compacted()
316    }
317
318    /// Restores a compact runtime and rehydrates its exact executable plugins.
319    pub fn from_checkpoint_and_journal_with_plugins(
320        checkpoint: SimulationCheckpoint,
321        segments: Vec<EvidenceJournalSegment>,
322        plugins: &[&dyn SimulationPlugin],
323    ) -> Result<Self, CanwuError> {
324        let mut simulation = Simulation::from_checkpoint_and_journal(checkpoint, segments)?;
325        for plugin in plugins {
326            simulation.register_plugin(*plugin)?;
327        }
328        simulation.ensure_runtime_ready()?;
329        simulation.into_compacted()
330    }
331}
332
333impl Simulation {
334    /// Converts this runtime into the opt-in compact journal interface.
335    ///
336    /// Conversion itself preserves the complete retained history. Call
337    /// [`CompactedSimulation::seal_evidence`] to release a validated segment
338    /// explicitly.
339    pub fn into_compacted(self) -> Result<CompactedSimulation, CanwuError> {
340        self.ensure_runtime_ready()?;
341        Ok(CompactedSimulation { simulation: self })
342    }
343
344    fn ensure_retained_evidence_is_sealable(&self) -> Result<(), CanwuError> {
345        if !self.state.scheduler.pending_ingress.is_empty() {
346            return Err(CanwuError::new(
347                ErrorCode::ArchiveNotReady,
348                "live evidence can be sealed only when the canonical ingress queue is empty",
349            ));
350        }
351        let reads_archived_evidence = |reads: &[super::StateKey]| {
352            reads.iter().any(|state| {
353                state == &super::StateKey::core_commands()
354                    || state == &super::StateKey::core_events()
355                    || state == &super::StateKey::core_ingress()
356            })
357        };
358        if self.plugins.descriptors().any(|descriptor| {
359            descriptor
360                .systems
361                .iter()
362                .any(|contract| reads_archived_evidence(&contract.reads))
363                || descriptor
364                    .boundary_systems
365                    .iter()
366                    .any(|contract| reads_archived_evidence(&contract.reads))
367                || descriptor
368                    .commands
369                    .iter()
370                    .any(|contract| reads_archived_evidence(&contract.reads))
371        }) {
372            return Err(CanwuError::new(
373                ErrorCode::ArchiveNotReady,
374                "live evidence sealing requires plugins whose declared reads use current state rather than historical command, event, or ingress records",
375            ));
376        }
377
378        let admitted_attempts: std::collections::BTreeSet<_> = self
379            .state
380            .evidence
381            .boundaries
382            .iter()
383            .flat_map(|record| record.admitted_attempts.iter().copied())
384            .collect();
385        if admitted_attempts.len() != self.state.evidence.command_attempts.len()
386            || self
387                .state
388                .evidence
389                .command_attempts
390                .iter()
391                .any(|attempt| !admitted_attempts.contains(&attempt.id))
392        {
393            return Err(CanwuError::new(
394                ErrorCode::ArchiveNotReady,
395                "live evidence sealing requires every retained command attempt to belong to a completed boundary",
396            ));
397        }
398        let admitted_commands: std::collections::BTreeSet<_> = self
399            .state
400            .evidence
401            .boundaries
402            .iter()
403            .flat_map(|record| record.admitted_commands.iter().copied())
404            .collect();
405        if admitted_commands.len() != self.state.evidence.commands.len()
406            || self
407                .state
408                .evidence
409                .commands
410                .iter()
411                .any(|command| !admitted_commands.contains(&command.id))
412        {
413            return Err(CanwuError::new(
414                ErrorCode::ArchiveNotReady,
415                "live evidence sealing requires every retained command to belong to a completed boundary",
416            ));
417        }
418        let admitted_ingress: std::collections::BTreeSet<_> = self
419            .state
420            .evidence
421            .boundaries
422            .iter()
423            .flat_map(|record| record.admitted_ingress.iter().copied())
424            .collect();
425        if admitted_ingress.len() != self.state.evidence.ingress.len()
426            || self
427                .state
428                .evidence
429                .ingress
430                .iter()
431                .any(|record| !admitted_ingress.contains(&record.id))
432        {
433            return Err(CanwuError::new(
434                ErrorCode::ArchiveNotReady,
435                "live evidence sealing requires every retained ingress record to belong to a completed boundary",
436            ));
437        }
438        let admitted_events: std::collections::BTreeSet<_> = self
439            .state
440            .evidence
441            .boundaries
442            .iter()
443            .flat_map(|record| record.admitted_events.iter().copied())
444            .collect();
445        if self.state.counters.admitted_event_count
446            != self
447                .state
448                .evidence
449                .archived
450                .event_count
451                .checked_add(
452                    u64::try_from(self.state.evidence.events.len()).map_err(|_| {
453                        CanwuError::new(
454                            ErrorCode::ArchiveNotReady,
455                            "retained event count exceeds the live archive cursor range",
456                        )
457                    })?,
458                )
459                .ok_or_else(|| {
460                    CanwuError::new(
461                        ErrorCode::ArchiveNotReady,
462                        "retained event cursor is exhausted",
463                    )
464                })?
465            || admitted_events.len() != self.state.evidence.events.len()
466            || self
467                .state
468                .evidence
469                .events
470                .iter()
471                .any(|event| !admitted_events.contains(&event.id))
472        {
473            return Err(CanwuError::new(
474                ErrorCode::ArchiveNotReady,
475                "live evidence sealing requires every retained event to be admitted by a later completed boundary",
476            ));
477        }
478        Ok(())
479    }
480
481    fn seal_retained_evidence(&mut self) -> Result<Option<EvidenceJournalSegment>, CanwuError> {
482        let start = self.state.evidence.archived;
483        let end = self.evidence_cursor()?;
484        if start == end {
485            return Ok(None);
486        }
487        self.ensure_retained_evidence_is_sealable()?;
488        let checkpoint_hash = self.state.metadata.checkpoint_hash.clone();
489        let commitment_roots = self.state.metadata.commitment_roots.clone();
490        let commitment_cache = self.state.metadata.commitment_cache.clone();
491        let prepared = (|| {
492            self.refresh_checkpoint_hash()?;
493
494            let mut archived_command_requests = Vec::new();
495            for attempt in &self.state.evidence.command_attempts {
496                let Some(request_id) = attempt.request_id else {
497                    continue;
498                };
499                let outcome = self.command_outcome_from_attempt(attempt)?;
500                archived_command_requests.push((
501                    request_id,
502                    ArchivedCommandRequestOutcome {
503                        input_hash: super::canonical_hash(
504                            "canwu.archive.command.request.v1",
505                            &(attempt.expected_revision, &attempt.envelope),
506                        )?,
507                        outcome,
508                    },
509                ));
510            }
511
512            let mut archived_ingress_requests = Vec::new();
513            for record in &self.state.evidence.ingress {
514                let IngressPayload::Command { request } = &record.payload else {
515                    continue;
516                };
517                archived_ingress_requests.push((
518                    request.request_id,
519                    ArchivedIngressRequest {
520                        input_hash: super::canonical_hash(
521                            "canwu.archive.ingress.command.v1",
522                            &(record.due_at, record.priority, request.as_ref()),
523                        )?,
524                        receipt: IngressReceipt {
525                            ingress_id: record.id,
526                            issued_at: record.issued_at,
527                            due_at: record.due_at,
528                        },
529                    },
530                ));
531            }
532            Ok::<_, CanwuError>((archived_command_requests, archived_ingress_requests))
533        })();
534        let (archived_command_requests, archived_ingress_requests) = match prepared {
535            Ok(prepared) => prepared,
536            Err(error) => {
537                self.state.metadata.checkpoint_hash = checkpoint_hash;
538                self.state.metadata.commitment_roots = commitment_roots;
539                self.state.metadata.commitment_cache = commitment_cache;
540                return Err(error);
541            }
542        };
543
544        self.state.evidence.archived_boundary_head = self
545            .state
546            .evidence
547            .boundaries
548            .last()
549            .map(|record| record.hash.clone())
550            .or_else(|| self.state.evidence.archived_boundary_head.clone());
551        self.state.evidence.archived_legacy_commands |= self
552            .state
553            .evidence
554            .commands
555            .iter()
556            .any(|record| record.attempt_id.is_none());
557        self.state.evidence.archived_tracked_attempts |=
558            !self.state.evidence.command_attempts.is_empty()
559                || !self.state.evidence.ingress.is_empty();
560        self.state.evidence.archived_unqueued_command_history |= has_unqueued_command_history(
561            &self.state.evidence.commands,
562            &self.state.evidence.command_attempts,
563            &self.state.evidence.ingress,
564        );
565        self.state
566            .evidence
567            .archived_command_requests
568            .extend(archived_command_requests);
569        self.state
570            .evidence
571            .archived_ingress_requests
572            .extend(archived_ingress_requests);
573        self.state.evidence.archived = end;
574        let segment = EvidenceJournalSegment {
575            format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
576            start,
577            end,
578            events: std::mem::take(&mut self.state.evidence.events),
579            commands: std::mem::take(&mut self.state.evidence.commands),
580            command_attempts: std::mem::take(&mut self.state.evidence.command_attempts),
581            ingress: std::mem::take(&mut self.state.evidence.ingress),
582            boundaries: std::mem::take(&mut self.state.evidence.boundaries),
583            random_draws: std::mem::take(&mut self.state.evidence.random_draws),
584        };
585        Ok(Some(segment))
586    }
587
588    pub(super) fn checkpoint_state(&self) -> SimulationSnapshot {
589        SimulationSnapshot {
590            engine_version: ENGINE_VERSION.to_owned(),
591            snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
592            run_manifest: Some(self.state.metadata.run_manifest.clone()),
593            run_manifest_hash: self.state.metadata.run_manifest_hash.clone(),
594            run_configuration: Some(self.state.metadata.run_configuration.clone()),
595            checkpoint_hash: self.state.metadata.checkpoint_hash.clone(),
596            commitment_format_version: self.state.metadata.commitment_format_version,
597            commitment_roots: self.state.metadata.commitment_roots.clone(),
598            revision_format_version: STATE_REVISION_FORMAT_VERSION,
599            state_revision: self.state.counters.state_revision,
600            replay_revision_format_version: self.state.metadata.replay_revision_format_version,
601            admission_cursor_format_version: ADMISSION_CURSOR_FORMAT_VERSION,
602            admitted_attempt_count: self.state.counters.admitted_attempt_count,
603            admitted_command_count: self.state.counters.admitted_command_count,
604            admitted_event_count: self.state.counters.admitted_event_count,
605            initial_time: self.state.scheduler.initial_time,
606            initial_scenario: self.bound_initial_scenario().cloned(),
607            now: self.state.scheduler.now,
608            plugin_registration_closed: self.state.metadata.plugin_registration_closed,
609            world: self.world(),
610            knowledge: self.state.current.knowledge.clone(),
611            events: Vec::new(),
612            commands: Vec::new(),
613            command_attempts: Vec::new(),
614            ingress: Vec::new(),
615            boundaries: Vec::new(),
616            plugin_components: self
617                .state
618                .current
619                .plugin_components
620                .values()
621                .cloned()
622                .collect(),
623            domain_records: self
624                .state
625                .current
626                .domain_records
627                .values()
628                .cloned()
629                .collect(),
630            plugin_descriptors: self.plugins.descriptors().cloned().collect(),
631            schema: self.schema.clone(),
632            root_seed: self.state.current.root_seed,
633            random_streams: self
634                .state
635                .current
636                .random_streams
637                .values()
638                .cloned()
639                .collect(),
640            random_draws: Vec::new(),
641            scheduled: self
642                .state
643                .scheduler
644                .actions
645                .iter()
646                .map(|(key, action)| ScheduledRecord {
647                    key: key.clone(),
648                    action: action.clone(),
649                })
650                .collect(),
651            legacy_rng: None,
652            next_event_id: self.state.counters.next_event_id,
653            next_command_id: self.state.counters.next_command_id,
654            next_command_attempt_id: self.state.counters.next_command_attempt_id,
655            next_ingress_id: self.state.counters.next_ingress_id,
656            next_boundary_id: self.state.counters.next_boundary_id,
657            next_random_draw_id: self.state.counters.next_random_draw_id,
658            next_schedule_sequence: self.state.counters.next_schedule_sequence,
659            next_correlation_id: self.state.counters.next_correlation_id,
660        }
661    }
662
663    /// Returns the current monotonic cut through every append-only journal.
664    pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
665        EvidenceCursor::from_evidence(&self.state.evidence)
666    }
667
668    /// Captures current authoritative state without cloning accumulated evidence.
669    pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
670        Ok(SimulationCheckpoint {
671            format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
672            journal_end: self.evidence_cursor()?,
673            state: self.checkpoint_state(),
674        })
675    }
676
677    /// Clones only evidence appended after a previously persisted cursor.
678    pub fn journal_segment_since(
679        &self,
680        start: EvidenceCursor,
681    ) -> Result<EvidenceJournalSegment, CanwuError> {
682        let end = self.evidence_cursor()?;
683        let cut = |value: u64, archived: u64, len: usize, label: &str| {
684            let value = value.checked_sub(archived).ok_or_else(|| {
685                CanwuError::new(
686                    ErrorCode::InvalidSnapshot,
687                    format!("{label} journal cursor precedes the retained live evidence window"),
688                )
689            })?;
690            let value = usize::try_from(value).map_err(|_| {
691                CanwuError::new(
692                    ErrorCode::InvalidSnapshot,
693                    format!("{label} journal cursor is not representable on this platform"),
694                )
695            })?;
696            if value > len {
697                return Err(CanwuError::new(
698                    ErrorCode::InvalidSnapshot,
699                    format!("{label} journal cursor exceeds the current evidence tail"),
700                ));
701            }
702            Ok(value)
703        };
704        let archived = self.state.evidence.archived;
705        let event_start = cut(
706            start.event_count,
707            archived.event_count,
708            self.state.evidence.events.len(),
709            "event",
710        )?;
711        let command_start = cut(
712            start.command_count,
713            archived.command_count,
714            self.state.evidence.commands.len(),
715            "command",
716        )?;
717        let attempt_start = cut(
718            start.command_attempt_count,
719            archived.command_attempt_count,
720            self.state.evidence.command_attempts.len(),
721            "command-attempt",
722        )?;
723        let ingress_start = cut(
724            start.ingress_count,
725            archived.ingress_count,
726            self.state.evidence.ingress.len(),
727            "ingress",
728        )?;
729        let boundary_start = cut(
730            start.boundary_count,
731            archived.boundary_count,
732            self.state.evidence.boundaries.len(),
733            "boundary",
734        )?;
735        let draw_start = cut(
736            start.random_draw_count,
737            archived.random_draw_count,
738            self.state.evidence.random_draws.len(),
739            "random-draw",
740        )?;
741        Ok(EvidenceJournalSegment {
742            format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
743            start,
744            end,
745            events: self.state.evidence.events[event_start..].to_vec(),
746            commands: self.state.evidence.commands[command_start..].to_vec(),
747            command_attempts: self.state.evidence.command_attempts[attempt_start..].to_vec(),
748            ingress: self.state.evidence.ingress[ingress_start..].to_vec(),
749            boundaries: self.state.evidence.boundaries[boundary_start..].to_vec(),
750            random_draws: self.state.evidence.random_draws[draw_start..].to_vec(),
751        })
752    }
753
754    /// Builds a portable full-save bundle with one segment from genesis.
755    pub fn checkpoint_journal(&self) -> Result<CheckpointJournal, CanwuError> {
756        if self.state.evidence.archived != EvidenceCursor::default() {
757            return Err(CanwuError::new(
758                ErrorCode::InvalidSnapshot,
759                "a compact live runtime requires its previously sealed evidence segments to build a portable save",
760            ));
761        }
762        let segment = self.journal_segment_since(EvidenceCursor::default())?;
763        Ok(CheckpointJournal {
764            checkpoint: self.checkpoint()?,
765            segments: (segment.start != segment.end)
766                .then_some(segment)
767                .into_iter()
768                .collect(),
769        })
770    }
771
772    /// Serializes the portable full-save checkpoint-journal bundle as JSON.
773    pub fn checkpoint_journal_json(&self) -> Result<String, CanwuError> {
774        serde_json::to_string_pretty(&self.checkpoint_journal()?).map_err(|error| {
775            CanwuError::new(
776                ErrorCode::InvalidSnapshot,
777                format!("could not serialize checkpoint journal: {error}"),
778            )
779        })
780    }
781
782    fn snapshot_from_checkpoint_and_journal(
783        checkpoint: SimulationCheckpoint,
784        segments: Vec<EvidenceJournalSegment>,
785    ) -> Result<SimulationSnapshot, CanwuError> {
786        if checkpoint.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
787            return Err(invalid_snapshot_error(format!(
788                "checkpoint-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
789                checkpoint.format_version
790            )));
791        }
792        let mut snapshot = checkpoint.state;
793        if snapshot.snapshot_format_version != SNAPSHOT_FORMAT_VERSION {
794            return Err(invalid_snapshot_error(format!(
795                "checkpoint-journal format {CHECKPOINT_JOURNAL_FORMAT_VERSION} requires snapshot format {SNAPSHOT_FORMAT_VERSION}"
796            )));
797        }
798        if !snapshot.events.is_empty()
799            || !snapshot.commands.is_empty()
800            || !snapshot.command_attempts.is_empty()
801            || !snapshot.ingress.is_empty()
802            || !snapshot.boundaries.is_empty()
803            || !snapshot.random_draws.is_empty()
804        {
805            return Err(invalid_snapshot_error(
806                "checkpoint current state must not duplicate append-only evidence",
807            ));
808        }
809
810        let mut cursor = EvidenceCursor::default();
811        for segment in segments {
812            if segment.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
813                return Err(invalid_snapshot_error(format!(
814                    "evidence-journal format {} is unsupported; this engine reads format {CHECKPOINT_JOURNAL_FORMAT_VERSION}",
815                    segment.format_version
816                )));
817            }
818            if segment.start != cursor {
819                return Err(invalid_snapshot_error(
820                    "evidence-journal segments must form one contiguous global prefix",
821                ));
822            }
823            let end = cursor.checked_advance(&segment)?;
824            if end == cursor {
825                return Err(invalid_snapshot_error(
826                    "evidence-journal segments must advance at least one journal cursor",
827                ));
828            }
829            if segment.end != end {
830                return Err(invalid_snapshot_error(
831                    "evidence-journal segment end does not match its encoded records",
832                ));
833            }
834            snapshot.events.extend(segment.events);
835            snapshot.commands.extend(segment.commands);
836            snapshot.command_attempts.extend(segment.command_attempts);
837            snapshot.ingress.extend(segment.ingress);
838            snapshot.boundaries.extend(segment.boundaries);
839            snapshot.random_draws.extend(segment.random_draws);
840            cursor = end;
841        }
842        if cursor != checkpoint.journal_end {
843            return Err(invalid_snapshot_error(
844                "evidence-journal segments do not reach the checkpoint journal cut",
845            ));
846        }
847        Ok(snapshot)
848    }
849
850    /// Restores a checkpoint after proving a contiguous journal prefix.
851    pub fn from_checkpoint_and_journal(
852        checkpoint: SimulationCheckpoint,
853        segments: Vec<EvidenceJournalSegment>,
854    ) -> Result<Self, CanwuError> {
855        Self::from_snapshot(Self::snapshot_from_checkpoint_and_journal(
856            checkpoint, segments,
857        )?)
858    }
859
860    /// Restores a portable checkpoint-journal bundle.
861    pub fn from_checkpoint_journal(bundle: CheckpointJournal) -> Result<Self, CanwuError> {
862        Self::from_checkpoint_and_journal(bundle.checkpoint, bundle.segments)
863    }
864
865    /// Restores a bundle and rehydrates its exact executable plugin contracts.
866    pub fn from_checkpoint_journal_with_plugins(
867        bundle: CheckpointJournal,
868        plugins: &[&dyn SimulationPlugin],
869    ) -> Result<Self, CanwuError> {
870        let mut simulation = Self::from_checkpoint_journal(bundle)?;
871        for plugin in plugins {
872            simulation.register_plugin(*plugin)?;
873        }
874        simulation.ensure_runtime_ready()?;
875        Ok(simulation)
876    }
877
878    /// Deserializes and restores a portable checkpoint-journal JSON bundle.
879    pub fn from_checkpoint_journal_json(json: &str) -> Result<Self, CanwuError> {
880        let bundle = serde_json::from_str(json).map_err(|error| {
881            CanwuError::new(
882                ErrorCode::InvalidSnapshot,
883                format!("could not deserialize checkpoint journal: {error}"),
884            )
885        })?;
886        Self::from_checkpoint_journal(bundle)
887    }
888
889    /// Deserializes a bundle and rehydrates its exact plugin contracts.
890    pub fn from_checkpoint_journal_json_with_plugins(
891        json: &str,
892        plugins: &[&dyn SimulationPlugin],
893    ) -> Result<Self, CanwuError> {
894        let bundle = serde_json::from_str(json).map_err(|error| {
895            CanwuError::new(
896                ErrorCode::InvalidSnapshot,
897                format!("could not deserialize checkpoint journal: {error}"),
898            )
899        })?;
900        Self::from_checkpoint_journal_with_plugins(bundle, plugins)
901    }
902}