Skip to main content

canwu_sim/runtime/
legacy_v4.rs

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}