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
14pub const CHECKPOINT_JOURNAL_FORMAT_VERSION: u32 = 1;
16
17#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
18pub 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)]
111pub 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)]
123pub 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)]
143pub struct CheckpointJournal {
145 pub checkpoint: SimulationCheckpoint,
146 pub segments: Vec<EvidenceJournalSegment>,
147}
148
149pub struct CompactedSimulation {
157 simulation: Simulation,
158}
159
160impl CompactedSimulation {
161 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
163 self.simulation.evidence_cursor()
164 }
165
166 pub fn checkpoint(&self) -> Result<SimulationCheckpoint, CanwuError> {
168 self.simulation.checkpoint()
169 }
170
171 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 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 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 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 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 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 pub fn evidence_cursor(&self) -> Result<EvidenceCursor, CanwuError> {
665 EvidenceCursor::from_evidence(&self.state.evidence)
666 }
667
668 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 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 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 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 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 pub fn from_checkpoint_journal(bundle: CheckpointJournal) -> Result<Self, CanwuError> {
862 Self::from_checkpoint_and_journal(bundle.checkpoint, bundle.segments)
863 }
864
865 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 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 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}