1use std::collections::HashMap;
7use std::ops::{Range, RangeInclusive};
8use std::time::Duration;
9
10use super::wall_clock::Instant;
11use tracing::instrument;
12
13use crate::SimulationError;
14use crate::observability::{Invariant, SimulationLayer, SimulationLayerHandle, TraceQuery};
15use crate::runner::fault_injector::FaultInjector;
16use crate::runner::locality::{LocalityConfig, MachineRegistry};
17use crate::runner::process::{Attrition, Process};
18use crate::runner::tags::TagDistribution;
19use crate::runner::workload::Workload;
20
21use super::orchestrator::{
22 GenerateReportInputs, IterationManager, MetricsCollector, OrchestrateInputs, OrchestrateOutput,
23 WorkloadOrchestrator,
24};
25
26#[derive(Debug, Clone, Copy)]
28pub(crate) struct WorkloadClientInfo {
29 pub(crate) client_id: usize,
31 pub(crate) client_count: usize,
33}
34
35struct RunOrchestratorInputs<'a> {
37 seed: u64,
38 iteration_count: usize,
39 workloads: Vec<Box<dyn Workload>>,
40 workload_info: Vec<(String, String)>,
41 client_info: Vec<WorkloadClientInfo>,
42 process_config: Option<super::orchestrator::ProcessConfig<'a>>,
43 sim: crate::sim::SimWorld,
44 fault_injectors: Vec<Box<dyn FaultInjector>>,
45 chaos_duration: Option<Duration>,
46 obs_handle: SimulationLayerHandle,
47 run_time_budget: Duration,
48}
49
50type OrchestrationOutcome = Result<OrchestrateOutput, (Vec<u64>, usize)>;
52
53struct FinalReportInputs {
55 converged: bool,
56 saturation: Option<super::report::SaturationReport>,
58 #[cfg(feature = "exploration")]
59 total_exploration_timelines: u64,
60 #[cfg(feature = "exploration")]
61 total_exploration_fork_points: u64,
62 #[cfg(feature = "exploration")]
63 total_exploration_bugs: u64,
64 #[cfg(feature = "exploration")]
65 bug_recipes: Vec<super::report::BugRecipe>,
66 #[cfg(feature = "exploration")]
67 per_seed_timelines: Vec<u64>,
68}
69
70struct ConvergenceState<'a> {
72 iteration_control: &'a IterationControl,
73 iteration_count: usize,
74 reached_sometimes: &'a std::collections::HashSet<String>,
75 all_sometimes_count: usize,
76 exploration_active: bool,
79 prev_signal: &'a mut usize,
80 plateau_count: &'a mut usize,
81 saturation: &'a mut Option<super::report::SaturationReport>,
83 already_converged: bool,
84}
85
86impl RunState {
87 fn new(builder: &SimulationBuilder) -> Self {
89 let iteration_manager =
90 IterationManager::new(builder.iteration_control.clone(), builder.seeds.clone());
91 let progress_milestone = iteration_manager
92 .max_iterations()
93 .map(|max| std::cmp::max(max / 10, 1));
94 Self {
95 iteration_manager,
96 metrics_collector: MetricsCollector::new(),
97 progress_milestone,
98 pending_return_map: Vec::new(),
99 #[cfg(feature = "exploration")]
100 total_exploration_timelines: 0,
101 #[cfg(feature = "exploration")]
102 total_exploration_fork_points: 0,
103 #[cfg(feature = "exploration")]
104 total_exploration_bugs: 0,
105 #[cfg(feature = "exploration")]
106 bug_recipes: Vec::new(),
107 #[cfg(feature = "exploration")]
108 per_seed_timelines: Vec::new(),
109 reached_sometimes: std::collections::HashSet::new(),
110 prev_signal: 0,
111 converged: false,
112 plateau_count: 0,
113 saturation: None,
114 }
115 }
116}
117
118struct RunState {
120 iteration_manager: IterationManager,
121 metrics_collector: MetricsCollector,
122 progress_milestone: Option<usize>,
124 pending_return_map: Vec<Option<usize>>,
128 #[cfg(feature = "exploration")]
130 total_exploration_timelines: u64,
131 #[cfg(feature = "exploration")]
132 total_exploration_fork_points: u64,
133 #[cfg(feature = "exploration")]
134 total_exploration_bugs: u64,
135 #[cfg(feature = "exploration")]
136 bug_recipes: Vec<super::report::BugRecipe>,
137 #[cfg(feature = "exploration")]
138 per_seed_timelines: Vec<u64>,
139 reached_sometimes: std::collections::HashSet<String>,
141 prev_signal: usize,
144 converged: bool,
145 plateau_count: usize,
146 saturation: Option<super::report::SaturationReport>,
148}
149
150struct ResolvedEntries {
152 workloads: Vec<Box<dyn Workload>>,
153 return_map: Vec<Option<usize>>,
156 client_info: Vec<WorkloadClientInfo>,
158}
159use super::report::{SimulationMetrics, SimulationReport};
160
161#[derive(Debug, Clone)]
165pub enum IterationControl {
166 FixedCount(usize),
168 TimeLimit(Duration),
170 UntilCoverageStable {
180 plateau_seeds: usize,
182 max_iterations: usize,
184 },
185}
186
187#[derive(Debug, Clone)]
202pub enum WorkloadCount {
203 Fixed(usize),
205 Random(Range<usize>),
208}
209
210impl WorkloadCount {
211 fn resolve(&self) -> usize {
214 match self {
215 WorkloadCount::Fixed(n) => *n,
216 WorkloadCount::Random(range) => crate::sim::sim_random_range(range.clone()),
217 }
218 }
219}
220
221#[derive(Debug, Clone, PartialEq)]
240pub enum ClientId {
241 Fixed(usize),
244 RandomRange(Range<usize>),
248}
249
250impl Default for ClientId {
251 fn default() -> Self {
252 Self::Fixed(0)
253 }
254}
255
256impl ClientId {
257 fn resolve(&self, index: usize) -> usize {
259 match self {
260 ClientId::Fixed(base) => base + index,
261 ClientId::RandomRange(range) => crate::sim::sim_random_range(range.clone()),
262 }
263 }
264}
265
266#[derive(Debug, Clone, PartialEq)]
281pub enum ProcessCount {
282 Fixed(usize),
284 Range(RangeInclusive<usize>),
287}
288
289impl ProcessCount {
290 pub(crate) fn resolve(&self) -> usize {
292 match self {
293 ProcessCount::Fixed(n) => *n,
294 ProcessCount::Range(range) => {
295 let start = *range.start();
296 let end = *range.end() + 1; if start >= end {
298 return start;
299 }
300 crate::sim::sim_random_range(start..end)
301 }
302 }
303 }
304}
305
306impl From<usize> for ProcessCount {
307 fn from(n: usize) -> Self {
308 ProcessCount::Fixed(n)
309 }
310}
311
312impl From<RangeInclusive<usize>> for ProcessCount {
313 fn from(range: RangeInclusive<usize>) -> Self {
314 ProcessCount::Range(range)
315 }
316}
317
318pub(crate) struct ProcessEntry {
320 pub(crate) count: ProcessCount,
321 pub(crate) factory: Box<dyn Fn() -> Box<dyn Process>>,
322 pub(crate) tags: TagDistribution,
323 pub(crate) name: String,
324 pub(crate) locality: Option<LocalityConfig>,
327}
328
329enum WorkloadEntry {
331 Instance(Option<Box<dyn Workload>>, ClientId),
333 Factory {
335 count: WorkloadCount,
336 client_id: ClientId,
337 factory: Box<dyn Fn(usize) -> Box<dyn Workload>>,
338 },
339}
340
341#[derive(Debug, Clone, Copy, PartialEq, Eq)]
348pub enum ChaosMode {
349 Random,
354 Swarm,
360}
361
362#[derive(Debug, Clone, PartialEq)]
368pub enum Chaos {
369 Network(ChaosMode),
371 Storage(ChaosMode),
373 Attrition {
376 config: Attrition,
378 mode: ChaosMode,
380 },
381 BuggifyKnobs,
387}
388
389pub struct SimulationBuilder {
391 iteration_control: IterationControl,
392 entries: Vec<WorkloadEntry>,
393 process_entry: Option<ProcessEntry>,
394 attrition: Option<Attrition>,
395 attrition_mode: ChaosMode,
396 seeds: Vec<u64>,
397 network_chaos: Option<ChaosMode>,
398 storage_chaos: Option<ChaosMode>,
399 buggify_knobs: bool,
403 swarm_operations: bool,
404 invariants: Vec<Box<dyn Invariant + Send>>,
405 fault_injectors: Vec<Box<dyn FaultInjector>>,
406 chaos_duration: Option<Duration>,
407 exploration_config: Option<crate::chaos::exploration_glue::ExplorationConfig>,
408 before_iteration_hooks: Vec<Box<dyn FnMut()>>,
409 seed_warning_timeout: Option<Duration>,
410 run_time_budget: Duration,
411}
412
413impl Default for SimulationBuilder {
414 fn default() -> Self {
415 Self::new()
416 }
417}
418
419impl SimulationBuilder {
420 #[must_use]
422 pub fn new() -> Self {
423 Self {
424 iteration_control: IterationControl::UntilCoverageStable {
425 plateau_seeds: 10,
426 max_iterations: 1000,
427 },
428 entries: Vec::new(),
429 process_entry: None,
430 attrition: None,
431 attrition_mode: ChaosMode::Random,
432 seeds: Vec::new(),
433 network_chaos: None,
434 storage_chaos: None,
435 buggify_knobs: false,
436 swarm_operations: false,
437 invariants: Vec::new(),
438 fault_injectors: Vec::new(),
439 chaos_duration: None,
440 exploration_config: None,
441 before_iteration_hooks: Vec::new(),
442 seed_warning_timeout: None,
443 run_time_budget: super::orchestrator::DEFAULT_RUN_TIME_BUDGET,
444 }
445 }
446
447 #[must_use]
452 pub fn workload(mut self, w: impl Workload) -> Self {
453 self.entries.push(WorkloadEntry::Instance(
454 Some(Box::new(w)),
455 ClientId::default(),
456 ));
457 self
458 }
459
460 #[must_use]
482 pub fn processes(
483 mut self,
484 count: impl Into<ProcessCount>,
485 factory: impl Fn() -> Box<dyn Process> + 'static,
486 ) -> Self {
487 let sample = factory();
488 let name = sample.name().to_string();
489 drop(sample);
490 self.process_entry = Some(ProcessEntry {
491 count: count.into(),
492 factory: Box::new(factory),
493 tags: TagDistribution::new(),
494 name,
495 locality: None,
496 });
497 self
498 }
499
500 #[must_use]
521 pub fn cluster(
522 mut self,
523 config: LocalityConfig,
524 factory: impl Fn() -> Box<dyn Process> + 'static,
525 ) -> Self {
526 let sample = factory();
527 let name = sample.name().to_string();
528 drop(sample);
529 self.process_entry = Some(ProcessEntry {
530 count: ProcessCount::Fixed(0),
532 factory: Box::new(factory),
533 tags: TagDistribution::new(),
534 name,
535 locality: Some(config),
536 });
537 self
538 }
539
540 pub fn tags(mut self, dimensions: &[(&str, &[&str])]) -> Result<Self, SimulationError> {
561 let entry = self.process_entry.as_mut().ok_or_else(|| {
562 SimulationError::InvalidState("tags() must be called after processes()".into())
563 })?;
564 for (key, values) in dimensions {
565 entry.tags.add(key, values);
566 }
567 Ok(self)
568 }
569
570 #[must_use]
581 pub fn attrition(mut self, config: Attrition) -> Self {
582 self.attrition = Some(config);
583 self
584 }
585
586 #[must_use]
605 pub fn workloads(
606 mut self,
607 count: WorkloadCount,
608 factory: impl Fn(usize) -> Box<dyn Workload> + 'static,
609 ) -> Self {
610 self.entries.push(WorkloadEntry::Factory {
611 count,
612 client_id: ClientId::default(),
613 factory: Box::new(factory),
614 });
615 self
616 }
617
618 #[must_use]
620 pub fn invariant<I: Invariant>(mut self, i: I) -> Self {
621 self.invariants.push(Box::new(i));
622 self
623 }
624
625 #[must_use]
627 pub fn invariant_fn(
628 mut self,
629 name: impl Into<String>,
630 f: impl Fn(&dyn TraceQuery, u64) + Send + 'static,
631 ) -> Self {
632 self.invariants
633 .push(crate::observability::invariant_fn(name, f));
634 self
635 }
636
637 #[must_use]
639 pub fn fault(mut self, f: impl FaultInjector) -> Self {
640 self.fault_injectors.push(Box::new(f));
641 self
642 }
643
644 #[must_use]
651 pub fn chaos_duration(mut self, duration: Duration) -> Self {
652 self.chaos_duration = Some(duration);
653 self
654 }
655
656 #[must_use]
658 pub fn set_iterations(mut self, iterations: usize) -> Self {
659 self.iteration_control = IterationControl::FixedCount(iterations);
660 self
661 }
662
663 #[must_use]
668 pub fn seed_warning_timeout(mut self, timeout: Duration) -> Self {
669 self.seed_warning_timeout = Some(timeout);
670 self
671 }
672
673 #[must_use]
692 pub fn run_time_budget(mut self, budget: Duration) -> Self {
693 self.run_time_budget = budget;
694 self
695 }
696
697 #[must_use]
708 pub fn until_coverage_stable(mut self, plateau_seeds: usize, max_iterations: usize) -> Self {
709 self.iteration_control = IterationControl::UntilCoverageStable {
710 plateau_seeds,
711 max_iterations,
712 };
713 self
714 }
715
716 #[must_use]
721 pub fn before_iteration(mut self, f: impl FnMut() + 'static) -> Self {
722 self.before_iteration_hooks.push(Box::new(f));
723 self
724 }
725
726 #[must_use]
728 pub fn set_debug_seeds(mut self, seeds: Vec<u64>) -> Self {
729 self.seeds = seeds;
730 self
731 }
732
733 #[must_use]
760 pub fn enable_chaos(mut self, surfaces: impl IntoIterator<Item = Chaos>) -> Self {
761 for surface in surfaces {
762 match surface {
763 Chaos::Network(mode) => self.network_chaos = Some(mode),
764 Chaos::Storage(mode) => self.storage_chaos = Some(mode),
765 Chaos::Attrition { config, mode } => {
766 self.attrition = Some(config);
767 self.attrition_mode = mode;
768 }
769 Chaos::BuggifyKnobs => self.buggify_knobs = true,
770 }
771 }
772 self
773 }
774
775 #[must_use]
783 pub fn swarm_operations(mut self) -> Self {
784 self.swarm_operations = true;
785 self
786 }
787
788 #[cfg(feature = "exploration")]
794 #[must_use]
795 pub fn enable_exploration(
796 mut self,
797 config: crate::chaos::exploration_glue::ExplorationConfig,
798 ) -> Self {
799 self.exploration_config = Some(config);
800 self
801 }
802
803 fn resolve_entries(&mut self) -> ResolvedEntries {
805 let mut workloads = Vec::new();
806 let mut return_map = Vec::new();
807 let mut client_info = Vec::new();
808
809 for (entry_idx, entry) in self.entries.iter_mut().enumerate() {
810 match entry {
811 WorkloadEntry::Instance(opt, cid) => {
812 if let Some(w) = opt.take() {
813 return_map.push(Some(entry_idx));
814 client_info.push(WorkloadClientInfo {
815 client_id: cid.resolve(0),
816 client_count: 1,
817 });
818 workloads.push(w);
819 }
820 }
821 WorkloadEntry::Factory {
822 count,
823 client_id,
824 factory,
825 } => {
826 let n = count.resolve();
827 for i in 0..n {
828 return_map.push(None);
829 client_info.push(WorkloadClientInfo {
830 client_id: client_id.resolve(i),
831 client_count: n,
832 });
833 workloads.push(factory(i));
834 }
835 }
836 }
837 }
838
839 ResolvedEntries {
840 workloads,
841 return_map,
842 client_info,
843 }
844 }
845
846 fn return_entries(
848 &mut self,
849 workloads: Vec<Box<dyn Workload>>,
850 return_map: Vec<Option<usize>>,
851 ) {
852 for (w, slot) in workloads.into_iter().zip(return_map) {
853 if let Some(entry_idx) = slot
854 && let WorkloadEntry::Instance(opt, _) = &mut self.entries[entry_idx]
855 {
856 *opt = Some(w);
857 }
858 }
860 }
861
862 fn run_orchestrator_blocking(inputs: RunOrchestratorInputs<'_>) -> OrchestrationOutcome {
865 let RunOrchestratorInputs {
866 seed,
867 iteration_count,
868 workloads,
869 workload_info,
870 client_info,
871 process_config,
872 sim,
873 fault_injectors,
874 chaos_duration,
875 obs_handle,
876 run_time_budget,
877 } = inputs;
878 let mut executor = crate::executor::Executor::new(seed);
882 executor.block_on(async move {
883 WorkloadOrchestrator::orchestrate_workloads(OrchestrateInputs {
884 workloads,
885 fault_injectors,
886 obs: obs_handle,
887 workload_info: &workload_info,
888 client_info: &client_info,
889 process_config,
890 seed,
891 sim,
892 chaos_duration,
893 iteration_count,
894 run_time_budget,
895 })
896 .await
897 })
898 }
899
900 fn build_sim_for_iteration(
907 network_chaos: Option<ChaosMode>,
908 storage_chaos: Option<ChaosMode>,
909 buggify_knobs: bool,
910 seed: u64,
911 ) -> crate::sim::SimWorld {
912 let mut network_config = match network_chaos {
913 Some(ChaosMode::Swarm) => crate::NetworkConfiguration::swarm_for_seed(),
914 Some(ChaosMode::Random) => crate::NetworkConfiguration::random_for_seed(),
915 None => crate::NetworkConfiguration::default(),
916 };
917 let mut storage_config = match storage_chaos {
918 Some(ChaosMode::Swarm) => crate::storage::StorageConfiguration::swarm_for_seed(),
919 Some(ChaosMode::Random) => crate::storage::StorageConfiguration::random_for_seed(),
920 None => crate::storage::StorageConfiguration::default(),
921 };
922 if buggify_knobs {
927 if network_chaos.is_some() {
928 network_config.chaos.apply_buggify_knobs();
929 }
930 if storage_chaos.is_some() {
931 storage_config.apply_buggify_knobs();
932 }
933 }
934 let mut sim = crate::sim::SimWorld::new_with_network_config_and_seed(network_config, seed);
935 sim.set_storage_config(storage_config);
936 sim
937 }
938
939 fn collect_fault_injectors(
942 user_injectors: &mut Vec<Box<dyn FaultInjector>>,
943 attrition: Option<&Attrition>,
944 ) -> Vec<Box<dyn FaultInjector>> {
945 let mut fault_injectors = std::mem::take(user_injectors);
946 if let Some(attrition) = attrition {
947 fault_injectors.push(Box::new(
948 crate::runner::fault_injector::AttritionInjector::new(attrition.clone()),
949 ));
950 }
951 fault_injectors
952 }
953
954 fn build_early_exit_report(
957 metrics_collector: MetricsCollector,
958 iteration_count: usize,
959 seeds_used: Vec<u64>,
960 ) -> SimulationReport {
961 let assertion_results = crate::chaos::assertion_results();
962 let (assertion_violations, coverage_violations) =
963 crate::chaos::validate_assertion_contracts();
964 crate::chaos::buggify_reset();
965 metrics_collector.generate_report(GenerateReportInputs {
966 iteration_count,
967 seeds_used,
968 assertion_results,
969 assertion_violations,
970 coverage_violations,
971 exploration: None,
972 assertion_details: Vec::new(),
973 bucket_summaries: Vec::new(),
974 convergence_timeout: false,
975 saturation: None,
976 })
977 }
978
979 fn check_convergence_or_plateau(state: ConvergenceState<'_>) -> bool {
988 let ConvergenceState {
989 iteration_control,
990 iteration_count,
991 reached_sometimes,
992 all_sometimes_count,
993 exploration_active,
994 prev_signal,
995 plateau_count,
996 saturation,
997 already_converged,
998 } = state;
999 if already_converged {
1000 return true;
1001 }
1002 let IterationControl::UntilCoverageStable { plateau_seeds, .. } = iteration_control else {
1003 return false;
1004 };
1005
1006 let edges = crate::chaos::exploration_glue::code_coverage_edges(exploration_active);
1009 let (signal, current) = match edges {
1010 Some(n) => (super::report::SaturationSignal::CodeCoverage, n),
1011 None => (
1012 super::report::SaturationSignal::AssertionCoverage,
1013 reached_sometimes.len(),
1014 ),
1015 };
1016
1017 if iteration_count == 1 {
1018 *prev_signal = current;
1019 } else if current == *prev_signal {
1020 *plateau_count += 1;
1021 } else {
1022 *plateau_count = 0;
1023 *prev_signal = current;
1024 }
1025
1026 let all_reached = all_sometimes_count > 0 && reached_sometimes.len() >= all_sometimes_count;
1027
1028 let edges_total = crate::chaos::exploration_glue::code_coverage_total().unwrap_or_default();
1029 *saturation = Some(super::report::SaturationReport {
1030 signal,
1031 edges_covered: edges.unwrap_or_default(),
1032 edges_total,
1033 sometimes_hit: reached_sometimes.len(),
1034 sometimes_total: all_sometimes_count,
1035 plateau_seeds: *plateau_seeds,
1036 });
1037
1038 tracing::warn!(
1039 "saturation: seed={} sometimes={}/{} signal={:?}={} quiet_seeds={}/{}",
1040 iteration_count,
1041 reached_sometimes.len(),
1042 all_sometimes_count,
1043 signal,
1044 current,
1045 *plateau_count,
1046 plateau_seeds,
1047 );
1048 if *plateau_count >= *plateau_seeds && all_reached {
1049 tracing::info!(
1050 "Saturated after {} seeds: all {} sometimes reached, {:?} stable ({}) for {} seeds",
1051 iteration_count,
1052 all_sometimes_count,
1053 signal,
1054 current,
1055 *plateau_count,
1056 );
1057 return true;
1058 }
1059 false
1060 }
1061
1062 fn log_slow_seed(seed: u64, wall_time: Duration, threshold: Option<Duration>) {
1064 if let Some(threshold) = threshold
1065 && wall_time > threshold
1066 {
1067 tracing::warn!(
1068 seed,
1069 wall_time_ms = u64::try_from(wall_time.as_millis()).unwrap_or(u64::MAX),
1070 threshold_ms = u64::try_from(threshold.as_millis()).unwrap_or(u64::MAX),
1071 "seed took {:.2}s (threshold: {}s)",
1072 wall_time.as_secs_f64(),
1073 threshold.as_secs(),
1074 );
1075 }
1076 }
1077
1078 fn log_progress_milestone(
1080 progress_milestone: Option<usize>,
1081 iteration_count: usize,
1082 max: usize,
1083 ) {
1084 if let Some(interval) = progress_milestone
1085 && iteration_count.is_multiple_of(interval)
1086 {
1087 let iteration_f64 = u32::try_from(iteration_count).map_or(f64::INFINITY, f64::from);
1088 let max_f64 = u32::try_from(max).map_or(f64::INFINITY, f64::from);
1089 let pct = (iteration_f64 / max_f64) * 100.0;
1090 tracing::info!(
1091 iteration = iteration_count,
1092 total = max,
1093 "[{}/{}] {:.0}% complete",
1094 iteration_count,
1095 max,
1096 pct,
1097 );
1098 }
1099 }
1100
1101 fn reset_per_iteration_state(
1103 seed: u64,
1104 swarm_operations: bool,
1105 obs_handle: &SimulationLayerHandle,
1106 ) {
1107 obs_handle.reset_for_seed();
1108 crate::sim::reset_sim_rng();
1109 crate::sim::set_sim_seed(seed);
1110 crate::sim::set_config_seed(seed);
1113 crate::sim::set_select_seed(seed);
1116 crate::sim::set_swarm_op_seed(swarm_operations.then_some(seed));
1119 crate::chaos::reset_always_violations();
1120 crate::chaos::buggify_init(0.5, 0.25);
1122 }
1123
1124 fn resolve_process_config(entry: &ProcessEntry) -> super::orchestrator::ProcessConfig<'_> {
1127 let localities = entry
1130 .locality
1131 .as_ref()
1132 .map(LocalityConfig::resolve_topology);
1133 let count = localities
1134 .as_ref()
1135 .map_or_else(|| entry.count.resolve(), Vec::len);
1136
1137 let mut registry = crate::runner::tags::TagRegistry::new();
1138 let mut machine_registry = MachineRegistry::new();
1139 let mut ips = Vec::with_capacity(count);
1140 let mut info = Vec::with_capacity(count);
1141 let base_name = &entry.name;
1142 for i in 0..count {
1143 let ip = format!("10.0.1.{}", i + 1);
1144 let ip_addr: std::net::IpAddr = ip.parse().expect("valid process IP");
1145 let tags = entry.tags.resolve(i);
1146 registry.register(ip_addr, tags);
1147 if let Some(localities) = &localities {
1148 machine_registry.register(ip_addr, localities[i].clone());
1149 }
1150 ips.push(ip.clone());
1151 let name = if count == 1 {
1152 base_name.clone()
1153 } else {
1154 format!("{base_name}-{i}")
1155 };
1156 info.push((name, ip));
1157 }
1158 super::orchestrator::ProcessConfig {
1159 factory: &*entry.factory,
1160 info,
1161 ips,
1162 tag_registry: registry,
1163 machine_registry,
1164 }
1165 }
1166
1167 fn init_assertions_and_exploration(
1170 exploration_config: Option<&crate::chaos::exploration_glue::ExplorationConfig>,
1171 ) {
1172 crate::chaos::exploration_glue::init_assertion_region();
1173 let _ = exploration_config;
1174 #[cfg(feature = "exploration")]
1175 if let Some(config) = exploration_config {
1176 moonpool_explorer::set_rng_hooks(crate::sim::rng_call_count, |seed| {
1177 crate::sim::set_sim_seed(seed);
1178 crate::sim::reset_rng_call_count();
1179 });
1180 if let Err(e) = moonpool_explorer::init(config) {
1181 tracing::error!("Failed to initialize exploration: {}", e);
1182 }
1183 }
1184 }
1185
1186 #[cfg(feature = "exploration")]
1189 fn build_exploration_report(
1190 total_timelines: u64,
1191 total_fork_points: u64,
1192 total_bugs: u64,
1193 bug_recipes: Vec<super::report::BugRecipe>,
1194 converged: bool,
1195 per_seed_timelines: Vec<u64>,
1196 ) -> super::report::ExplorationReport {
1197 let final_stats = moonpool_explorer::exploration_stats();
1198 let coverage_bits = moonpool_explorer::explored_map_bits_set().unwrap_or(0);
1199 super::report::ExplorationReport {
1200 total_timelines,
1201 fork_points: total_fork_points,
1202 bugs_found: total_bugs,
1203 bug_recipes,
1204 energy_remaining: final_stats.as_ref().map_or(0, |s| s.global_energy),
1205 realloc_pool_remaining: final_stats.as_ref().map_or(0, |s| s.realloc_pool_remaining),
1206 coverage_bits,
1207 coverage_total: u32::try_from(moonpool_explorer::coverage::COVERAGE_MAP_SIZE * 8)
1208 .expect("coverage map size fits in u32"),
1209 sancov_edges_total: final_stats.as_ref().map_or(0, |s| s.sancov_edges_total),
1210 sancov_edges_covered: final_stats.as_ref().map_or(0, |s| s.sancov_edges_covered),
1211 converged,
1212 per_seed_timelines,
1213 }
1214 }
1215
1216 #[cfg(feature = "exploration")]
1220 fn accumulate_exploration_stats(
1221 seed: u64,
1222 per_seed_timelines: &mut Vec<u64>,
1223 total_timelines: &mut u64,
1224 total_fork_points: &mut u64,
1225 total_bugs: &mut u64,
1226 bug_recipes: &mut Vec<super::report::BugRecipe>,
1227 ) {
1228 if let Some(stats) = moonpool_explorer::exploration_stats() {
1229 per_seed_timelines.push(stats.total_timelines);
1230 *total_timelines += stats.total_timelines;
1231 *total_fork_points += stats.fork_points;
1232 *total_bugs += stats.bug_found;
1233 } else {
1234 per_seed_timelines.push(0);
1235 }
1236 if let Some(recipe) = moonpool_explorer::bug_recipe() {
1237 bug_recipes.push(super::report::BugRecipe { seed, recipe });
1238 }
1239 }
1240
1241 fn scan_assertion_slots(reached: &mut std::collections::HashSet<String>) -> usize {
1246 let slots = moonpool_assertions::assertion_read_all();
1247 for slot in &slots {
1248 if let Some(kind) = moonpool_assertions::AssertKind::from_u8(slot.kind)
1249 && matches!(
1250 kind,
1251 moonpool_assertions::AssertKind::Sometimes
1252 | moonpool_assertions::AssertKind::Reachable
1253 )
1254 {
1255 if slot.pass_count > 0 {
1256 reached.insert(slot.msg.clone());
1257 } else if !reached.contains(&slot.msg) {
1258 tracing::warn!(
1259 "UNREACHED slot: kind={:?} msg={:?} pass={} fail={}",
1260 kind,
1261 slot.msg,
1262 slot.pass_count,
1263 slot.fail_count
1264 );
1265 }
1266 }
1267 }
1268 slots
1269 .iter()
1270 .filter(|s| {
1271 moonpool_assertions::AssertKind::from_u8(s.kind).is_some_and(|k| {
1272 matches!(
1273 k,
1274 moonpool_assertions::AssertKind::Sometimes
1275 | moonpool_assertions::AssertKind::Reachable
1276 )
1277 })
1278 })
1279 .map(|s| s.msg.clone())
1280 .collect::<std::collections::HashSet<_>>()
1281 .len()
1282 }
1283
1284 fn empty_report() -> SimulationReport {
1286 SimulationReport {
1287 iterations: 0,
1288 successful_runs: 0,
1289 failed_runs: 0,
1290 metrics: SimulationMetrics::default(),
1291 individual_metrics: Vec::new(),
1292 seeds_used: Vec::new(),
1293 seeds_failing: Vec::new(),
1294 assertion_results: HashMap::new(),
1295 assertion_violations: Vec::new(),
1296 coverage_violations: Vec::new(),
1297 exploration: None,
1298 assertion_details: Vec::new(),
1299 bucket_summaries: Vec::new(),
1300 convergence_timeout: false,
1301 saturation: None,
1302 }
1303 }
1304
1305 #[instrument(skip_all)]
1306 pub fn run(mut self) -> SimulationReport {
1316 if self.entries.is_empty() {
1317 return Self::empty_report();
1318 }
1319
1320 struct SelectOverrideReset;
1326 impl Drop for SelectOverrideReset {
1327 fn drop(&mut self) {
1328 crate::sim::reset_select_rng();
1329 }
1330 }
1331 let _select_reset = SelectOverrideReset;
1332
1333 let layer = SimulationLayer::new();
1337 let (obs_handle, _obs_guard) = layer.install();
1338 for inv in self.invariants.drain(..) {
1339 obs_handle.register(inv);
1340 }
1341
1342 Self::init_assertions_and_exploration(self.exploration_config.as_ref());
1343
1344 let mut state = RunState::new(&self);
1345
1346 while state.iteration_manager.should_continue() {
1347 if let Some(report) = self.execute_iteration(&mut state, &obs_handle) {
1348 return report;
1349 }
1350 if state.converged {
1351 break;
1352 }
1353 }
1354
1355 Self::build_final_report(
1356 state.metrics_collector,
1357 &state.iteration_manager,
1358 self.exploration_config.as_ref(),
1359 &self.iteration_control,
1360 &FinalReportInputs {
1361 converged: state.converged,
1362 saturation: state.saturation,
1363 #[cfg(feature = "exploration")]
1364 total_exploration_timelines: state.total_exploration_timelines,
1365 #[cfg(feature = "exploration")]
1366 total_exploration_fork_points: state.total_exploration_fork_points,
1367 #[cfg(feature = "exploration")]
1368 total_exploration_bugs: state.total_exploration_bugs,
1369 #[cfg(feature = "exploration")]
1370 bug_recipes: state.bug_recipes,
1371 #[cfg(feature = "exploration")]
1372 per_seed_timelines: state.per_seed_timelines,
1373 },
1374 )
1375 }
1376
1377 fn execute_iteration(
1380 &mut self,
1381 state: &mut RunState,
1382 obs_handle: &SimulationLayerHandle,
1383 ) -> Option<SimulationReport> {
1384 let seed = state.iteration_manager.next_iteration();
1385 let iteration_count = state.iteration_manager.current_iteration();
1386
1387 self.prepare_iteration(obs_handle, seed, iteration_count);
1388
1389 let (orchestration_result, start_time) =
1390 self.run_orchestrator_for_iteration(state, obs_handle, seed, iteration_count);
1391
1392 if let Err(report) = self.handle_orchestration_result(
1393 state,
1394 orchestration_result,
1395 seed,
1396 iteration_count,
1397 start_time,
1398 ) {
1399 return Some(*report);
1400 }
1401
1402 self.finish_iteration(state, seed, iteration_count);
1403 None
1404 }
1405
1406 fn prepare_iteration(
1409 &mut self,
1410 obs_handle: &SimulationLayerHandle,
1411 seed: u64,
1412 iteration_count: usize,
1413 ) {
1414 if iteration_count > 1 {
1418 #[cfg(feature = "exploration")]
1419 if let Some(ref config) = self.exploration_config {
1420 moonpool_explorer::prepare_next_seed(config.global_energy);
1421 }
1422 crate::chaos::assertions::skip_next_assertion_reset();
1423 }
1424
1425 for hook in &mut self.before_iteration_hooks {
1426 hook();
1427 }
1428
1429 Self::reset_per_iteration_state(seed, self.swarm_operations, obs_handle);
1430 }
1431
1432 fn run_orchestrator_for_iteration(
1438 &mut self,
1439 state: &mut RunState,
1440 obs_handle: &SimulationLayerHandle,
1441 seed: u64,
1442 iteration_count: usize,
1443 ) -> (OrchestrationOutcome, Instant) {
1444 let ResolvedEntries {
1445 workloads,
1446 return_map,
1447 client_info,
1448 } = self.resolve_entries();
1449 state.pending_return_map = return_map;
1450
1451 let workload_info: Vec<(String, String)> = workloads
1452 .iter()
1453 .enumerate()
1454 .map(|(i, w)| (w.name().to_string(), format!("10.0.0.{}", i + 1)))
1455 .collect();
1456
1457 let process_config = self
1458 .process_entry
1459 .as_ref()
1460 .map(Self::resolve_process_config);
1461
1462 let sim = Self::build_sim_for_iteration(
1463 self.network_chaos,
1464 self.storage_chaos,
1465 self.buggify_knobs,
1466 seed,
1467 );
1468 let start_time = Instant::now();
1469 let attrition = match (self.attrition.as_ref(), self.attrition_mode) {
1473 (Some(base), ChaosMode::Swarm) => Some(base.swarm_for_seed()),
1474 (Some(base), ChaosMode::Random) => Some(base.clone()),
1475 (None, _) => None,
1476 };
1477 let fault_injectors =
1478 Self::collect_fault_injectors(&mut self.fault_injectors, attrition.as_ref());
1479 let outcome = Self::run_orchestrator_blocking(RunOrchestratorInputs {
1480 seed,
1481 iteration_count,
1482 workloads,
1483 workload_info,
1484 client_info,
1485 process_config,
1486 sim,
1487 fault_injectors,
1488 chaos_duration: self.chaos_duration,
1489 obs_handle: obs_handle.clone(),
1490 run_time_budget: self.run_time_budget,
1491 });
1492 (outcome, start_time)
1493 }
1494
1495 fn handle_orchestration_result(
1498 &mut self,
1499 state: &mut RunState,
1500 result: OrchestrationOutcome,
1501 seed: u64,
1502 iteration_count: usize,
1503 start_time: Instant,
1504 ) -> Result<(), Box<SimulationReport>> {
1505 let max_iterations = state
1506 .iteration_manager
1507 .max_iterations()
1508 .unwrap_or(iteration_count);
1509 let seeds_used_snapshot = state.iteration_manager.seeds_used().to_vec();
1510 match result {
1511 Ok(OrchestrateOutput {
1512 workloads: returned_workloads,
1513 fault_injectors: returned_injectors,
1514 results: all_results,
1515 metrics: sim_metrics,
1516 }) => {
1517 let return_map = std::mem::take(&mut state.pending_return_map);
1518 self.return_entries(returned_workloads, return_map);
1519 self.fault_injectors = returned_injectors;
1520 let wall_time = start_time.elapsed();
1521 state.metrics_collector.record_iteration(
1522 seed,
1523 wall_time,
1524 &all_results,
1525 crate::chaos::has_always_violations(),
1526 sim_metrics,
1527 );
1528 Self::log_slow_seed(seed, wall_time, self.seed_warning_timeout);
1529 Self::log_progress_milestone(
1530 state.progress_milestone,
1531 iteration_count,
1532 max_iterations,
1533 );
1534 Ok(())
1535 }
1536 Err((faulty_seeds_from_deadlock, failed_count)) => {
1537 state
1538 .metrics_collector
1539 .add_faulty_seeds(faulty_seeds_from_deadlock);
1540 state.metrics_collector.add_failed_runs(failed_count);
1541 let metrics_collector =
1542 std::mem::replace(&mut state.metrics_collector, MetricsCollector::new());
1543 Err(Box::new(Self::build_early_exit_report(
1544 metrics_collector,
1545 iteration_count,
1546 seeds_used_snapshot,
1547 )))
1548 }
1549 }
1550 }
1551
1552 fn finish_iteration(&self, state: &mut RunState, seed: u64, iteration_count: usize) {
1555 #[cfg(not(feature = "exploration"))]
1557 let _ = seed;
1558 #[cfg(feature = "exploration")]
1559 if self.exploration_config.is_some() {
1560 Self::accumulate_exploration_stats(
1561 seed,
1562 &mut state.per_seed_timelines,
1563 &mut state.total_exploration_timelines,
1564 &mut state.total_exploration_fork_points,
1565 &mut state.total_exploration_bugs,
1566 &mut state.bug_recipes,
1567 );
1568 }
1569
1570 let needs_assertion_scan = matches!(
1571 self.iteration_control,
1572 IterationControl::UntilCoverageStable { .. }
1573 );
1574 if needs_assertion_scan {
1575 let all_sometimes_count = Self::scan_assertion_slots(&mut state.reached_sometimes);
1576 state.converged = Self::check_convergence_or_plateau(ConvergenceState {
1577 iteration_control: &self.iteration_control,
1578 iteration_count,
1579 reached_sometimes: &state.reached_sometimes,
1580 all_sometimes_count,
1581 exploration_active: self.exploration_config.is_some(),
1582 prev_signal: &mut state.prev_signal,
1583 plateau_count: &mut state.plateau_count,
1584 saturation: &mut state.saturation,
1585 already_converged: state.converged,
1586 });
1587 }
1588
1589 crate::chaos::buggify_reset();
1590 }
1591
1592 fn build_final_report(
1594 metrics_collector: MetricsCollector,
1595 iteration_manager: &IterationManager,
1596 exploration_config: Option<&crate::chaos::exploration_glue::ExplorationConfig>,
1597 iteration_control: &IterationControl,
1598 inputs: &FinalReportInputs,
1599 ) -> SimulationReport {
1600 let converged = inputs.converged;
1601
1602 #[cfg(feature = "exploration")]
1607 let exploration_report = if exploration_config.is_some() {
1608 Some(Self::build_exploration_report(
1609 inputs.total_exploration_timelines,
1610 inputs.total_exploration_fork_points,
1611 inputs.total_exploration_bugs,
1612 inputs.bug_recipes.clone(),
1613 converged,
1614 inputs.per_seed_timelines.clone(),
1615 ))
1616 } else {
1617 None
1618 };
1619 #[cfg(not(feature = "exploration"))]
1620 let exploration_report: Option<super::report::ExplorationReport> = None;
1621
1622 let assertion_results = crate::chaos::assertion_results();
1624 let (assertion_violations, coverage_violations) =
1625 crate::chaos::validate_assertion_contracts();
1626 let raw_assertion_slots = moonpool_assertions::assertion_read_all();
1627 let raw_each_buckets = moonpool_assertions::each_bucket_read_all();
1628
1629 let did_exploration_cleanup = {
1633 #[cfg(feature = "exploration")]
1634 {
1635 if exploration_config.is_some() {
1636 moonpool_explorer::cleanup();
1637 true
1638 } else {
1639 false
1640 }
1641 }
1642 #[cfg(not(feature = "exploration"))]
1643 {
1644 let _ = exploration_config;
1645 false
1646 }
1647 };
1648 if !did_exploration_cleanup {
1649 crate::chaos::exploration_glue::cleanup_assertion_region();
1650 }
1651
1652 let assertion_details = build_assertion_details(&raw_assertion_slots);
1653 let bucket_summaries = build_bucket_summaries(&raw_each_buckets);
1654 let iteration_count = iteration_manager.current_iteration();
1655
1656 let convergence_timeout = matches!(
1658 iteration_control,
1659 IterationControl::UntilCoverageStable { .. }
1660 ) && !converged;
1661
1662 crate::chaos::buggify_reset();
1663
1664 metrics_collector.generate_report(GenerateReportInputs {
1665 iteration_count,
1666 seeds_used: iteration_manager.seeds_used().to_vec(),
1667 assertion_results,
1668 assertion_violations,
1669 coverage_violations,
1670 exploration: exploration_report,
1671 assertion_details,
1672 bucket_summaries,
1673 convergence_timeout,
1674 saturation: inputs.saturation.clone(),
1675 })
1676 }
1677}
1678
1679fn build_assertion_details(
1681 slots: &[moonpool_assertions::AssertionSlotSnapshot],
1682) -> Vec<super::report::AssertionDetail> {
1683 use super::report::{AssertionDetail, AssertionStatus};
1684 use moonpool_assertions::AssertKind;
1685
1686 slots
1687 .iter()
1688 .filter_map(|slot| {
1689 let kind = AssertKind::from_u8(slot.kind)?;
1690 let total = slot.pass_count.saturating_add(slot.fail_count);
1691
1692 if total == 0 && slot.frontier == 0 {
1694 return None;
1695 }
1696
1697 let status = match kind {
1698 AssertKind::Always
1699 | AssertKind::AlwaysOrUnreachable
1700 | AssertKind::NumericAlways => {
1701 if slot.fail_count > 0 {
1702 AssertionStatus::Fail
1703 } else {
1704 AssertionStatus::Pass
1705 }
1706 }
1707 AssertKind::Sometimes | AssertKind::NumericSometimes | AssertKind::Reachable => {
1708 if slot.pass_count > 0 {
1709 AssertionStatus::Pass
1710 } else {
1711 AssertionStatus::Miss
1712 }
1713 }
1714 AssertKind::Unreachable => {
1715 if slot.pass_count > 0 {
1716 AssertionStatus::Fail
1717 } else {
1718 AssertionStatus::Pass
1719 }
1720 }
1721 AssertKind::BooleanSometimesAll => {
1722 if slot.frontier > 0 {
1723 AssertionStatus::Pass
1724 } else {
1725 AssertionStatus::Miss
1726 }
1727 }
1728 };
1729
1730 Some(AssertionDetail {
1731 msg: slot.msg.clone(),
1732 kind,
1733 pass_count: slot.pass_count,
1734 fail_count: slot.fail_count,
1735 watermark: slot.watermark,
1736 frontier: slot.frontier,
1737 status,
1738 })
1739 })
1740 .collect()
1741}
1742
1743fn build_bucket_summaries(
1745 buckets: &[moonpool_assertions::EachBucket],
1746) -> Vec<super::report::BucketSiteSummary> {
1747 use super::report::BucketSiteSummary;
1748 use std::collections::HashMap;
1749
1750 let mut sites: HashMap<u32, BucketSiteSummary> = HashMap::new();
1751
1752 for bucket in buckets {
1753 let entry = sites
1754 .entry(bucket.site_hash)
1755 .or_insert_with(|| BucketSiteSummary {
1756 msg: bucket.msg_str().to_string(),
1757 buckets_discovered: 0,
1758 total_hits: 0,
1759 });
1760
1761 entry.buckets_discovered += 1;
1762 entry.total_hits += u64::from(bucket.pass_count);
1763 }
1764
1765 let mut summaries: Vec<_> = sites.into_values().collect();
1766 summaries.sort_by_key(|s| std::cmp::Reverse(s.total_hits));
1767 summaries
1768}
1769
1770#[cfg(test)]
1771mod tests {
1772 use super::*;
1773 use async_trait::async_trait;
1774 use moonpool_core::RandomProvider;
1775
1776 use crate::SimulationResult;
1777 use crate::runner::context::SimContext;
1778
1779 struct BasicWorkload;
1780
1781 #[async_trait]
1782 impl Workload for BasicWorkload {
1783 fn name(&self) -> &'static str {
1784 "test_workload"
1785 }
1786
1787 async fn run(&mut self, _ctx: &SimContext) -> SimulationResult<()> {
1788 Ok(())
1789 }
1790 }
1791
1792 #[test]
1793 fn test_simulation_builder_basic() {
1794 let report = SimulationBuilder::new()
1795 .workload(BasicWorkload)
1796 .set_iterations(3)
1797 .set_debug_seeds(vec![1, 2, 3])
1798 .run();
1799
1800 assert_eq!(report.iterations, 3);
1801 assert_eq!(report.successful_runs, 3);
1802 assert_eq!(report.failed_runs, 0);
1803 assert!((report.success_rate() - 100.0).abs() < f64::EPSILON);
1804 assert_eq!(report.seeds_used, vec![1, 2, 3]);
1805 }
1806
1807 struct FailingWorkload;
1808
1809 #[async_trait]
1810 impl Workload for FailingWorkload {
1811 fn name(&self) -> &'static str {
1812 "failing_workload"
1813 }
1814
1815 async fn run(&mut self, ctx: &SimContext) -> SimulationResult<()> {
1816 let random_num: u32 = ctx.random().random_range(0..100);
1818 if random_num.is_multiple_of(2) {
1819 return Err(crate::SimulationError::InvalidState(
1820 "Test failure".to_string(),
1821 ));
1822 }
1823 Ok(())
1824 }
1825 }
1826
1827 #[test]
1828 fn test_simulation_builder_with_failures() {
1829 let report = SimulationBuilder::new()
1834 .workload(FailingWorkload)
1835 .set_debug_seeds((1..=10).collect())
1836 .set_iterations(10)
1837 .run();
1838
1839 assert_eq!(report.iterations, 10);
1840 assert_eq!(
1841 report.successful_runs + report.failed_runs,
1842 10,
1843 "all iterations should be accounted for"
1844 );
1845 assert!(
1846 report.failed_runs > 0,
1847 "expected at least one failure across 10 seeds"
1848 );
1849 assert!(
1850 report.successful_runs > 0,
1851 "expected at least one success across 10 seeds"
1852 );
1853 }
1854}