Skip to main content

canwu_sim/runtime/
replay.rs

1use super::{
2    BOUNDARY_STATE_HASH_V1_PREFIX, BoundaryRecord, BoundaryRequest, COMMITMENT_FORMAT_VERSION,
3    CanwuError, CauseRef, CommandAttemptOutcome, CommandAttemptRecord, CommandIngress,
4    CommandOutcome, CommandRecord, ENGINE_VERSION, ErrorCode, IngressPayload, IngressRecord,
5    PluginIngressRequest, PluginRegistry, ReplayJournal, RunConfiguration, RunManifest,
6    SNAPSHOT_FORMAT_VERSION, STATE_REVISION_FORMAT_VERSION, Scenario, SimDuration, SimTime,
7    Simulation, SimulationPlugin, authoritative_revision_count, boundary_state_hash_format,
8    is_canonical_hash, manifest,
9};
10use std::collections::BTreeSet;
11
12impl Simulation {
13    /// Reconstructs caller-supplied core commands without proving a recorded
14    /// package environment. Use [`Self::replay_from_journal`] for exact replay.
15    pub fn replay(
16        seed: u64,
17        scenario: Scenario,
18        commands: &[CommandRecord],
19        final_time: SimTime,
20    ) -> Result<Self, CanwuError> {
21        Self::replay_with_plugins(seed, scenario, &[], commands, final_time)
22    }
23
24    /// Reconstructs caller-supplied inputs under caller-supplied plugins.
25    /// This is not an exact replay identity check.
26    pub fn replay_with_plugins(
27        seed: u64,
28        scenario: Scenario,
29        plugins: &[&dyn SimulationPlugin],
30        commands: &[CommandRecord],
31        final_time: SimTime,
32    ) -> Result<Self, CanwuError> {
33        Self::replay_with_boundaries(seed, scenario, plugins, commands, &[], final_time)
34    }
35
36    /// Reconstructs caller-supplied inputs and compares supplied boundaries.
37    /// Use [`Self::replay_from_journal`] when command-only runs must also bind
38    /// their recorded run and plugin identities.
39    pub fn replay_with_boundaries(
40        seed: u64,
41        scenario: Scenario,
42        plugins: &[&dyn SimulationPlugin],
43        commands: &[CommandRecord],
44        boundaries: &[BoundaryRecord],
45        final_time: SimTime,
46    ) -> Result<Self, CanwuError> {
47        let run_manifest = RunManifest::for_scenario("canwu.inline", "scenario", "1", &scenario)?;
48        Self::replay_with_run_manifest(
49            seed,
50            scenario,
51            run_manifest,
52            plugins,
53            commands,
54            boundaries,
55            final_time,
56        )
57    }
58
59    /// Reconstructs caller-supplied inputs under a caller-supplied run manifest.
60    /// This is useful for fixtures; it does not establish recorded identity.
61    pub fn replay_with_run_manifest(
62        seed: u64,
63        scenario: Scenario,
64        run_manifest: RunManifest,
65        plugins: &[&dyn SimulationPlugin],
66        commands: &[CommandRecord],
67        boundaries: &[BoundaryRecord],
68        final_time: SimTime,
69    ) -> Result<Self, CanwuError> {
70        let simulation =
71            Self::new_with_manifest_and_plugins(seed, scenario, run_manifest, plugins)?;
72        Self::replay_records(simulation, commands, &[], &[], boundaries, final_time)
73    }
74
75    /// Reconstructs caller-supplied inputs under a caller-supplied declared run
76    /// configuration. This does not establish recorded-environment identity.
77    #[allow(clippy::too_many_arguments)]
78    pub fn replay_with_run_configuration(
79        seed: u64,
80        scenario: Scenario,
81        run_manifest: RunManifest,
82        run_configuration: RunConfiguration,
83        plugins: &[&dyn SimulationPlugin],
84        commands: &[CommandRecord],
85        command_attempts: &[CommandAttemptRecord],
86        boundaries: &[BoundaryRecord],
87        final_time: SimTime,
88    ) -> Result<Self, CanwuError> {
89        if command_attempts
90            .iter()
91            .any(|attempt| attempt.ingress == CommandIngress::FrozenReplay)
92        {
93            return Err(CanwuError::new(
94                ErrorCode::ReplayEnvironmentMismatch,
95                "frozen replay attempts require an environment-bound replay journal",
96            ));
97        }
98        let simulation = Self::new_with_run_configuration_and_plugins(
99            seed,
100            scenario,
101            run_manifest,
102            run_configuration,
103            plugins,
104        )?;
105        Self::replay_records(
106            simulation,
107            commands,
108            command_attempts,
109            &[],
110            boundaries,
111            final_time,
112        )
113    }
114
115    /// Replays only after the recorded engine, run, seed, and plugin manifests
116    /// match, then verifies the final checkpoint commitment.
117    pub fn replay_from_journal(
118        scenario: Scenario,
119        plugins: &[&dyn SimulationPlugin],
120        journal: &ReplayJournal,
121    ) -> Result<Self, CanwuError> {
122        if !matches!(
123            journal.commitment_format_version,
124            0 | COMMITMENT_FORMAT_VERSION
125        ) {
126            return Err(CanwuError::new(
127                ErrorCode::ReplayEnvironmentMismatch,
128                format!(
129                    "replay journal commitment format {} is unsupported; this engine reads legacy format 0 and current format {COMMITMENT_FORMAT_VERSION}",
130                    journal.commitment_format_version
131                ),
132            ));
133        }
134        if journal.commitment_format_version == 0
135            && journal.boundaries.iter().any(|boundary| {
136                boundary
137                    .state_hash
138                    .as_deref()
139                    .is_some_and(|hash| hash.starts_with(BOUNDARY_STATE_HASH_V1_PREFIX))
140            })
141        {
142            return Err(CanwuError::new(
143                ErrorCode::ReplayEnvironmentMismatch,
144                "legacy commitment journals cannot contain boundary state commitment v1",
145            ));
146        }
147        if journal.revision_format_version != STATE_REVISION_FORMAT_VERSION {
148            return Err(CanwuError::new(
149                ErrorCode::LegacyReplayUnavailable,
150                "legacy revision histories can continue after snapshot migration but cannot claim exact replay because historical state commitments cannot be reconstructed without their original runtime",
151            ));
152        }
153        let expected_final_revision = authoritative_revision_count(
154            journal.commands.len(),
155            journal.command_attempts.len(),
156            journal.boundaries.len(),
157        )?;
158        if journal.final_revision != expected_final_revision {
159            return Err(CanwuError::new(
160                ErrorCode::ReplayEnvironmentMismatch,
161                "replay journal final revision is inconsistent with its committed evidence",
162            ));
163        }
164        let normalized = journal;
165        if matches!(journal.run_manifest, RunManifest::MigratedLegacy { .. }) {
166            return Err(CanwuError::new(
167                ErrorCode::LegacyReplayUnavailable,
168                "legacy checkpoints can continue after migration but lack enough recorded identity for exact replay",
169            ));
170        }
171        manifest::validate(&normalized.run_manifest, Some(&scenario), false)?;
172        manifest::validate_run_configuration(
173            &normalized.run_manifest,
174            &normalized.run_configuration,
175        )?;
176        if normalized.engine_version != ENGINE_VERSION
177            || normalized.snapshot_format_version != SNAPSHOT_FORMAT_VERSION
178            || !is_canonical_hash(&normalized.run_manifest_hash)
179            || manifest::hash(&normalized.run_manifest)? != normalized.run_manifest_hash
180            || !is_canonical_hash(&normalized.checkpoint_hash)
181        {
182            return Err(CanwuError::new(
183                ErrorCode::ReplayEnvironmentMismatch,
184                "replay journal engine, format, or run identity does not match this runtime",
185            ));
186        }
187        PluginRegistry::from_descriptors(normalized.plugin_descriptors.clone()).map_err(
188            |error| {
189                CanwuError::new(
190                    ErrorCode::ReplayEnvironmentMismatch,
191                    format!("replay journal plugin manifest is invalid: {error}"),
192                )
193            },
194        )?;
195
196        let mut simulation = Self::new_with_configuration_snapshot(
197            normalized.root_seed,
198            scenario,
199            normalized.run_manifest.clone(),
200            normalized.run_configuration.clone(),
201        )?;
202        if normalized.commitment_format_version == 0 {
203            simulation.state.metadata.commitment_format_version = 0;
204            simulation.state.metadata.commitment_roots = None;
205            simulation.state.metadata.commitment_cache = None;
206            simulation.refresh_checkpoint_hash()?;
207        }
208        let simulation = Self::activate_initial_plugins(simulation, plugins)?;
209        let actual_descriptors: Vec<_> = simulation.plugin_descriptors().cloned().collect();
210        if actual_descriptors != normalized.plugin_descriptors {
211            return Err(CanwuError::new(
212                ErrorCode::ReplayEnvironmentMismatch,
213                "active plugin identities and contracts do not match the replay journal",
214            ));
215        }
216
217        let mut simulation = Self::replay_records(
218            simulation,
219            &normalized.commands,
220            &normalized.command_attempts,
221            &normalized.ingress,
222            &normalized.boundaries,
223            normalized.final_time,
224        )?;
225        if normalized.plugin_registration_closed
226            && !simulation.state.metadata.plugin_registration_closed
227        {
228            simulation.advance(SimDuration::ZERO)?;
229        }
230        if simulation.state.metadata.plugin_registration_closed
231            != normalized.plugin_registration_closed
232        {
233            return Err(CanwuError::new(
234                ErrorCode::ReplayMismatch,
235                "replayed plugin-registration lifecycle does not match the recorded journal",
236            ));
237        }
238        if simulation.revision() != normalized.final_revision {
239            return Err(CanwuError::new(
240                ErrorCode::ReplayMismatch,
241                "replayed final state revision does not match the recorded journal",
242            ));
243        }
244        if simulation.checkpoint_hash() != normalized.checkpoint_hash {
245            return Err(CanwuError::new(
246                ErrorCode::ReplayMismatch,
247                "replayed final checkpoint does not match the recorded journal",
248            ));
249        }
250        Ok(simulation)
251    }
252
253    fn replay_records(
254        mut simulation: Self,
255        commands: &[CommandRecord],
256        attempts: &[CommandAttemptRecord],
257        ingress: &[IngressRecord],
258        boundaries: &[BoundaryRecord],
259        final_time: SimTime,
260    ) -> Result<Self, CanwuError> {
261        simulation.ensure_runtime_ready()?;
262        if !attempts.is_empty() {
263            return Self::replay_attempt_records(
264                simulation, commands, attempts, ingress, boundaries, final_time,
265            );
266        }
267        let mut next_ingress = 0;
268        let mut next_command = 0;
269        for (boundary_index, expected_boundary) in boundaries.iter().enumerate() {
270            enqueue_replay_ingress_cut(
271                &mut simulation,
272                ingress,
273                &mut next_ingress,
274                boundary_index,
275            )?;
276            for admitted in &expected_boundary.admitted_commands {
277                let Some(record) = commands.get(next_command) else {
278                    return Err(CanwuError::new(
279                        ErrorCode::ReplayMismatch,
280                        "boundary replay admits a command absent from the journal",
281                    ));
282                };
283                if record.id != *admitted {
284                    return Err(CanwuError::new(
285                        ErrorCode::ReplayMismatch,
286                        "boundary replay command admission does not match journal order",
287                    ));
288                }
289                replay_command_record(&mut simulation, record, expected_boundary.at)?;
290                next_command += 1;
291            }
292            let receipt = simulation.settle_boundary_with_state_hash_format(
293                BoundaryRequest {
294                    at: expected_boundary.at,
295                    cadences: expected_boundary.cadences.clone(),
296                },
297                boundary_state_hash_format(expected_boundary.state_hash.as_deref())?,
298            )?;
299            let Some(actual_boundary) = simulation.boundaries().last() else {
300                return Err(CanwuError::new(
301                    ErrorCode::ReplayMismatch,
302                    "boundary replay did not append settlement evidence",
303                ));
304            };
305            if receipt.boundary_id != expected_boundary.id || actual_boundary != expected_boundary {
306                return Err(CanwuError::new(
307                    ErrorCode::ReplayMismatch,
308                    format!(
309                        "regenerated boundary {} did not match its journal evidence",
310                        expected_boundary.id
311                    ),
312                ));
313            }
314        }
315        for record in &commands[next_command..] {
316            replay_command_record(&mut simulation, record, final_time)?;
317        }
318        enqueue_replay_ingress_cut(
319            &mut simulation,
320            ingress,
321            &mut next_ingress,
322            boundaries.len(),
323        )?;
324        if next_ingress != ingress.len() {
325            return Err(CanwuError::new(
326                ErrorCode::ReplayMismatch,
327                "ingress journal contains an impossible future boundary issue cut",
328            ));
329        }
330        if final_time < simulation.time() {
331            return Err(CanwuError::new(
332                ErrorCode::InvalidDuration,
333                "replay final time cannot precede the last command",
334            ));
335        }
336        if final_time > simulation.time() {
337            simulation.ensure_legacy_advance_does_not_cross_ingress(final_time)?;
338            simulation.advance_to(final_time)?;
339        }
340        Ok(simulation)
341    }
342
343    fn replay_attempt_records(
344        mut simulation: Self,
345        commands: &[CommandRecord],
346        attempts: &[CommandAttemptRecord],
347        ingress: &[IngressRecord],
348        boundaries: &[BoundaryRecord],
349        final_time: SimTime,
350    ) -> Result<Self, CanwuError> {
351        let command_ingress_requests: BTreeSet<_> = ingress
352            .iter()
353            .filter_map(|record| match &record.payload {
354                IngressPayload::Command { request } => Some(request.request_id),
355                IngressPayload::Decision { request } => {
356                    request.command.as_ref().map(|command| command.request_id)
357                }
358                IngressPayload::Plugin { .. } | IngressPayload::Calendar { .. } => None,
359            })
360            .collect();
361        let mut next_ingress = 0;
362        let mut next_attempt = 0;
363        for (boundary_index, expected_boundary) in boundaries.iter().enumerate() {
364            enqueue_replay_ingress_cut(
365                &mut simulation,
366                ingress,
367                &mut next_ingress,
368                boundary_index,
369            )?;
370            let mut admitted_commands = Vec::new();
371            for admitted in &expected_boundary.admitted_attempts {
372                let Some(record) = attempts.get(next_attempt) else {
373                    return Err(CanwuError::new(
374                        ErrorCode::ReplayMismatch,
375                        "boundary replay admits a command attempt absent from the journal",
376                    ));
377                };
378                if record.id != *admitted {
379                    return Err(CanwuError::new(
380                        ErrorCode::ReplayMismatch,
381                        "boundary replay attempt admission does not match journal order",
382                    ));
383                }
384                let queued = record
385                    .request_id
386                    .is_some_and(|request| command_ingress_requests.contains(&request));
387                if !queued {
388                    replay_attempt_record(&mut simulation, record, commands, expected_boundary.at)?;
389                }
390                if let CommandAttemptOutcome::Accepted { command_id } = record.outcome {
391                    admitted_commands.push(command_id);
392                }
393                next_attempt += 1;
394            }
395            if admitted_commands != expected_boundary.admitted_commands {
396                return Err(CanwuError::new(
397                    ErrorCode::ReplayMismatch,
398                    "boundary replay accepted-command cut disagrees with admitted attempts",
399                ));
400            }
401            let receipt = simulation.settle_boundary_with_state_hash_format(
402                BoundaryRequest {
403                    at: expected_boundary.at,
404                    cadences: expected_boundary.cadences.clone(),
405                },
406                boundary_state_hash_format(expected_boundary.state_hash.as_deref())?,
407            )?;
408            let Some(actual_boundary) = simulation.boundaries().last() else {
409                return Err(CanwuError::new(
410                    ErrorCode::ReplayMismatch,
411                    "boundary replay did not append settlement evidence",
412                ));
413            };
414            if receipt.boundary_id != expected_boundary.id || actual_boundary != expected_boundary {
415                return Err(CanwuError::new(
416                    ErrorCode::ReplayMismatch,
417                    format!(
418                        "regenerated boundary {} did not match its journal evidence",
419                        expected_boundary.id
420                    ),
421                ));
422            }
423        }
424        for record in &attempts[next_attempt..] {
425            replay_attempt_record(&mut simulation, record, commands, final_time)?;
426        }
427        enqueue_replay_ingress_cut(
428            &mut simulation,
429            ingress,
430            &mut next_ingress,
431            boundaries.len(),
432        )?;
433        if next_ingress != ingress.len() {
434            return Err(CanwuError::new(
435                ErrorCode::ReplayMismatch,
436                "ingress journal contains an impossible future boundary issue cut",
437            ));
438        }
439        if simulation.command_log() != commands {
440            return Err(CanwuError::new(
441                ErrorCode::ReplayMismatch,
442                "replayed accepted command journal does not match its recorded evidence",
443            ));
444        }
445        if final_time < simulation.time() {
446            return Err(CanwuError::new(
447                ErrorCode::InvalidDuration,
448                "replay final time cannot precede the last command attempt",
449            ));
450        }
451        if final_time > simulation.time() {
452            simulation.ensure_legacy_advance_does_not_cross_ingress(final_time)?;
453            simulation.advance_to(final_time)?;
454        }
455        Ok(simulation)
456    }
457}
458
459fn enqueue_replay_ingress_cut(
460    simulation: &mut Simulation,
461    ingress: &[IngressRecord],
462    next_ingress: &mut usize,
463    boundary_count: usize,
464) -> Result<(), CanwuError> {
465    let expected_boundary_count = u64::try_from(boundary_count).map_err(|_| {
466        CanwuError::new(
467            ErrorCode::ReplayMismatch,
468            "boundary count exceeds ingress range",
469        )
470    })?;
471    while let Some(record) = ingress.get(*next_ingress) {
472        if record.eligible_boundary_count < expected_boundary_count {
473            return Err(CanwuError::new(
474                ErrorCode::ReplayMismatch,
475                "ingress journal skipped its recorded issue boundary",
476            ));
477        }
478        if record.eligible_boundary_count > expected_boundary_count {
479            break;
480        }
481        if record.issued_at < simulation.time() {
482            return Err(CanwuError::new(
483                ErrorCode::ReplayMismatch,
484                "ingress journal issue time precedes replay state",
485            ));
486        }
487        if record.issued_at > simulation.time() {
488            simulation.ensure_legacy_advance_does_not_cross_ingress(record.issued_at)?;
489            simulation.advance_to(record.issued_at)?;
490        }
491        if let Some(actual) = simulation.state.evidence.ingress.get(*next_ingress) {
492            if actual != record {
493                return Err(CanwuError::new(
494                    ErrorCode::ReplayMismatch,
495                    "plugin-generated ingress does not match journal evidence",
496                ));
497            }
498            *next_ingress += 1;
499            continue;
500        }
501        if matches!(record.cause, Some(CauseRef::Boundary(_))) {
502            return Err(CanwuError::new(
503                ErrorCode::ReplayMismatch,
504                "recorded boundary-generated ingress was not reproduced by its plugin system",
505            ));
506        }
507        let receipt = match &record.payload {
508            IngressPayload::Command { request } => simulation.enqueue_command(
509                record.due_at,
510                record.priority,
511                request.as_ref().clone(),
512            )?,
513            IngressPayload::Plugin {
514                plugin,
515                packet_type,
516                payload,
517                affected_entities,
518            } => {
519                let mut request = PluginIngressRequest::new(
520                    plugin.clone(),
521                    packet_type.clone(),
522                    record.due_at,
523                    payload.clone(),
524                )
525                .with_priority(record.priority);
526                request.affected_entities.clone_from(affected_entities);
527                request.cause.clone_from(&record.cause);
528                simulation.enqueue_plugin_ingress(request)?
529            }
530            IngressPayload::Calendar { cadences } => {
531                simulation.schedule_calendar_boundary(record.due_at, cadences.clone())?
532            }
533            IngressPayload::Decision { request } => simulation.enqueue_decision(
534                record.due_at,
535                record.priority,
536                request.as_ref().clone(),
537            )?,
538        };
539        if receipt.ingress_id != record.id
540            || simulation.state.evidence.ingress.last() != Some(record)
541        {
542            return Err(CanwuError::new(
543                ErrorCode::ReplayMismatch,
544                "regenerated ingress record does not match journal evidence",
545            ));
546        }
547        *next_ingress += 1;
548    }
549    Ok(())
550}
551
552fn replay_command_record(
553    simulation: &mut Simulation,
554    record: &CommandRecord,
555    latest_time: SimTime,
556) -> Result<(), CanwuError> {
557    if record.accepted_at < simulation.time() || record.accepted_at > latest_time {
558        return Err(CanwuError::new(
559            ErrorCode::ReplayMismatch,
560            "replay command timestamps do not match authoritative operation order",
561        ));
562    }
563    simulation.ensure_legacy_advance_does_not_cross_ingress(record.accepted_at)?;
564    simulation.advance_to(record.accepted_at)?;
565    let CommandOutcome::Accepted { receipt } = simulation.admit_command(
566        None,
567        None,
568        record.envelope.clone(),
569        CommandIngress::LegacyDirect,
570        None,
571        false,
572    )?
573    else {
574        return Err(CanwuError::new(
575            ErrorCode::ReplayMismatch,
576            "legacy replay command was rejected",
577        ));
578    };
579    if receipt.command_id != record.id {
580        return Err(CanwuError::new(
581            ErrorCode::ReplayMismatch,
582            "replay command IDs did not match the journal",
583        ));
584    }
585    Ok(())
586}
587
588fn replay_attempt_record(
589    simulation: &mut Simulation,
590    record: &CommandAttemptRecord,
591    commands: &[CommandRecord],
592    latest_time: SimTime,
593) -> Result<(), CanwuError> {
594    if record.at < simulation.time() || record.at > latest_time {
595        return Err(CanwuError::new(
596            ErrorCode::ReplayMismatch,
597            "replay command-attempt timestamps do not match authoritative operation order",
598        ));
599    }
600    simulation.ensure_legacy_advance_does_not_cross_ingress(record.at)?;
601    simulation.advance_to(record.at)?;
602    let outcome = simulation.admit_command(
603        record.request_id,
604        record.expected_revision,
605        record.envelope.clone(),
606        record.ingress,
607        None,
608        true,
609    )?;
610    if simulation.command_attempts().last() != Some(record) {
611        return Err(CanwuError::new(
612            ErrorCode::ReplayMismatch,
613            format!(
614                "regenerated command attempt {} did not match its journal evidence",
615                record.id
616            ),
617        ));
618    }
619    match (&record.outcome, outcome) {
620        (CommandAttemptOutcome::Accepted { command_id }, CommandOutcome::Accepted { receipt })
621            if receipt.command_id == *command_id =>
622        {
623            let index = usize::try_from(command_id.get().saturating_sub(1)).map_err(|_| {
624                CanwuError::new(
625                    ErrorCode::ReplayMismatch,
626                    "replayed command ID exceeds the journal index range",
627                )
628            })?;
629            if simulation.command_log().last() != commands.get(index) {
630                return Err(CanwuError::new(
631                    ErrorCode::ReplayMismatch,
632                    "regenerated command record did not match its journal evidence",
633                ));
634            }
635        }
636        (
637            CommandAttemptOutcome::Rejected { error: expected },
638            CommandOutcome::Rejected { rejection },
639        ) if rejection.error == *expected => {}
640        _ => {
641            return Err(CanwuError::new(
642                ErrorCode::ReplayMismatch,
643                "replayed command-attempt outcome differs from its journal evidence",
644            ));
645        }
646    }
647    Ok(())
648}