1use super::{
2 ADMISSION_CURSOR_FORMAT_VERSION, BoundaryRecord, BoundarySystemContract,
3 CHECKPOINT_JOURNAL_FORMAT_VERSION, COMMITMENT_FORMAT_VERSION, CanwuError, CheckpointJournal,
4 CommandAttemptRecord, CommandRecord, CommitmentRoots, DecisionState, DeterministicRng,
5 DomainRecord, DomainRecordClass, DomainRecordSchema, DomainReferenceSchema, ENGINE_VERSION,
6 ErrorCode, EvidenceCursor, EvidenceJournalSegment, IngressRecord, KnowledgeSnapshot,
7 PayloadSchema, PluginActionDescriptor, PluginComponentRecord, PluginDescriptor,
8 PluginIngressDescriptor, RandomDrawAddress, RandomDrawOutcome, RandomDrawProducer,
9 RandomDrawRecord, RandomStreamKey, RandomStreamState, ReservationRef, RunConfigurationSnapshot,
10 RunManifest, SNAPSHOT_FORMAT_VERSION, STATE_REVISION_FORMAT_VERSION, Scenario, ScheduledRecord,
11 SchemaRegistry, SimEvent, SimTime, Simulation, SimulationCheckpoint, SimulationSnapshot,
12 StateKey, StateVisibility, SystemCadence, SystemContract, WorldSnapshot,
13 boundary_state_hash_for_commitments, canonical_hash, checkpoint_hash_for_commitments,
14 commitment_roots_are_canonical, compute_boundary_hash, invalid_snapshot,
15 invalid_snapshot_error, is_canonical_hash, manifest, random_stream_commitment_root,
16 snapshot_checkpoint_hash, snapshot_commitment_roots,
17};
18use canwu_core::{DomainRecordKind, RandomDrawId};
19use canwu_event::{CauseRef, EventAudience};
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22use std::collections::BTreeMap;
23
24pub(super) const LEGACY_V4_ENGINE_VERSION: &str = "0.4.0";
25
26fn reject_unknown_fields(input: &Value, encoded: &Value, path: &str) -> Result<(), CanwuError> {
27 match (input, encoded) {
28 (Value::Object(input), Value::Object(encoded)) => {
29 for (key, value) in input {
30 let next = if path.is_empty() {
31 key.clone()
32 } else {
33 format!("{path}.{key}")
34 };
35 let Some(expected) = encoded.get(key) else {
36 return invalid_snapshot(format!(
37 "strict legacy wire contains unknown field `{next}`"
38 ));
39 };
40 reject_unknown_fields(value, expected, &next)?;
41 }
42 }
43 (Value::Array(input), Value::Array(encoded)) => {
44 if input.len() != encoded.len() {
45 return invalid_snapshot(format!(
46 "strict legacy wire array `{path}` changed shape during decoding"
47 ));
48 }
49 for (index, (value, expected)) in input.iter().zip(encoded).enumerate() {
50 reject_unknown_fields(value, expected, &format!("{path}[{index}]"))?;
51 }
52 }
53 _ => {}
54 }
55 Ok(())
56}
57
58fn deserialize_strict<T>(value: &Value, label: &str) -> Result<T, CanwuError>
59where
60 T: for<'de> Deserialize<'de> + Serialize,
61{
62 let decoded: T = serde_json::from_value(value.clone()).map_err(|error| {
63 invalid_snapshot_error(format!("could not deserialize strict {label}: {error}"))
64 })?;
65 let encoded = serde_json::to_value(&decoded).map_err(|error| {
66 invalid_snapshot_error(format!("could not re-encode strict {label}: {error}"))
67 })?;
68 reject_unknown_fields(value, &encoded, "")?;
69 Ok(decoded)
70}
71#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
72#[serde(deny_unknown_fields)]
73pub(super) struct LegacyV4RandomDrawRecord {
74 pub id: RandomDrawId,
75 pub at: SimTime,
76 pub stream: RandomStreamKey,
77 pub position: u64,
78 pub upper_exclusive: u64,
79 pub value: u64,
80 pub purpose: String,
81 pub producer: RandomDrawProducer,
82 #[serde(default)]
83 pub outcome: Option<RandomDrawOutcome>,
84 pub cause: CauseRef,
85 pub correlation_id: u64,
86}
87
88impl From<LegacyV4RandomDrawRecord> for RandomDrawRecord {
89 fn from(value: LegacyV4RandomDrawRecord) -> Self {
90 Self {
91 id: value.id,
92 at: value.at,
93 stream: value.stream,
94 address: RandomDrawAddress::Sequential {
95 position: value.position,
96 },
97 operation_evidence: None,
98 upper_exclusive: value.upper_exclusive,
99 value: value.value,
100 purpose: value.purpose,
101 producer: value.producer,
102 outcome: value.outcome,
103 cause: value.cause,
104 correlation_id: value.correlation_id,
105 }
106 }
107}
108
109#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
110#[serde(deny_unknown_fields)]
111pub(super) struct LegacyV4DomainRecordSchema {
112 kind: DomainRecordKind,
113 class: DomainRecordClass,
114 payload_schema: PayloadSchema,
115 references: Vec<DomainReferenceSchema>,
116}
117
118impl From<LegacyV4DomainRecordSchema> for DomainRecordSchema {
119 fn from(value: LegacyV4DomainRecordSchema) -> Self {
120 let mut current = Self::new(value.kind, value.class);
121 current.payload_schema = value.payload_schema;
122 current.references = value.references;
123 current
124 }
125}
126
127#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
128#[serde(deny_unknown_fields)]
129pub(super) struct LegacyV4BoundarySystemContract {
130 name: String,
131 phase: super::BoundaryPhase,
132 cadence: SystemCadence,
133 reads: Vec<StateKey>,
134 writes: Vec<StateKey>,
135 emits: Vec<String>,
136 reservation_offers: Vec<StateKey>,
137 reservation_requests: Vec<StateKey>,
138 reservation_reads: Vec<ReservationRef>,
139 #[serde(default)]
140 random_streams: Vec<RandomStreamKey>,
141 visibility: StateVisibility,
142}
143
144impl From<LegacyV4BoundarySystemContract> for BoundarySystemContract {
145 fn from(value: LegacyV4BoundarySystemContract) -> Self {
146 Self {
147 name: value.name,
148 phase: value.phase,
149 cadence: value.cadence,
150 reads: value.reads,
151 writes: value.writes,
152 emits: value.emits,
153 reservation_offers: value.reservation_offers,
154 reservation_requests: value.reservation_requests,
155 reservation_reads: value.reservation_reads,
156 random_streams: value.random_streams,
157 knowledge_writes: Vec::new(),
158 plugin_ingress_targets: Vec::new(),
159 visibility: value.visibility,
160 }
161 }
162}
163
164#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
165#[serde(deny_unknown_fields)]
166pub(super) struct LegacyV4PluginDescriptor {
167 name: String,
168 #[serde(default)]
169 version: String,
170 #[serde(default)]
171 semantic_hash: String,
172 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
173 event_audiences: BTreeMap<String, EventAudience>,
174 systems: Vec<SystemContract>,
175 #[serde(default)]
176 boundary_systems: Vec<LegacyV4BoundarySystemContract>,
177 commands: Vec<PluginActionDescriptor>,
178 #[serde(default, skip_serializing_if = "Vec::is_empty")]
179 ingress: Vec<PluginIngressDescriptor>,
180 schema_types: Vec<String>,
181 #[serde(default, skip_serializing_if = "Vec::is_empty")]
182 record_schemas: Vec<LegacyV4DomainRecordSchema>,
183}
184
185impl From<LegacyV4PluginDescriptor> for PluginDescriptor {
186 fn from(value: LegacyV4PluginDescriptor) -> Self {
187 Self {
188 name: value.name,
189 version: value.version,
190 semantic_hash: value.semantic_hash,
191 event_audiences: value.event_audiences,
192 systems: value.systems,
193 boundary_systems: value.boundary_systems.into_iter().map(Into::into).collect(),
194 commands: value.commands,
195 ingress: value.ingress,
196 schema_types: value.schema_types,
197 record_schemas: value.record_schemas.into_iter().map(Into::into).collect(),
198 knowledge_schemas: Vec::new(),
199 }
200 }
201}
202
203fn legacy_sorted_hash_by<T, K, F>(
204 domain: &str,
205 values: &[T],
206 mut key: F,
207) -> Result<String, CanwuError>
208where
209 T: Serialize,
210 K: Ord,
211 F: FnMut(&T) -> K,
212{
213 let mut ordered: Vec<_> = values.iter().collect();
214 ordered.sort_by_key(|value| key(value));
215 canonical_hash(domain, &ordered)
216}
217
218#[derive(Serialize)]
219struct LegacyV4IdentityCommitmentMaterial<'a> {
220 engine_version: &'a str,
221 snapshot_format_version: u32,
222 run_manifest: &'a RunManifest,
223 run_manifest_hash: &'a str,
224 initial_time: SimTime,
225 #[serde(skip_serializing_if = "Option::is_none")]
226 initial_scenario: Option<&'a Scenario>,
227 plugin_descriptors: String,
228 schema: &'a SchemaRegistry,
229}
230
231fn legacy_identity_commitment_root(
232 run_manifest: &RunManifest,
233 run_manifest_hash: &str,
234 initial_time: SimTime,
235 initial_scenario: Option<&Scenario>,
236 plugin_descriptors: &[LegacyV4PluginDescriptor],
237 schema: &SchemaRegistry,
238) -> Result<String, CanwuError> {
239 canonical_hash(
240 "canwu.commitment.identity.v1",
241 &LegacyV4IdentityCommitmentMaterial {
242 engine_version: LEGACY_V4_ENGINE_VERSION,
243 snapshot_format_version: 4,
244 run_manifest,
245 run_manifest_hash,
246 initial_time,
247 initial_scenario,
248 plugin_descriptors: legacy_sorted_hash_by(
249 "canwu.commitment.identity.plugins.v1",
250 plugin_descriptors,
251 |descriptor| descriptor.name.clone(),
252 )?,
253 schema,
254 },
255 )
256}
257
258#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
259#[serde(deny_unknown_fields)]
260pub(super) struct LegacyV4SimulationSnapshot {
261 pub engine_version: String,
262 pub snapshot_format_version: u32,
263 #[serde(default)]
264 pub run_manifest: Option<RunManifest>,
265 #[serde(default)]
266 pub run_manifest_hash: String,
267 #[serde(default, skip_serializing_if = "Option::is_none")]
268 pub run_configuration: Option<RunConfigurationSnapshot>,
269 #[serde(default)]
270 pub checkpoint_hash: String,
271 #[serde(default, skip_serializing_if = "super::is_zero_u32")]
272 pub commitment_format_version: u32,
273 #[serde(default, skip_serializing_if = "Option::is_none")]
274 pub commitment_roots: Option<CommitmentRoots>,
275 #[serde(default)]
276 pub revision_format_version: u32,
277 #[serde(default, skip_serializing_if = "super::is_zero_u64")]
278 pub state_revision: u64,
279 #[serde(default, skip_serializing_if = "super::is_zero_u32")]
280 pub replay_revision_format_version: u32,
281 #[serde(default, skip_serializing_if = "super::is_zero_u32")]
282 pub admission_cursor_format_version: u32,
283 #[serde(default, skip_serializing_if = "super::is_zero_u64")]
284 pub admitted_attempt_count: u64,
285 #[serde(default, skip_serializing_if = "super::is_zero_u64")]
286 pub admitted_command_count: u64,
287 #[serde(default, skip_serializing_if = "super::is_zero_u64")]
288 pub admitted_event_count: u64,
289 pub initial_time: SimTime,
290 #[serde(default, skip_serializing_if = "Option::is_none")]
291 pub initial_scenario: Option<Scenario>,
292 pub now: SimTime,
293 pub plugin_registration_closed: bool,
294 pub world: WorldSnapshot,
295 pub knowledge: KnowledgeSnapshot,
296 pub events: Vec<SimEvent>,
297 pub commands: Vec<CommandRecord>,
298 #[serde(default, skip_serializing_if = "Vec::is_empty")]
299 pub command_attempts: Vec<CommandAttemptRecord>,
300 #[serde(default, skip_serializing_if = "Vec::is_empty")]
301 pub ingress: Vec<IngressRecord>,
302 #[serde(default)]
303 pub boundaries: Vec<BoundaryRecord>,
304 pub plugin_components: Vec<PluginComponentRecord>,
305 #[serde(default, skip_serializing_if = "Vec::is_empty")]
306 pub domain_records: Vec<DomainRecord>,
307 #[serde(default, skip_serializing_if = "DecisionState::is_empty")]
308 pub decisions: DecisionState,
309 pub plugin_descriptors: Vec<LegacyV4PluginDescriptor>,
310 pub schema: SchemaRegistry,
311 #[serde(default)]
312 pub root_seed: u64,
313 #[serde(default)]
314 pub random_streams: Vec<RandomStreamState>,
315 #[serde(default)]
316 pub random_draws: Vec<LegacyV4RandomDrawRecord>,
317 pub scheduled: Vec<ScheduledRecord>,
318 #[serde(default, rename = "rng", skip_serializing_if = "Option::is_none")]
319 pub legacy_rng: Option<DeterministicRng>,
320 pub next_event_id: u64,
321 pub next_command_id: u64,
322 #[serde(default = "super::one_u64", skip_serializing_if = "super::is_one_u64")]
323 pub next_command_attempt_id: u64,
324 #[serde(default = "super::one_u64", skip_serializing_if = "super::is_one_u64")]
325 pub next_ingress_id: u64,
326 #[serde(default)]
327 pub next_boundary_id: u64,
328 #[serde(default)]
329 pub next_random_draw_id: u64,
330 pub next_schedule_sequence: u64,
331 pub next_correlation_id: u64,
332 #[serde(default = "super::one_u64", skip_serializing_if = "super::is_one_u64")]
333 pub next_decision_trace_id: u64,
334}
335
336impl LegacyV4SimulationSnapshot {
337 fn into_current(self) -> SimulationSnapshot {
338 SimulationSnapshot {
339 engine_version: self.engine_version,
340 snapshot_format_version: self.snapshot_format_version,
341 run_manifest: self.run_manifest,
342 run_manifest_hash: self.run_manifest_hash,
343 run_configuration: self.run_configuration,
344 checkpoint_hash: self.checkpoint_hash,
345 commitment_format_version: self.commitment_format_version,
346 commitment_roots: self.commitment_roots,
347 revision_format_version: self.revision_format_version,
348 state_revision: self.state_revision,
349 replay_revision_format_version: self.replay_revision_format_version,
350 admission_cursor_format_version: self.admission_cursor_format_version,
351 admitted_attempt_count: self.admitted_attempt_count,
352 admitted_command_count: self.admitted_command_count,
353 admitted_event_count: self.admitted_event_count,
354 initial_time: self.initial_time,
355 initial_scenario: self.initial_scenario,
356 now: self.now,
357 plugin_registration_closed: self.plugin_registration_closed,
358 world: self.world,
359 knowledge: self.knowledge,
360 events: self.events,
361 commands: self.commands,
362 command_attempts: self.command_attempts,
363 ingress: self.ingress,
364 boundaries: self.boundaries,
365 plugin_components: self.plugin_components,
366 domain_records: self.domain_records,
367 decisions: self.decisions,
368 plugin_descriptors: self
369 .plugin_descriptors
370 .into_iter()
371 .map(Into::into)
372 .collect(),
373 schema: self.schema,
374 root_seed: self.root_seed,
375 random_streams: self.random_streams,
376 random_draws: self.random_draws.into_iter().map(Into::into).collect(),
377 scheduled: self.scheduled,
378 legacy_rng: self.legacy_rng,
379 next_event_id: self.next_event_id,
380 next_command_id: self.next_command_id,
381 next_command_attempt_id: self.next_command_attempt_id,
382 next_ingress_id: self.next_ingress_id,
383 next_boundary_id: self.next_boundary_id,
384 next_random_draw_id: self.next_random_draw_id,
385 next_knowledge_record_id: 1,
386 next_schedule_sequence: self.next_schedule_sequence,
387 next_correlation_id: self.next_correlation_id,
388 next_decision_trace_id: self.next_decision_trace_id,
389 }
390 }
391}
392
393#[derive(Serialize)]
394struct LegacyRandomCommitmentMaterial {
395 root_seed: u64,
396 streams: String,
397 draws: String,
398}
399
400fn legacy_random_commitment_root(
401 root_seed: u64,
402 streams: &[RandomStreamState],
403 draws: &[LegacyV4RandomDrawRecord],
404) -> Result<String, CanwuError> {
405 let mut ordered: Vec<_> = draws.iter().collect();
406 ordered.sort_by_key(|draw| draw.id);
407 canonical_hash(
408 "canwu.commitment.random.v1",
409 &LegacyRandomCommitmentMaterial {
410 root_seed,
411 streams: random_stream_commitment_root(streams)?,
412 draws: canonical_hash("canwu.commitment.random.draws.v1", &ordered)?,
413 },
414 )
415}
416
417fn validate_boundary_chain(boundaries: &[BoundaryRecord]) -> Result<(), CanwuError> {
418 let mut previous = super::GENESIS_BOUNDARY_HASH;
419 for boundary in boundaries {
420 if boundary.previous_hash != previous || compute_boundary_hash(boundary)? != boundary.hash {
421 return invalid_snapshot("legacy format-4 boundary hash chain is inconsistent");
422 }
423 previous = &boundary.hash;
424 }
425 Ok(())
426}
427
428fn validate_legacy_commitments(
429 legacy: &LegacyV4SimulationSnapshot,
430 shadow: &SimulationSnapshot,
431) -> Result<(), CanwuError> {
432 if legacy.commitment_format_version != COMMITMENT_FORMAT_VERSION {
433 return invalid_snapshot("legacy format-4 snapshot must use commitment format 1");
434 }
435 let stored = legacy.commitment_roots.as_ref().ok_or_else(|| {
436 invalid_snapshot_error("legacy format-4 snapshot is missing commitment roots")
437 })?;
438 if !commitment_roots_are_canonical(stored) {
439 return invalid_snapshot("legacy format-4 commitment roots are not canonical");
440 }
441 let mut expected = snapshot_commitment_roots(shadow)?;
442 expected.random = legacy_random_commitment_root(
443 legacy.root_seed,
444 &legacy.random_streams,
445 &legacy.random_draws,
446 )?;
447 expected.identity = legacy_identity_commitment_root(
448 legacy.run_manifest.as_ref().ok_or_else(|| {
449 invalid_snapshot_error("legacy format-4 snapshot is missing its run manifest")
450 })?,
451 &legacy.run_manifest_hash,
452 legacy.initial_time,
453 legacy.initial_scenario.as_ref(),
454 &legacy.plugin_descriptors,
455 &legacy.schema,
456 )?;
457 if &expected != stored {
458 return invalid_snapshot(
459 "legacy format-4 commitment roots do not match the persisted state",
460 );
461 }
462 let expected_checkpoint = checkpoint_hash_for_commitments(
463 stored,
464 &legacy.run_manifest_hash,
465 legacy.commitment_format_version,
466 legacy.revision_format_version,
467 legacy.state_revision,
468 legacy.replay_revision_format_version,
469 )?;
470 if !is_canonical_hash(&legacy.checkpoint_hash) || expected_checkpoint != legacy.checkpoint_hash
471 {
472 return invalid_snapshot("legacy format-4 checkpoint hash is inconsistent");
473 }
474 if let Some(boundary) = legacy.boundaries.last()
475 && let Some(state_hash) = boundary.state_hash.as_deref()
476 {
477 let Some(hash) = state_hash.strip_prefix(super::BOUNDARY_STATE_HASH_V1_PREFIX) else {
478 return invalid_snapshot(
479 "legacy format-4 migration requires the final boundary state hash to use v1 commitments",
480 );
481 };
482 let mut boundary_roots = stored.clone();
483 boundary_roots.boundary_chain = canonical_hash(
484 "canwu.commitment.boundary-chain.v1",
485 boundary.previous_hash.as_str(),
486 )?;
487 if !is_canonical_hash(hash)
488 || boundary_state_hash_for_commitments(&boundary_roots)? != state_hash
489 {
490 return invalid_snapshot("legacy format-4 final boundary state hash is inconsistent");
491 }
492 }
493 Ok(())
494}
495
496fn validate_legacy_v4(
497 legacy: &LegacyV4SimulationSnapshot,
498) -> Result<SimulationSnapshot, CanwuError> {
499 if legacy.engine_version != LEGACY_V4_ENGINE_VERSION || legacy.snapshot_format_version != 4 {
500 return Err(CanwuError::new(
501 ErrorCode::UnsupportedSnapshotVersion,
502 "legacy migration accepts only engine 0.4.0 snapshot format 4",
503 ));
504 }
505 if legacy.legacy_rng.is_some() {
506 return invalid_snapshot("legacy format-4 snapshots cannot contain the pre-format-4 RNG");
507 }
508 if legacy.revision_format_version != STATE_REVISION_FORMAT_VERSION
509 || legacy.admission_cursor_format_version != ADMISSION_CURSOR_FORMAT_VERSION
510 || legacy.replay_revision_format_version > STATE_REVISION_FORMAT_VERSION
511 {
512 return invalid_snapshot("legacy format-4 revision or admission format is unsupported");
513 }
514 let mut shadow = legacy.clone().into_current();
515 if shadow.run_configuration.is_none() {
516 let manifest = shadow.run_manifest.as_ref().ok_or_else(|| {
517 invalid_snapshot_error("legacy format-4 snapshot is missing its run manifest")
518 })?;
519 shadow.run_configuration = Some(super::migration::inferred_run_configuration(manifest)?);
520 }
521 let run_manifest = shadow.run_manifest.as_ref().ok_or_else(|| {
522 invalid_snapshot_error("legacy format-4 snapshot is missing its run manifest")
523 })?;
524 manifest::validate(run_manifest, shadow.initial_scenario.as_ref(), true)?;
525 if manifest::hash(run_manifest)? != shadow.run_manifest_hash {
526 return invalid_snapshot("legacy format-4 run manifest hash is inconsistent");
527 }
528 validate_boundary_chain(&legacy.boundaries)?;
529 validate_legacy_commitments(legacy, &shadow)?;
530 Ok(shadow)
531}
532
533fn rebase_boundary_chain(snapshot: &mut SimulationSnapshot) -> Result<(), CanwuError> {
534 let mut previous = super::GENESIS_BOUNDARY_HASH.to_owned();
535 for boundary in &mut snapshot.boundaries {
536 boundary.previous_hash.clone_from(&previous);
537 boundary.hash = compute_boundary_hash(boundary)?;
538 previous.clone_from(&boundary.hash);
539 }
540 Ok(())
541}
542
543fn refresh_migrated_boundary_head_state_hash(
544 snapshot: &mut SimulationSnapshot,
545) -> Result<(), CanwuError> {
546 let Some(previous_hash) = snapshot
547 .boundaries
548 .last()
549 .map(|boundary| boundary.previous_hash.clone())
550 else {
551 return Ok(());
552 };
553 let mut roots = snapshot_commitment_roots(snapshot)?;
554 roots.boundary_chain =
555 canonical_hash("canwu.commitment.boundary-chain.v1", previous_hash.as_str())?;
556 let state_hash = boundary_state_hash_for_commitments(&roots)?;
557 let boundary = snapshot
558 .boundaries
559 .last_mut()
560 .expect("a captured boundary head must still exist");
561 boundary.state_hash = Some(state_hash);
562 boundary.hash = compute_boundary_hash(boundary)?;
563 Ok(())
564}
565
566pub(super) fn migrate_legacy_v4(
567 legacy: &LegacyV4SimulationSnapshot,
568) -> Result<SimulationSnapshot, CanwuError> {
569 let mut snapshot = validate_legacy_v4(legacy)?;
570 ENGINE_VERSION.clone_into(&mut snapshot.engine_version);
571 snapshot.snapshot_format_version = SNAPSHOT_FORMAT_VERSION;
572 snapshot.replay_revision_format_version = 0;
573 snapshot.commitment_roots = None;
574 snapshot.checkpoint_hash.clear();
575 rebase_boundary_chain(&mut snapshot)?;
576 refresh_migrated_boundary_head_state_hash(&mut snapshot)?;
577 snapshot.commitment_roots = Some(snapshot_commitment_roots(&snapshot)?);
578 snapshot.checkpoint_hash = snapshot_checkpoint_hash(&snapshot)?;
579 Ok(snapshot)
580}
581
582pub(super) fn deserialize_snapshot_json(json: &str) -> Result<SimulationSnapshot, CanwuError> {
583 let value: Value = serde_json::from_str(json).map_err(|error| {
584 invalid_snapshot_error(format!("could not deserialize snapshot envelope: {error}"))
585 })?;
586 let object = value
587 .as_object()
588 .ok_or_else(|| invalid_snapshot_error("snapshot envelope must be an object"))?;
589 let format = object
590 .get("snapshot_format_version")
591 .and_then(Value::as_u64)
592 .and_then(|value| u32::try_from(value).ok())
593 .ok_or_else(|| invalid_snapshot_error("snapshot format selector is missing or invalid"))?;
594 let engine = object
595 .get("engine_version")
596 .and_then(Value::as_str)
597 .ok_or_else(|| invalid_snapshot_error("snapshot engine selector is missing or invalid"))?;
598 match format {
599 SNAPSHOT_FORMAT_VERSION if engine == ENGINE_VERSION => {
600 deserialize_strict(&value, "format-5 snapshot")
601 }
602 4 if engine == LEGACY_V4_ENGINE_VERSION => {
603 let legacy: LegacyV4SimulationSnapshot =
604 deserialize_strict(&value, "legacy format-4 snapshot")?;
605 migrate_legacy_v4(&legacy)
606 }
607 _ => Err(CanwuError::new(
608 ErrorCode::UnsupportedSnapshotVersion,
609 format!(
610 "snapshot format {format} from engine {engine} is unsupported; this engine reads its own format {SNAPSHOT_FORMAT_VERSION} and engine {LEGACY_V4_ENGINE_VERSION} format 4"
611 ),
612 )),
613 }
614}
615
616#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
617#[serde(deny_unknown_fields)]
618struct LegacyV4SimulationCheckpoint {
619 format_version: u32,
620 journal_end: EvidenceCursor,
621 state: LegacyV4SimulationSnapshot,
622}
623
624#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
625#[serde(deny_unknown_fields)]
626struct LegacyV4EvidenceJournalSegment {
627 format_version: u32,
628 start: EvidenceCursor,
629 end: EvidenceCursor,
630 #[serde(default, skip_serializing_if = "Vec::is_empty")]
631 events: Vec<SimEvent>,
632 #[serde(default, skip_serializing_if = "Vec::is_empty")]
633 commands: Vec<CommandRecord>,
634 #[serde(default, skip_serializing_if = "Vec::is_empty")]
635 command_attempts: Vec<CommandAttemptRecord>,
636 #[serde(default, skip_serializing_if = "Vec::is_empty")]
637 ingress: Vec<IngressRecord>,
638 #[serde(default, skip_serializing_if = "Vec::is_empty")]
639 boundaries: Vec<BoundaryRecord>,
640 #[serde(default, skip_serializing_if = "Vec::is_empty")]
641 random_draws: Vec<LegacyV4RandomDrawRecord>,
642}
643
644#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
645#[serde(deny_unknown_fields)]
646struct LegacyV4CheckpointJournal {
647 checkpoint: LegacyV4SimulationCheckpoint,
648 segments: Vec<LegacyV4EvidenceJournalSegment>,
649}
650
651fn advance_cursor(
652 cursor: EvidenceCursor,
653 segment: &LegacyV4EvidenceJournalSegment,
654) -> Result<EvidenceCursor, CanwuError> {
655 let add = |value: u64, length: usize, label: &str| {
656 value
657 .checked_add(u64::try_from(length).map_err(|_| {
658 invalid_snapshot_error(format!("legacy {label} segment is too large"))
659 })?)
660 .ok_or_else(|| invalid_snapshot_error(format!("legacy {label} cursor overflow")))
661 };
662 Ok(EvidenceCursor {
663 event_count: add(cursor.event_count, segment.events.len(), "event")?,
664 command_count: add(cursor.command_count, segment.commands.len(), "command")?,
665 command_attempt_count: add(
666 cursor.command_attempt_count,
667 segment.command_attempts.len(),
668 "command-attempt",
669 )?,
670 ingress_count: add(cursor.ingress_count, segment.ingress.len(), "ingress")?,
671 boundary_count: add(cursor.boundary_count, segment.boundaries.len(), "boundary")?,
672 random_draw_count: add(
673 cursor.random_draw_count,
674 segment.random_draws.len(),
675 "random-draw",
676 )?,
677 })
678}
679
680fn clear_snapshot_evidence(snapshot: &mut SimulationSnapshot) {
681 snapshot.events.clear();
682 snapshot.commands.clear();
683 snapshot.command_attempts.clear();
684 snapshot.ingress.clear();
685 snapshot.boundaries.clear();
686 snapshot.random_draws.clear();
687}
688
689fn migrate_legacy_checkpoint_journal(
690 legacy: LegacyV4CheckpointJournal,
691) -> Result<CheckpointJournal, CanwuError> {
692 if legacy.checkpoint.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION {
693 return invalid_snapshot("legacy checkpoint-journal format is unsupported");
694 }
695 let mut full = legacy.checkpoint.state.clone();
696 if !full.events.is_empty()
697 || !full.commands.is_empty()
698 || !full.command_attempts.is_empty()
699 || !full.ingress.is_empty()
700 || !full.boundaries.is_empty()
701 || !full.random_draws.is_empty()
702 {
703 return invalid_snapshot("legacy checkpoint state duplicates append-only evidence");
704 }
705 let mut cursor = EvidenceCursor::default();
706 for segment in &legacy.segments {
707 if segment.format_version != CHECKPOINT_JOURNAL_FORMAT_VERSION || segment.start != cursor {
708 return invalid_snapshot("legacy checkpoint-journal segments are not contiguous");
709 }
710 let end = advance_cursor(cursor, segment)?;
711 if end == cursor || segment.end != end {
712 return invalid_snapshot("legacy checkpoint-journal segment end is invalid");
713 }
714 full.events.extend(segment.events.iter().cloned());
715 full.commands.extend(segment.commands.iter().cloned());
716 full.command_attempts
717 .extend(segment.command_attempts.iter().cloned());
718 full.ingress.extend(segment.ingress.iter().cloned());
719 full.boundaries.extend(segment.boundaries.iter().cloned());
720 full.random_draws
721 .extend(segment.random_draws.iter().cloned());
722 cursor = end;
723 }
724 if cursor != legacy.checkpoint.journal_end {
725 return invalid_snapshot("legacy checkpoint-journal does not reach its declared cut");
726 }
727
728 let migrated = migrate_legacy_v4(&full)?;
729 let continuation_checkpoint = Simulation::from_snapshot(migrated.clone())?.checkpoint()?;
730 let mut checkpoint_state = migrated.clone();
731 clear_snapshot_evidence(&mut checkpoint_state);
732 let mut segments = Vec::with_capacity(legacy.segments.len());
733 let mut event_at = 0usize;
734 let mut command_at = 0usize;
735 let mut attempt_at = 0usize;
736 let mut ingress_at = 0usize;
737 let mut boundary_at = 0usize;
738 let mut draw_at = 0usize;
739 for segment in legacy.segments {
740 let event_end = event_at + segment.events.len();
741 let command_end = command_at + segment.commands.len();
742 let attempt_end = attempt_at + segment.command_attempts.len();
743 let ingress_end = ingress_at + segment.ingress.len();
744 let boundary_end = boundary_at + segment.boundaries.len();
745 let draw_end = draw_at + segment.random_draws.len();
746 segments.push(EvidenceJournalSegment {
747 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
748 start: segment.start,
749 end: segment.end,
750 events: migrated.events[event_at..event_end].to_vec(),
751 commands: migrated.commands[command_at..command_end].to_vec(),
752 command_attempts: migrated.command_attempts[attempt_at..attempt_end].to_vec(),
753 ingress: migrated.ingress[ingress_at..ingress_end].to_vec(),
754 boundaries: migrated.boundaries[boundary_at..boundary_end].to_vec(),
755 random_draws: migrated.random_draws[draw_at..draw_end].to_vec(),
756 archive: None,
757 });
758 event_at = event_end;
759 command_at = command_end;
760 attempt_at = attempt_end;
761 ingress_at = ingress_end;
762 boundary_at = boundary_end;
763 draw_at = draw_end;
764 }
765 Ok(CheckpointJournal {
766 checkpoint: SimulationCheckpoint {
767 format_version: CHECKPOINT_JOURNAL_FORMAT_VERSION,
768 journal_end: legacy.checkpoint.journal_end,
769 state: checkpoint_state,
770 archived_segment_headers: Vec::new(),
771 archived_segment_manifest_root: None,
772 archived_evidence_receipts: Vec::new(),
773 archived_receipt_root: None,
774 evidence_dependencies: continuation_checkpoint.evidence_dependencies,
775 evidence_dependency_root: continuation_checkpoint.evidence_dependency_root,
776 keyed_draw_reservations: Vec::new(),
777 keyed_reservation_root: None,
778 },
779 segments,
780 })
781}
782
783pub(super) fn deserialize_checkpoint_journal_json(
784 json: &str,
785) -> Result<CheckpointJournal, CanwuError> {
786 let value: Value = serde_json::from_str(json).map_err(|error| {
787 invalid_snapshot_error(format!("could not deserialize checkpoint journal: {error}"))
788 })?;
789 let state = value
790 .get("checkpoint")
791 .and_then(|value| value.get("state"))
792 .and_then(Value::as_object)
793 .ok_or_else(|| invalid_snapshot_error("checkpoint journal state selector is missing"))?;
794 let format = state
795 .get("snapshot_format_version")
796 .and_then(Value::as_u64)
797 .and_then(|value| u32::try_from(value).ok())
798 .ok_or_else(|| invalid_snapshot_error("checkpoint journal format selector is invalid"))?;
799 let engine = state
800 .get("engine_version")
801 .and_then(Value::as_str)
802 .ok_or_else(|| invalid_snapshot_error("checkpoint journal engine selector is invalid"))?;
803 match format {
804 SNAPSHOT_FORMAT_VERSION if engine == ENGINE_VERSION => {
805 deserialize_strict(&value, "format-5 checkpoint journal")
806 }
807 4 if engine == LEGACY_V4_ENGINE_VERSION => {
808 let legacy: LegacyV4CheckpointJournal =
809 deserialize_strict(&value, "legacy format-4 checkpoint journal")?;
810 migrate_legacy_checkpoint_journal(legacy)
811 }
812 _ => Err(CanwuError::new(
813 ErrorCode::UnsupportedSnapshotVersion,
814 "checkpoint journal engine or snapshot format is unsupported",
815 )),
816 }
817}
818
819#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
820#[serde(deny_unknown_fields)]
821pub(super) struct LegacyV4ReplayJournalWire {
822 pub engine_version: String,
823 pub snapshot_format_version: u32,
824 pub root_seed: u64,
825 pub run_manifest: RunManifest,
826 pub run_manifest_hash: String,
827 #[serde(default)]
828 pub run_configuration: Option<RunConfigurationSnapshot>,
829 pub plugin_descriptors: Vec<LegacyV4PluginDescriptor>,
830 pub plugin_registration_closed: bool,
831 pub commands: Vec<CommandRecord>,
832 #[serde(default)]
833 pub command_attempts: Vec<CommandAttemptRecord>,
834 #[serde(default)]
835 pub ingress: Vec<IngressRecord>,
836 pub boundaries: Vec<BoundaryRecord>,
837 pub final_time: SimTime,
838 pub checkpoint_hash: String,
839 #[serde(default)]
840 pub commitment_format_version: u32,
841 #[serde(default)]
842 pub revision_format_version: u32,
843 #[serde(default)]
844 pub final_revision: u64,
845}
846
847pub(super) fn validate_legacy_replay_wire(wire: &LegacyV4ReplayJournalWire) -> Result<(), String> {
848 if wire.engine_version != LEGACY_V4_ENGINE_VERSION || wire.snapshot_format_version != 4 {
849 return Err("legacy replay accepts only engine 0.4.0 format 4".to_owned());
850 }
851 if manifest::hash(&wire.run_manifest).map_err(|error| error.to_string())?
852 != wire.run_manifest_hash
853 {
854 return Err("legacy replay manifest hash is inconsistent".to_owned());
855 }
856 let expected_revision = super::authoritative_revision_count(
857 wire.commands.len(),
858 wire.command_attempts.len(),
859 wire.boundaries.len(),
860 )
861 .map_err(|error| error.to_string())?;
862 if wire.final_revision != expected_revision {
863 return Err("legacy replay final revision is inconsistent".to_owned());
864 }
865 validate_boundary_chain(&wire.boundaries).map_err(|error| error.to_string())
866}
867pub(super) fn deserialize_replay_value(value: &Value) -> Result<super::ReplayJournal, String> {
868 let object = value
869 .as_object()
870 .ok_or_else(|| "replay journal envelope must be an object".to_owned())?;
871 let format = object
872 .get("snapshot_format_version")
873 .and_then(Value::as_u64)
874 .and_then(|value| u32::try_from(value).ok())
875 .ok_or_else(|| "replay journal format selector is missing or invalid".to_owned())?;
876 let engine = object
877 .get("engine_version")
878 .and_then(Value::as_str)
879 .ok_or_else(|| "replay journal engine selector is missing or invalid".to_owned())?;
880 match format {
881 SNAPSHOT_FORMAT_VERSION if engine == ENGINE_VERSION => {
882 let wire: super::persistence::ReplayJournalWire =
883 deserialize_strict(value, "format-5 replay journal")
884 .map_err(|error| error.to_string())?;
885 let run_configuration = wire
886 .run_configuration
887 .map_or_else(|| super::inferred_run_configuration(&wire.run_manifest), Ok)
888 .map_err(|error| error.to_string())?;
889 Ok(super::ReplayJournal {
890 engine_version: wire.engine_version,
891 snapshot_format_version: wire.snapshot_format_version,
892 root_seed: wire.root_seed,
893 run_manifest: wire.run_manifest,
894 run_manifest_hash: wire.run_manifest_hash,
895 run_configuration,
896 plugin_descriptors: wire.plugin_descriptors,
897 plugin_registration_closed: wire.plugin_registration_closed,
898 commands: wire.commands,
899 command_attempts: wire.command_attempts,
900 ingress: wire.ingress,
901 boundaries: wire.boundaries,
902 final_time: wire.final_time,
903 checkpoint_hash: wire.checkpoint_hash,
904 commitment_format_version: wire.commitment_format_version,
905 revision_format_version: wire.revision_format_version,
906 final_revision: wire.final_revision,
907 })
908 }
909 4 if engine == LEGACY_V4_ENGINE_VERSION => {
910 let wire: LegacyV4ReplayJournalWire =
911 deserialize_strict(value, "legacy format-4 replay journal")
912 .map_err(|error| error.to_string())?;
913 validate_legacy_replay_wire(&wire)?;
914 let run_configuration = wire
915 .run_configuration
916 .clone()
917 .map_or_else(|| super::inferred_run_configuration(&wire.run_manifest), Ok)
918 .map_err(|error| error.to_string())?;
919 Ok(super::ReplayJournal {
920 engine_version: wire.engine_version,
921 snapshot_format_version: wire.snapshot_format_version,
922 root_seed: wire.root_seed,
923 run_manifest: wire.run_manifest,
924 run_manifest_hash: wire.run_manifest_hash,
925 run_configuration,
926 plugin_descriptors: wire
927 .plugin_descriptors
928 .into_iter()
929 .map(Into::into)
930 .collect(),
931 plugin_registration_closed: wire.plugin_registration_closed,
932 commands: wire.commands,
933 command_attempts: wire.command_attempts,
934 ingress: wire.ingress,
935 boundaries: wire.boundaries,
936 final_time: wire.final_time,
937 checkpoint_hash: wire.checkpoint_hash,
938 commitment_format_version: wire.commitment_format_version,
939 revision_format_version: 0,
940 final_revision: wire.final_revision,
941 })
942 }
943 _ => Err(format!(
944 "replay journal format {format} from engine {engine} is unsupported"
945 )),
946 }
947}
948#[cfg(test)]
949mod tests {
950 use super::super::{Simulation, demo_scenario};
951 use super::*;
952
953 fn empty_legacy_value() -> Value {
954 let (scenario, _) = demo_scenario();
955 let simulation = Simulation::new(401, scenario).expect("fixture simulation should build");
956 let mut value =
957 serde_json::to_value(simulation.snapshot()).expect("current snapshot should serialize");
958 let object = value.as_object_mut().expect("snapshot should be an object");
959 object.insert(
960 "engine_version".to_owned(),
961 Value::String(LEGACY_V4_ENGINE_VERSION.to_owned()),
962 );
963 object.insert("snapshot_format_version".to_owned(), Value::from(4));
964 let mut legacy: LegacyV4SimulationSnapshot =
965 serde_json::from_value(value).expect("empty current wire should fit legacy V4");
966 let shadow = legacy.clone().into_current();
967 let mut roots = snapshot_commitment_roots(&shadow).expect("roots should compute");
968 roots.random = legacy_random_commitment_root(
969 legacy.root_seed,
970 &legacy.random_streams,
971 &legacy.random_draws,
972 )
973 .expect("legacy random root should compute");
974 roots.identity = legacy_identity_commitment_root(
975 legacy.run_manifest.as_ref().expect("manifest should exist"),
976 &legacy.run_manifest_hash,
977 legacy.initial_time,
978 legacy.initial_scenario.as_ref(),
979 &legacy.plugin_descriptors,
980 &legacy.schema,
981 )
982 .expect("legacy identity root should compute");
983 legacy.commitment_roots = Some(roots.clone());
984 legacy.checkpoint_hash = checkpoint_hash_for_commitments(
985 &roots,
986 &legacy.run_manifest_hash,
987 legacy.commitment_format_version,
988 legacy.revision_format_version,
989 legacy.state_revision,
990 legacy.replay_revision_format_version,
991 )
992 .expect("legacy checkpoint should compute");
993 serde_json::to_value(legacy).expect("legacy snapshot should serialize")
994 }
995
996 #[test]
997 fn strict_v4_snapshot_validates_before_migration_and_rejects_unknown_nested_fields() {
998 let value = empty_legacy_value();
999 let json = serde_json::to_string(&value).expect("fixture should serialize");
1000 let migrated = deserialize_snapshot_json(&json).expect("valid V4 should migrate");
1001 assert_eq!(migrated.engine_version, ENGINE_VERSION);
1002 assert_eq!(migrated.snapshot_format_version, SNAPSHOT_FORMAT_VERSION);
1003 Simulation::from_snapshot(migrated).expect("migrated snapshot should become live");
1004
1005 let mut tampered = value;
1006 tampered["world"]["format_5_only"] = Value::Bool(true);
1007 let error = deserialize_snapshot_json(
1008 &serde_json::to_string(&tampered).expect("tamper should serialize"),
1009 )
1010 .expect_err("unknown nested fields must fail before migration");
1011 assert_eq!(error.code, ErrorCode::InvalidSnapshot);
1012 assert!(error.message.contains("world.format_5_only"));
1013 }
1014
1015 #[test]
1016 fn legacy_replay_and_checkpoint_journal_are_validated_then_marked_historical() {
1017 let snapshot = empty_legacy_value();
1018 let snapshot_object = snapshot.as_object().expect("snapshot object");
1019 let replay = serde_json::json!({
1020 "engine_version": LEGACY_V4_ENGINE_VERSION,
1021 "snapshot_format_version": 4,
1022 "root_seed": snapshot_object["root_seed"],
1023 "run_manifest": snapshot_object["run_manifest"],
1024 "run_manifest_hash": snapshot_object["run_manifest_hash"],
1025 "run_configuration": snapshot_object["run_configuration"],
1026 "plugin_descriptors": snapshot_object["plugin_descriptors"],
1027 "plugin_registration_closed": snapshot_object["plugin_registration_closed"],
1028 "commands": [],
1029 "command_attempts": [],
1030 "ingress": [],
1031 "boundaries": [],
1032 "final_time": snapshot_object["now"],
1033 "checkpoint_hash": snapshot_object["checkpoint_hash"],
1034 "commitment_format_version": snapshot_object["commitment_format_version"],
1035 "revision_format_version": snapshot_object["revision_format_version"],
1036 "final_revision": 0
1037 });
1038 let journal: crate::ReplayJournal = serde_json::from_value(replay)
1039 .expect("valid legacy replay envelope should deserialize");
1040 assert_eq!(journal.revision_format_version, 0);
1041
1042 let bundle = serde_json::json!({
1043 "checkpoint": {
1044 "format_version": CHECKPOINT_JOURNAL_FORMAT_VERSION,
1045 "journal_end": EvidenceCursor::default(),
1046 "state": snapshot
1047 },
1048 "segments": []
1049 });
1050 let migrated = deserialize_checkpoint_journal_json(
1051 &serde_json::to_string(&bundle).expect("bundle should serialize"),
1052 )
1053 .expect("valid legacy checkpoint journal should migrate");
1054 assert_eq!(
1055 migrated.checkpoint.state.snapshot_format_version,
1056 SNAPSHOT_FORMAT_VERSION
1057 );
1058 }
1059}