Skip to main content

canwu_sim/runtime/
replay.rs

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