1use crate::{Config, Sim};
15use core::fmt::Write as _;
16use core::time::Duration;
17use recon_core::{NodeId, Protocol, Time};
18use std::collections::BTreeMap;
19
20#[derive(Debug, Clone, PartialEq, Eq)]
27pub enum Step<C> {
28 Command { node: NodeId, cmd: C },
30 Crash(NodeId),
32 Restart(NodeId),
34 Suspend(NodeId),
36 Resume(NodeId),
38 CrashOnNextWrite(NodeId),
40 Sever(NodeId, NodeId),
42 Reconnect(NodeId, NodeId),
44 Partition(Vec<Vec<NodeId>>),
46 Heal,
48 BreakSession(NodeId, NodeId),
50}
51
52impl<C> Step<C> {
53 pub fn nodes(&self) -> Vec<NodeId> {
55 match self {
56 Step::Command { node, .. }
57 | Step::Crash(node)
58 | Step::Restart(node)
59 | Step::Suspend(node)
60 | Step::Resume(node)
61 | Step::CrashOnNextWrite(node) => vec![*node],
62 Step::Sever(a, b) | Step::Reconnect(a, b) | Step::BreakSession(a, b) => vec![*a, *b],
63 Step::Partition(groups) => groups.iter().flatten().copied().collect(),
64 Step::Heal => Vec::new(),
65 }
66 }
67
68 pub fn mentions(&self, node: NodeId) -> bool {
70 self.nodes().contains(&node)
71 }
72
73 fn apply<P>(&self, sim: &mut Sim<P>)
74 where
75 P: Protocol<Cmd = C>,
76 C: Clone,
77 P::Cmd: Clone,
78 P::Msg: Clone + PartialEq,
79 P::Ind: Clone,
80 P::Meta: Clone,
81 P::Entry: Clone,
82 {
83 match self {
84 Step::Command { node, cmd } => {
85 let _ = sim.command(*node, cmd.clone());
88 }
89 Step::Crash(n) => sim.crash(*n),
90 Step::Restart(n) => sim.restart(*n),
91 Step::Suspend(n) => sim.suspend(*n),
92 Step::Resume(n) => sim.resume(*n),
93 Step::CrashOnNextWrite(n) => sim.crash_on_next_write(*n),
94 Step::Sever(a, b) => sim.sever(*a, *b),
95 Step::Reconnect(a, b) => sim.reconnect(*a, *b),
96 Step::Partition(groups) => {
97 let borrowed: Vec<&[NodeId]> = groups.iter().map(|g| g.as_slice()).collect();
98 sim.partition(&borrowed);
99 }
100 Step::Heal => sim.heal(),
101 Step::BreakSession(a, b) => sim.break_session(*a, *b),
102 }
103 }
104}
105
106#[derive(Debug, Clone, PartialEq)]
112pub struct Scenario<C> {
113 pub config: Config,
116 pub nodes: Vec<NodeId>,
118 pub steps: Vec<(Duration, Step<C>)>,
122 pub horizon: Duration,
124}
125
126impl<C> Scenario<C> {
127 pub fn new(config: Config, nodes: impl IntoIterator<Item = NodeId>) -> Self {
129 Scenario {
130 config,
131 nodes: nodes.into_iter().collect(),
132 steps: Vec::new(),
133 horizon: Duration::ZERO,
134 }
135 }
136
137 pub fn at(mut self, at: Duration, step: Step<C>) -> Self {
145 if let Some((last, _)) = self.steps.last() {
146 assert!(at >= *last, "scenario steps must be in non-decreasing time order");
147 }
148 self.steps.push((at, step));
149 self
150 }
151
152 pub fn horizon(mut self, horizon: Duration) -> Self {
154 self.horizon = horizon;
155 self
156 }
157
158 pub fn end(&self) -> Duration {
160 self.steps.last().map(|(at, _)| *at).unwrap_or(Duration::ZERO).max(self.horizon)
161 }
162}
163
164impl<C: Clone> Scenario<C> {
165 pub(crate) fn without_steps(&self, drop: &[usize]) -> Self {
167 Scenario {
168 config: self.config.clone(),
169 nodes: self.nodes.clone(),
170 steps: self
171 .steps
172 .iter()
173 .enumerate()
174 .filter(|(i, _)| !drop.contains(i))
175 .map(|(_, s)| s.clone())
176 .collect(),
177 horizon: self.horizon,
178 }
179 .repaired()
180 }
181
182 pub(crate) fn without_node(&self, node: NodeId) -> Self {
186 let mut steps = Vec::new();
187 for (at, step) in &self.steps {
188 match step {
189 Step::Partition(groups) => {
190 let kept: Vec<Vec<NodeId>> = groups
191 .iter()
192 .map(|g| g.iter().copied().filter(|n| *n != node).collect::<Vec<_>>())
193 .filter(|g| !g.is_empty())
194 .collect();
195 if !kept.is_empty() {
196 steps.push((*at, Step::Partition(kept)));
197 }
198 }
199 s if s.mentions(node) => {}
200 s => steps.push((*at, s.clone())),
201 }
202 }
203 Scenario {
204 config: self.config.clone(),
205 nodes: self.nodes.iter().copied().filter(|n| *n != node).collect(),
206 steps,
207 horizon: self.horizon,
208 }
209 .repaired()
210 }
211
212 pub(crate) fn with_horizon(&self, horizon: Duration) -> Self {
214 Scenario {
215 config: self.config.clone(),
216 nodes: self.nodes.clone(),
217 steps: self.steps.iter().filter(|(at, _)| *at <= horizon).cloned().collect(),
218 horizon,
219 }
220 .repaired()
221 }
222
223 pub(crate) fn with_step(&self, index: usize, step: Step<C>) -> Self {
225 let mut out = self.clone();
226 out.steps[index].1 = step;
227 out.repaired()
228 }
229}
230
231impl<C: Clone> Scenario<C> {
232 fn repaired(mut self) -> Self {
243 #[derive(PartialEq, Clone, Copy)]
244 enum Liveness {
245 Running,
246 Suspended,
247 Crashed,
248 }
249 let mut state: BTreeMap<NodeId, Liveness> = BTreeMap::new();
250 let mut kept = Vec::with_capacity(self.steps.len());
251 for (at, step) in self.steps.drain(..) {
252 let live = |m: &BTreeMap<NodeId, Liveness>, n: &NodeId| {
253 *m.get(n).unwrap_or(&Liveness::Running)
254 };
255 let keep = match &step {
256 Step::Suspend(n) => live(&state, n) == Liveness::Running,
257 Step::Resume(n) => live(&state, n) == Liveness::Suspended,
258 Step::Restart(n) => live(&state, n) == Liveness::Crashed,
259 _ => true,
260 };
261 if !keep {
262 continue;
263 }
264 match &step {
265 Step::Suspend(n) => {
266 state.insert(*n, Liveness::Suspended);
267 }
268 Step::Resume(n) | Step::Restart(n) => {
269 state.insert(*n, Liveness::Running);
270 }
271 Step::Crash(n) => {
272 state.insert(*n, Liveness::Crashed);
273 }
274 _ => {}
275 }
276 kept.push((at, step));
277 }
278 self.steps = kept;
279 self
280 }
281
282 pub fn is_well_formed(&self) -> bool
287 where
288 C: PartialEq,
289 {
290 self.clone().repaired().steps == self.steps
291 }
292}
293
294impl<P> Sim<P>
295where
296 P: Protocol,
297 P::Cmd: Clone,
298 P::Msg: Clone + PartialEq,
299 P::Ind: Clone,
300 P::Meta: Clone,
301 P::Entry: Clone,
302{
303 pub fn run_scenario(
316 scenario: &Scenario<P::Cmd>,
317 build: impl FnOnce(Config, &[NodeId]) -> Sim<P>,
318 ) -> Sim<P> {
319 let mut sim = build(scenario.config.clone(), &scenario.nodes);
320 for (at, step) in &scenario.steps {
321 sim.run_until(Time::from_offset(*at));
322 step.apply(&mut sim);
323 }
324 sim.run_until(Time::from_offset(scenario.horizon));
325 sim
326 }
327}
328
329impl<C: core::fmt::Debug> Scenario<C> {
330 pub fn to_rust(&self, name: &str) -> String {
337 let mut s = String::new();
338 let _ = writeln!(s, "fn {name}() -> Scenario<Cmd> {{");
339 let _ = writeln!(s, " Scenario {{");
340 let _ = writeln!(s, " config: {},", render_config(&self.config, 8));
341 let _ = writeln!(s, " nodes: {},", render_nodes(&self.nodes));
342 if self.steps.is_empty() {
343 let _ = writeln!(s, " steps: vec![],");
344 } else {
345 let _ = writeln!(s, " steps: vec![");
346 for (at, step) in &self.steps {
347 let _ =
348 writeln!(s, " ({}, {}),", render_duration(*at), render_step(step));
349 }
350 let _ = writeln!(s, " ],");
351 }
352 let _ = writeln!(s, " horizon: {},", render_duration(self.horizon));
353 let _ = writeln!(s, " }}");
354 let _ = writeln!(s, "}}");
355 s
356 }
357}
358
359fn render_duration(d: Duration) -> String {
362 format!("Duration::from_nanos({})", d.as_nanos())
363}
364
365fn render_node(n: NodeId) -> String {
366 format!("NodeId({})", n.0)
367}
368
369fn render_nodes(nodes: &[NodeId]) -> String {
370 let inner: Vec<String> = nodes.iter().map(|n| render_node(*n)).collect();
371 format!("vec![{}]", inner.join(", "))
372}
373
374fn render_step<C: core::fmt::Debug>(step: &Step<C>) -> String {
375 match step {
376 Step::Command { node, cmd } => {
377 format!("Step::Command {{ node: {}, cmd: {cmd:?} }}", render_node(*node))
378 }
379 Step::Crash(n) => format!("Step::Crash({})", render_node(*n)),
380 Step::Restart(n) => format!("Step::Restart({})", render_node(*n)),
381 Step::Suspend(n) => format!("Step::Suspend({})", render_node(*n)),
382 Step::Resume(n) => format!("Step::Resume({})", render_node(*n)),
383 Step::CrashOnNextWrite(n) => format!("Step::CrashOnNextWrite({})", render_node(*n)),
384 Step::Sever(a, b) => {
385 format!("Step::Sever({}, {})", render_node(*a), render_node(*b))
386 }
387 Step::Reconnect(a, b) => {
388 format!("Step::Reconnect({}, {})", render_node(*a), render_node(*b))
389 }
390 Step::BreakSession(a, b) => {
391 format!("Step::BreakSession({}, {})", render_node(*a), render_node(*b))
392 }
393 Step::Heal => "Step::Heal".to_string(),
394 Step::Partition(groups) => {
395 let inner: Vec<String> = groups.iter().map(|g| render_nodes(g)).collect();
396 format!("Step::Partition(vec![{}])", inner.join(", "))
397 }
398 }
399}
400
401fn render_config(c: &Config, indent: usize) -> String {
405 let pad = " ".repeat(indent + 4);
406 let close = " ".repeat(indent);
407 let mut s = String::from("Config {\n");
408 let _ = writeln!(s, "{pad}seed: {},", c.seed);
409 let _ = writeln!(s, "{pad}loss: {:?},", c.loss);
410 let _ = writeln!(s, "{pad}duplication: {:?},", c.duplication);
411 let _ = writeln!(s, "{pad}reorder: {:?},", c.reorder);
412 let _ = writeln!(s, "{pad}latency_min: {},", render_duration(c.latency_min));
413 let _ = writeln!(s, "{pad}latency_max: {},", render_duration(c.latency_max));
414 let _ = writeln!(s, "{pad}reorder_delay: {},", render_duration(c.reorder_delay));
415 match c.synchronous {
416 None => {
417 let _ = writeln!(s, "{pad}synchronous: None,");
418 }
419 Some(d) => {
420 let _ = writeln!(s, "{pad}synchronous: Some({}),", render_duration(d));
421 }
422 }
423 let _ = writeln!(s, "{pad}sessions: {},", c.sessions);
424 let _ = writeln!(s, "{pad}reconnect_interval: {},", render_duration(c.reconnect_interval));
425 let _ = writeln!(s, "{pad}max_steps: {},", c.max_steps);
426 let _ = write!(s, "{close}}}");
427 s
428}