Skip to main content

recon_sim/
scenario.rs

1//! A run described as a value.
2//!
3//! Everything the simulator can be told to do imperatively — commands, crashes, partitions,
4//! session breaks — can also be held in a [`Scenario`]: a configuration with its seed, a
5//! membership, a list of timed [`Step`]s, and a horizon to run to. A description can be compared,
6//! printed, taken apart, and above all made *smaller*, which is what [`crate::shrink()`] does with
7//! it.
8//!
9//! This does not replace the imperative form and is not meant to. A test that provokes one named
10//! condition reads better as a sequence of calls. Scenarios are for the searching kind, where the
11//! failing input was discovered rather than chosen, and the useful next question is "how much of
12//! this mattered?"
13
14use crate::{Config, Sim};
15use core::fmt::Write as _;
16use core::time::Duration;
17use recon_core::{NodeId, Protocol, Time};
18use std::collections::BTreeMap;
19
20/// One thing done to a run from outside it.
21///
22/// The vocabulary is exactly the simulator's own mutators, so that anything a hand-written test
23/// can do to a run is something a description can say. The two opt-ins that are not faults —
24/// the codec check and session-event delivery — belong to how the run is *built* rather than to
25/// what happens during it, and so live in the constructor a scenario is run with.
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub enum Step<C> {
28    /// Hand a command to a process.
29    Command { node: NodeId, cmd: C },
30    /// Crash a process: it loses everything volatile.
31    Crash(NodeId),
32    /// Restart a crashed process, which takes its startup branch.
33    Restart(NodeId),
34    /// Suspend a process: it stops, keeping its state and everything addressed to it.
35    Suspend(NodeId),
36    /// Resume a suspended process.
37    Resume(NodeId),
38    /// Arm a process's next durable write to be the one it dies inside.
39    CrashOnNextWrite(NodeId),
40    /// Cut one pair off from each other, leaving every other pair alone.
41    Sever(NodeId, NodeId),
42    /// Restore one pair, leaving every other severing in place.
43    Reconnect(NodeId, NodeId),
44    /// Split the network into groups, replacing whatever was severed before.
45    Partition(Vec<Vec<NodeId>>),
46    /// Restore full connectivity, discarding every severing however it was made.
47    Heal,
48    /// End the session between a pair, losing an unknown suffix of what was in flight.
49    BreakSession(NodeId, NodeId),
50}
51
52impl<C> Step<C> {
53    /// Every process this step names, in the order it names them.
54    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    /// Whether this step mentions `node`.
69    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                // The identity is for a caller who wants to find the operation again; a scenario
86                // step is replayed rather than followed, so it has no use for one.
87                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/// A whole run, as data.
107///
108/// Executing one twice produces the same trace, by the determinism the simulator already
109/// guarantees: the configuration carries the seed, and nothing else in a run is drawn from
110/// anywhere but the generator that seed starts.
111#[derive(Debug, Clone, PartialEq)]
112pub struct Scenario<C> {
113    /// Network conditions and the seed. The seed lives here rather than beside it, because it is
114    /// what the simulator already treats as part of a configuration.
115    pub config: Config,
116    /// The processes in the run.
117    pub nodes: Vec<NodeId>,
118    /// What happens, and when — each time measured from the start of the run, in non-decreasing
119    /// order. Two steps at the same time happen in the order given, with nothing dispatched
120    /// between them.
121    pub steps: Vec<(Duration, Step<C>)>,
122    /// How long to run after the last step.
123    pub horizon: Duration,
124}
125
126impl<C> Scenario<C> {
127    /// An empty run over `nodes` with no steps and no horizon.
128    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    /// Add a step at `at`, measured from the start of the run.
138    ///
139    /// # Panics
140    ///
141    /// If `at` precedes the last step already added. A description is executed in the order it
142    /// is written, and the clock does not go backwards, so an out-of-order step would silently
143    /// happen at the wrong moment rather than where it reads.
144    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    /// Run until `horizon` after the start.
153    pub fn horizon(mut self, horizon: Duration) -> Self {
154        self.horizon = horizon;
155        self
156    }
157
158    /// The time the run ends: the horizon, or the last step if that is later.
159    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    /// The same scenario without the steps at `drop`, given as indices into [`Scenario::steps`].
166    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    /// The same scenario without `node`, and without every step that mentions it.
183    ///
184    /// A partition step keeps its other groups; a group emptied by the removal goes with it.
185    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    /// The same scenario run only to `horizon`, dropping any step scheduled after it.
213    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    /// The same scenario with step `index` replaced.
224    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    /// Drop the steps that a reduction has left dangling.
233    ///
234    /// A `Resume` belongs to a `Suspend` and a `Restart` to a `Crash`, and the simulator refuses
235    /// each without its partner — deliberately, since resuming a crashed process would be a
236    /// pause pretending to be a recovery. Deleting steps is exactly what a reduction does, so a
237    /// reduction that did not repair the pairing would spend most of its candidates on runs that
238    /// panic rather than on runs that answer the question.
239    ///
240    /// Repairing rather than rejecting: a candidate that has lost a `Suspend` is still a
241    /// candidate, it is just one where the process never stopped.
242    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    /// Whether every `Resume` has a `Suspend` and every `Restart` a `Crash`, in order.
283    ///
284    /// A scenario written by hand can be wrong; one produced by a reduction cannot, because every
285    /// reduction repairs. Exposed so a test can say which it is holding.
286    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    /// Execute a description, and return the run it produced.
304    ///
305    /// `build` is handed the scenario's configuration and membership and returns a simulator over
306    /// them — normally `Sim::new(config, nodes, make)`, plus whichever opt-ins the protocol needs
307    /// (`deliver_session_events`, `enable_codec_check`). It takes both rather than closing over
308    /// them because the shrinker will hand it *smaller* ones, and a builder that ignored its
309    /// arguments would quietly keep running the original.
310    ///
311    /// The clock is advanced to each step's time and the step applied, exactly as a hand-written
312    /// test would: `run_until(at)` then the call. Two steps at the same time are applied with
313    /// nothing dispatched between them, which is what makes "command, then break the session
314    /// carrying it" expressible.
315    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    /// Render as Rust that reconstructs this scenario, as a function named `name`.
331    ///
332    /// The end of a reduction should be something to paste, not something to transcribe. The
333    /// command is rendered with its `Debug`, which is valid Rust for the derived implementations
334    /// this repository's commands all use, provided their variants are in scope where the output
335    /// is pasted.
336    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
359/// `Duration::from_nanos(n)` — exact for every value the simulator can hold, where
360/// `from_millis` would round a latency drawn in nanoseconds into a different run.
361fn 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
401/// Every field, rather than the builder calls that would have produced them: `synchronous` and
402/// `sessions` overwrite the fault knobs, so builder order is not recoverable from a value and a
403/// rendering that guessed at it could produce a different run.
404fn 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}