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 #[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 #[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 #[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 #[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 #[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 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 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}