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::Plugin { .. } | IngressPayload::Calendar { .. } => None,
356 })
357 .collect();
358 let mut next_ingress = 0;
359 let mut next_attempt = 0;
360 for (boundary_index, expected_boundary) in boundaries.iter().enumerate() {
361 enqueue_replay_ingress_cut(
362 &mut simulation,
363 ingress,
364 &mut next_ingress,
365 boundary_index,
366 )?;
367 let mut admitted_commands = Vec::new();
368 for admitted in &expected_boundary.admitted_attempts {
369 let Some(record) = attempts.get(next_attempt) else {
370 return Err(CanwuError::new(
371 ErrorCode::ReplayMismatch,
372 "boundary replay admits a command attempt absent from the journal",
373 ));
374 };
375 if record.id != *admitted {
376 return Err(CanwuError::new(
377 ErrorCode::ReplayMismatch,
378 "boundary replay attempt admission does not match journal order",
379 ));
380 }
381 let queued = record
382 .request_id
383 .is_some_and(|request| command_ingress_requests.contains(&request));
384 if !queued {
385 replay_attempt_record(&mut simulation, record, commands, expected_boundary.at)?;
386 }
387 if let CommandAttemptOutcome::Accepted { command_id } = record.outcome {
388 admitted_commands.push(command_id);
389 }
390 next_attempt += 1;
391 }
392 if admitted_commands != expected_boundary.admitted_commands {
393 return Err(CanwuError::new(
394 ErrorCode::ReplayMismatch,
395 "boundary replay accepted-command cut disagrees with admitted attempts",
396 ));
397 }
398 let receipt = simulation.settle_boundary_with_state_hash_format(
399 BoundaryRequest {
400 at: expected_boundary.at,
401 cadences: expected_boundary.cadences.clone(),
402 },
403 boundary_state_hash_format(expected_boundary.state_hash.as_deref())?,
404 )?;
405 let Some(actual_boundary) = simulation.boundaries().last() else {
406 return Err(CanwuError::new(
407 ErrorCode::ReplayMismatch,
408 "boundary replay did not append settlement evidence",
409 ));
410 };
411 if receipt.boundary_id != expected_boundary.id || actual_boundary != expected_boundary {
412 return Err(CanwuError::new(
413 ErrorCode::ReplayMismatch,
414 format!(
415 "regenerated boundary {} did not match its journal evidence",
416 expected_boundary.id
417 ),
418 ));
419 }
420 }
421 for record in &attempts[next_attempt..] {
422 replay_attempt_record(&mut simulation, record, commands, final_time)?;
423 }
424 enqueue_replay_ingress_cut(
425 &mut simulation,
426 ingress,
427 &mut next_ingress,
428 boundaries.len(),
429 )?;
430 if next_ingress != ingress.len() {
431 return Err(CanwuError::new(
432 ErrorCode::ReplayMismatch,
433 "ingress journal contains an impossible future boundary issue cut",
434 ));
435 }
436 if simulation.command_log() != commands {
437 return Err(CanwuError::new(
438 ErrorCode::ReplayMismatch,
439 "replayed accepted command journal does not match its recorded evidence",
440 ));
441 }
442 if final_time < simulation.time() {
443 return Err(CanwuError::new(
444 ErrorCode::InvalidDuration,
445 "replay final time cannot precede the last command attempt",
446 ));
447 }
448 if final_time > simulation.time() {
449 simulation.ensure_legacy_advance_does_not_cross_ingress(final_time)?;
450 simulation.advance_to(final_time)?;
451 }
452 Ok(simulation)
453 }
454}
455
456fn enqueue_replay_ingress_cut(
457 simulation: &mut Simulation,
458 ingress: &[IngressRecord],
459 next_ingress: &mut usize,
460 boundary_count: usize,
461) -> Result<(), CanwuError> {
462 let expected_boundary_count = u64::try_from(boundary_count).map_err(|_| {
463 CanwuError::new(
464 ErrorCode::ReplayMismatch,
465 "boundary count exceeds ingress range",
466 )
467 })?;
468 while let Some(record) = ingress.get(*next_ingress) {
469 if record.eligible_boundary_count < expected_boundary_count {
470 return Err(CanwuError::new(
471 ErrorCode::ReplayMismatch,
472 "ingress journal skipped its recorded issue boundary",
473 ));
474 }
475 if record.eligible_boundary_count > expected_boundary_count {
476 break;
477 }
478 if record.issued_at < simulation.time() {
479 return Err(CanwuError::new(
480 ErrorCode::ReplayMismatch,
481 "ingress journal issue time precedes replay state",
482 ));
483 }
484 if record.issued_at > simulation.time() {
485 simulation.ensure_legacy_advance_does_not_cross_ingress(record.issued_at)?;
486 simulation.advance_to(record.issued_at)?;
487 }
488 if let Some(actual) = simulation.state.evidence.ingress.get(*next_ingress) {
489 if actual != record {
490 return Err(CanwuError::new(
491 ErrorCode::ReplayMismatch,
492 "plugin-generated ingress does not match journal evidence",
493 ));
494 }
495 *next_ingress += 1;
496 continue;
497 }
498 if matches!(record.cause, Some(CauseRef::Boundary(_))) {
499 return Err(CanwuError::new(
500 ErrorCode::ReplayMismatch,
501 "recorded boundary-generated ingress was not reproduced by its plugin system",
502 ));
503 }
504 let receipt = match &record.payload {
505 IngressPayload::Command { request } => simulation.enqueue_command(
506 record.due_at,
507 record.priority,
508 request.as_ref().clone(),
509 )?,
510 IngressPayload::Plugin {
511 plugin,
512 packet_type,
513 payload,
514 affected_entities,
515 } => {
516 let mut request = PluginIngressRequest::new(
517 plugin.clone(),
518 packet_type.clone(),
519 record.due_at,
520 payload.clone(),
521 )
522 .with_priority(record.priority);
523 request.affected_entities.clone_from(affected_entities);
524 request.cause.clone_from(&record.cause);
525 simulation.enqueue_plugin_ingress(request)?
526 }
527 IngressPayload::Calendar { cadences } => {
528 simulation.schedule_calendar_boundary(record.due_at, cadences.clone())?
529 }
530 };
531 if receipt.ingress_id != record.id
532 || simulation.state.evidence.ingress.last() != Some(record)
533 {
534 return Err(CanwuError::new(
535 ErrorCode::ReplayMismatch,
536 "regenerated ingress record does not match journal evidence",
537 ));
538 }
539 *next_ingress += 1;
540 }
541 Ok(())
542}
543
544fn replay_command_record(
545 simulation: &mut Simulation,
546 record: &CommandRecord,
547 latest_time: SimTime,
548) -> Result<(), CanwuError> {
549 if record.accepted_at < simulation.time() || record.accepted_at > latest_time {
550 return Err(CanwuError::new(
551 ErrorCode::ReplayMismatch,
552 "replay command timestamps do not match authoritative operation order",
553 ));
554 }
555 simulation.ensure_legacy_advance_does_not_cross_ingress(record.accepted_at)?;
556 simulation.advance_to(record.accepted_at)?;
557 let CommandOutcome::Accepted { receipt } = simulation.admit_command(
558 None,
559 None,
560 record.envelope.clone(),
561 CommandIngress::LegacyDirect,
562 false,
563 )?
564 else {
565 return Err(CanwuError::new(
566 ErrorCode::ReplayMismatch,
567 "legacy replay command was rejected",
568 ));
569 };
570 if receipt.command_id != record.id {
571 return Err(CanwuError::new(
572 ErrorCode::ReplayMismatch,
573 "replay command IDs did not match the journal",
574 ));
575 }
576 Ok(())
577}
578
579fn replay_attempt_record(
580 simulation: &mut Simulation,
581 record: &CommandAttemptRecord,
582 commands: &[CommandRecord],
583 latest_time: SimTime,
584) -> Result<(), CanwuError> {
585 if record.at < simulation.time() || record.at > latest_time {
586 return Err(CanwuError::new(
587 ErrorCode::ReplayMismatch,
588 "replay command-attempt timestamps do not match authoritative operation order",
589 ));
590 }
591 simulation.ensure_legacy_advance_does_not_cross_ingress(record.at)?;
592 simulation.advance_to(record.at)?;
593 let outcome = simulation.admit_command(
594 record.request_id,
595 record.expected_revision,
596 record.envelope.clone(),
597 record.ingress,
598 true,
599 )?;
600 if simulation.command_attempts().last() != Some(record) {
601 return Err(CanwuError::new(
602 ErrorCode::ReplayMismatch,
603 format!(
604 "regenerated command attempt {} did not match its journal evidence",
605 record.id
606 ),
607 ));
608 }
609 match (&record.outcome, outcome) {
610 (CommandAttemptOutcome::Accepted { command_id }, CommandOutcome::Accepted { receipt })
611 if receipt.command_id == *command_id =>
612 {
613 let index = usize::try_from(command_id.get().saturating_sub(1)).map_err(|_| {
614 CanwuError::new(
615 ErrorCode::ReplayMismatch,
616 "replayed command ID exceeds the journal index range",
617 )
618 })?;
619 if simulation.command_log().last() != commands.get(index) {
620 return Err(CanwuError::new(
621 ErrorCode::ReplayMismatch,
622 "regenerated command record did not match its journal evidence",
623 ));
624 }
625 }
626 (
627 CommandAttemptOutcome::Rejected { error: expected },
628 CommandOutcome::Rejected { rejection },
629 ) if rejection.error == *expected => {}
630 _ => {
631 return Err(CanwuError::new(
632 ErrorCode::ReplayMismatch,
633 "replayed command-attempt outcome differs from its journal evidence",
634 ));
635 }
636 }
637 Ok(())
638}