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