Skip to main content

recon_sim/
sim.rs

1//! The deterministic execution environment.
2//!
3//! Runs a set of processes in one thread with a virtual clock and a seeded generator. It *is*
4//! the fair-loss link layer: messages may be lost, duplicated, delayed and reordered, and the
5//! protocols above are responsible for recovering from that.
6
7use crate::config::Config;
8use crate::narrate::{Render, render};
9use crate::trace::{DropReason, NotBegun, OpId, ProtoTrace, ProtoTraceEvent, Trace, TraceEvent};
10use core::time::Duration;
11use rand::{Rng, SeedableRng};
12use rand_chacha::ChaCha8Rng;
13use recon_core::error::CodecError;
14use recon_core::{
15    Cx, Effect, MemStore, NoNotes, NodeId, NoteSink, Position, ProtoCx, ProtoEffect, Protocol,
16    SessionEvent, Store, Time, TimerId, WriteKind,
17};
18use std::collections::{BTreeMap, BTreeSet};
19
20/// Round-trips one message through the wire codec, when codec checking is enabled.
21type CodecCheck<M> = fn(&M) -> Result<M, CodecError>;
22
23/// Something scheduled to happen at a point in virtual time.
24enum Scheduled<P: Protocol> {
25    Deliver {
26        from: NodeId,
27        to: NodeId,
28        msg: P::Msg,
29    },
30    Timer {
31        node: NodeId,
32        id: TimerId,
33    },
34    Command {
35        node: NodeId,
36        cmd: P::Cmd,
37        /// Names this operation in the trace. Carried so that whether it was handled or discarded,
38        /// the record says which operation it was.
39        op: OpId,
40    },
41    ScopeEvent {
42        node: NodeId,
43        scope: P::Scope,
44    },
45    /// Retry establishing every session that is not up. A deployed link keeps trying on its own
46    /// rather than waiting for the layers above to transmit, so the model does too.
47    Reconnect,
48}
49
50/// Whether a process is handling events, and if not, why.
51#[derive(Debug, Clone, Copy, PartialEq, Eq)]
52enum Liveness {
53    Running,
54    /// Stopped with its state preserved — a pause, not a failure.
55    Suspended,
56    /// Stopped having lost everything volatile.
57    Crashed,
58}
59
60/// A process in the run.
61struct Node<P> {
62    protocol: P,
63    liveness: Liveness,
64}
65
66/// A deterministic run of `P` across several processes.
67///
68/// Ordering is total and reproducible: the queue is keyed by `(time, sequence)`, so events at
69/// the same virtual instant are processed in the order they were scheduled, on every run.
70pub struct Sim<P: Protocol> {
71    now: Time,
72    rng: ChaCha8Rng,
73    config: Config,
74    seq: u64,
75    steps: u64,
76    queue: BTreeMap<(Time, u64), Scheduled<P>>,
77    nodes: BTreeMap<NodeId, Node<P>>,
78    /// Which pairs cannot reach each other, normalised so `(a, b)` and `(b, a)` are one entry.
79    ///
80    /// A set of pairs rather than a grouping, because a grouping makes reachability an equivalence
81    /// relation and real networks do not: `A` may reach `B` and `B` reach `C` while `A` cannot reach
82    /// `C`. [`Sim::partition`] is the special case in which the severed pairs happen to be exactly
83    /// those spanning two groups.
84    severed: BTreeSet<(NodeId, NodeId)>,
85    /// Everything that came due for a process while it was suspended, in the order it came due:
86    /// timers, deliveries carried by a live session, and scope events.
87    ///
88    /// A suspension is a stall, not a failure, so nothing addressed to a suspended process is
89    /// dropped — dropping a delivery while its session stays up would lose a message with no
90    /// `SessionEnded` to say so, which is the one thing `docs/conditional-guarantees.md` forbids
91    /// of every layer and therefore of the simulator too. These are re-dispatched by `resume`.
92    /// A crash destroys them instead — see `discard_pending_of`.
93    deferred: Vec<(NodeId, Scheduled<P>)>,
94    /// Rebuilds a process after a crash, since a crash loses volatile state.
95    make: Box<dyn FnMut(NodeId) -> P>,
96    /// The current epoch of the session between each pair, keyed by the pair in sorted order.
97    /// Absent means no session has been established yet.
98    sessions: BTreeMap<(NodeId, NodeId), u64>,
99    /// The next epoch to hand out for each pair, so epochs increase across re-establishment.
100    next_epoch: BTreeMap<(NodeId, NodeId), u64>,
101    /// When each pair's last session ended, so that none re-opens in the instant it closed.
102    ended_at: BTreeMap<(NodeId, NodeId), Time>,
103    /// The last time anything was delivered from one process to another, so that a message is
104    /// never delivered before one sent earlier the same way. This is what makes a session FIFO.
105    last_delivery: BTreeMap<(NodeId, NodeId), Time>,
106    /// Turns a session ending into whatever the protocol calls a scope. Absent unless the
107    /// protocol opted in, exactly as the codec check does.
108    session_scope: Option<fn(SessionEvent) -> P::Scope>,
109    /// What each process can read: everything it has written, durable or not. A protocol reads
110    /// its own writes back at once, which is what makes the interface synchronous.
111    storage: BTreeMap<NodeId, MemStore<P::Meta, P::Entry>>,
112    /// Processes whose next write is fatal — dying mid-`fsync`.
113    doomed: BTreeSet<NodeId>,
114    trace: ProtoTrace<P>,
115    /// One source of operation identities for the whole run, as `next_timer` is for timers.
116    next_op: u64,
117    /// What the current handler narrated, flushed into the trace when it returns. Empty, and never
118    /// written to, unless the run was asked to record notes.
119    notes: Vec<P::Note>,
120    /// Whether anything is listening. Off by default, so an ordinary run pays nothing.
121    record_notes: bool,
122    /// One source of timer identities for the whole run: two layers of one process must never
123    /// be handed the same handle, or each would accept the other's expiry.
124    next_timer: u64,
125    effects: Vec<ProtoEffect<P>>,
126    codec_check: Option<CodecCheck<P::Msg>>,
127    /// Renders each event as it is recorded. Absent unless the run was asked for it, exactly as
128    /// the codec check is.
129    render: Option<Render<P>>,
130}
131
132impl<P> Sim<P>
133where
134    P: Protocol,
135    P::Cmd: Clone,
136    P::Msg: Clone + PartialEq,
137    P::Ind: Clone,
138    P::Meta: Clone,
139    P::Entry: Clone,
140{
141    /// Build a run over `nodes`, constructing each process with `make`.
142    pub fn new(
143        config: Config,
144        nodes: &[NodeId],
145        mut make: impl FnMut(NodeId) -> P + 'static,
146    ) -> Self {
147        let rng = ChaCha8Rng::seed_from_u64(config.seed);
148        let mut map = BTreeMap::new();
149        for &id in nodes {
150            map.insert(id, Node { protocol: make(id), liveness: Liveness::Running });
151        }
152        let mut sim = Sim {
153            next_timer: 0,
154            now: Time::ZERO,
155            rng,
156            config,
157            seq: 0,
158            steps: 0,
159            queue: BTreeMap::new(),
160            nodes: map,
161            severed: BTreeSet::new(),
162            deferred: Vec::new(),
163            make: Box::new(make),
164            sessions: BTreeMap::new(),
165            next_epoch: BTreeMap::new(),
166            ended_at: BTreeMap::new(),
167            last_delivery: BTreeMap::new(),
168            session_scope: None,
169            storage: BTreeMap::new(),
170            doomed: BTreeSet::new(),
171            trace: Trace::default(),
172            next_op: 0,
173            notes: Vec::new(),
174            record_notes: false,
175            effects: Vec::new(),
176            codec_check: None,
177            render: None,
178        };
179        if sim.config.is_session_based() {
180            // The link starts trying immediately and keeps trying, so a session comes up as soon
181            // as one is possible rather than when something above happens to transmit.
182            sim.schedule(Time::ZERO, Scheduled::Reconnect);
183        }
184        // A first start, for every process: nothing has been written down, so this is the branch
185        // the book takes with ⟨ Init ⟩ rather than ⟨ Recovery ⟩.
186        let ids: Vec<NodeId> = sim.nodes.keys().copied().collect();
187        for node in ids {
188            sim.run_handler(node, |p, cx| p.on_init(cx));
189        }
190        sim
191    }
192
193    /// The upper bound on delivery, when the run is synchronous.
194    ///
195    /// A protocol whose correctness rests on this bound should be configured from it rather than
196    /// from a timeout that happens to work.
197    pub fn delivery_bound(&self) -> Option<Duration> {
198        self.config.delivery_bound()
199    }
200
201    /// The current virtual time.
202    pub fn now(&self) -> Time {
203        self.now
204    }
205
206    /// The record of what has happened so far.
207    pub fn trace(&self) -> &ProtoTrace<P> {
208        &self.trace
209    }
210
211    /// Record what protocols narrate, so the trace holds their decisions beside what happened.
212    ///
213    /// Off by default: a run pays nothing for an audience it does not have, and the protocol's own
214    /// code is the same either way — it calls `Cx::note` regardless, and the sink discards. That is
215    /// what makes narrating unable to change the run.
216    pub fn record_notes(&mut self) {
217        self.record_notes = true;
218    }
219
220    /// Record one event: render it if anything is listening, then keep it.
221    ///
222    /// Rendered *as* it is recorded rather than by a later walk over the trace, because a run that
223    /// fails to terminate is one of the things worth reading.
224    fn record(&mut self, event: ProtoTraceEvent<P>) {
225        if let Some(render) = self.render {
226            render(&event);
227        }
228        self.trace.push(event);
229    }
230
231    /// Borrow a process, for inspecting state a trace cannot show.
232    pub fn protocol(&self, node: NodeId) -> Option<&P> {
233        self.nodes.get(&node).map(|n| &n.protocol)
234    }
235
236    /// [`Sim::protocol`] for a process the test knows is running. Panics naming the process
237    /// otherwise, which is more use in a failure than `unwrap`'s line number.
238    pub fn at(&self, node: NodeId) -> &P {
239        self.protocol(node).unwrap_or_else(|| panic!("{node} is not running"))
240    }
241
242    /// Borrow what a process has written down. `None` if it has written nothing.
243    pub fn storage(&self, node: NodeId) -> Option<&MemStore<P::Meta, P::Entry>> {
244        self.storage.get(&node)
245    }
246
247    /// The processes in this run, in a stable order.
248    pub fn nodes(&self) -> impl Iterator<Item = NodeId> + '_ {
249        self.nodes.keys().copied()
250    }
251
252    /// Hand `cmd` to `node` at the current time, and take an identity naming that operation.
253    ///
254    /// The identity is a return value rather than a parameter, so a caller with no interest in it
255    /// carries on as before. It names the operation in the trace: [`Trace::invoked_at`] says when
256    /// the process handled it, and [`Trace::why_not_begun`] says why it never did.
257    pub fn command(&mut self, node: NodeId, cmd: P::Cmd) -> OpId {
258        self.command_at(node, Duration::ZERO, cmd)
259    }
260
261    /// Hand `cmd` to `node` after `after` has elapsed, and take an identity naming that operation.
262    pub fn command_at(&mut self, node: NodeId, after: Duration, cmd: P::Cmd) -> OpId {
263        let op = OpId(self.next_op);
264        self.next_op += 1;
265        self.schedule(self.now + after, Scheduled::Command { node, cmd, op });
266        op
267    }
268
269    /// Arm the next write by `node` to be the one it dies inside.
270    ///
271    /// Whether that write landed is decided by the seed, and the process cannot tell: what it
272    /// reads on recovering is the only evidence.
273    pub fn crash_on_next_write(&mut self, node: NodeId) {
274        self.doomed.insert(node);
275    }
276
277    /// Crash `node`: it stops handling events and loses everything volatile.
278    ///
279    /// Its protocol state is replaced with a freshly initialised one, and its pending timers and
280    /// anything held for it are discarded, so a restart resumes having forgotten what it
281    /// delivered. This is what a real process gets. For a stall that preserves state, use
282    /// [`Sim::suspend`].
283    pub fn crash(&mut self, node: NodeId) {
284        let fresh = (self.make)(node);
285        if let Some(n) = self.nodes.get_mut(&node) {
286            n.protocol = fresh;
287            n.liveness = Liveness::Crashed;
288        } else {
289            return;
290        }
291        self.discard_pending_of(node);
292        if self.config.is_session_based() {
293            self.end_sessions_of(node);
294        }
295        let at = self.now;
296        self.record(TraceEvent::Crashed { at, node });
297    }
298
299    /// Suspend `node`: it stops handling events but keeps its state, its timers, and everything
300    /// addressed to it while it is away.
301    ///
302    /// Not what a crash does. This is a *stall* — the process is descheduled and comes back
303    /// having missed nothing, because the timers, deliveries and scope events that came due
304    /// meanwhile were held rather than dropped. Losing them would be losing a message inside a
305    /// session that never ended, which is the one thing this model forbids of a layer.
306    ///
307    /// For unreachability use [`Sim::partition`], and for failure [`Sim::crash`]. Resume with
308    /// [`Sim::resume`]; [`Sim::restart`] is for a crashed process and re-runs its startup branch.
309    pub fn suspend(&mut self, node: NodeId) {
310        if let Some(n) = self.nodes.get_mut(&node) {
311            n.liveness = Liveness::Suspended;
312            let at = self.now;
313            self.record(TraceEvent::Suspended { at, node });
314        }
315    }
316
317    /// Resume a *suspended* `node`: everything held while it was away is dispatched now.
318    ///
319    /// No startup branch runs. Nothing was lost, so there is nothing to recover and nothing to
320    /// initialise — replaying `on_init` or `on_recovery` over intact volatile state would be
321    /// telling a process it restarted when it did not. That is what separates this from
322    /// [`Sim::restart`].
323    ///
324    /// Note what a resumed process is *not* told: that time passed. Its clock ran while it could
325    /// not read it, so anything that measures silence — a failure detector — comes back with
326    /// stale evidence and a timer due immediately. That is what a stall does to a real process,
327    /// and it is why the synchronous model excludes one.
328    pub fn resume(&mut self, node: NodeId) {
329        match self.nodes.get_mut(&node) {
330            Some(n) if n.liveness == Liveness::Suspended => n.liveness = Liveness::Running,
331            Some(_) => panic!("resume({node}): not suspended — restart() is for a crash"),
332            None => return,
333        }
334        let at = self.now;
335        self.record(TraceEvent::Resumed { at, node });
336        self.release_deferred(node);
337    }
338
339    /// Restart a *crashed* `node`, which takes its startup branch.
340    ///
341    /// A crash discarded its volatile state, its timers, and anything that was on its way, so
342    /// there is nothing held to release. What survived is in storage, and is handed back through
343    /// `on_recovery`. For a suspension use [`Sim::resume`].
344    pub fn restart(&mut self, node: NodeId) {
345        match self.nodes.get_mut(&node) {
346            Some(n) if n.liveness == Liveness::Crashed => n.liveness = Liveness::Running,
347            Some(_) => panic!("restart({node}): not crashed — resume() is for a suspension"),
348            None => return,
349        }
350        let at = self.now;
351        self.record(TraceEvent::Restarted { at, node });
352
353        // What survived, handed back as an event rather than through the constructor: the
354        // algorithms that need it re-announce their log and re-send what was pending, and those
355        // are effects, which a constructor cannot emit.
356        // Exactly one branch, as the book has it: something in storage means recovery, nothing
357        // means this incarnation is starting afresh and takes the first-start path instead.
358        // The disk did not survive either. A fault in its own right — storage is not guaranteed to
359        // come back — and the model already has the event that says so, since a restart with
360        // nothing written takes the first-start branch and records `had_state: false`.
361        //
362        // It exists for the audit in `scripts/check-durability-tests.sh`: a test asserting that
363        // something survived a restart should fail when nothing did, and one that passes anyway is
364        // reading the network. That is not hypothetical — it found three such tests, one of them
365        // guarding the composed durable record, and a leak in a fourth written the day before to
366        // close exactly this hole. Compile-time rather than a knob, because the audit asks the
367        // question of every existing test at once without editing any of them.
368        #[cfg(feature = "lose-storage-on-restart")]
369        self.storage.remove(&node);
370
371        let survived = self.storage.get(&node).map(|s| !s.is_empty()).unwrap_or(false);
372        self.record(TraceEvent::Recovered { at, node, had_state: survived });
373        if survived {
374            self.run_handler(node, |p, cx| p.on_recovery(cx));
375        } else {
376            self.run_handler(node, |p, cx| p.on_init(cx));
377        }
378    }
379
380    /// Re-dispatch everything held while `node` was suspended, in the order it came due.
381    fn release_deferred(&mut self, node: NodeId) {
382        let at = self.now;
383        let mut keep = Vec::new();
384        let mut due = Vec::new();
385        for (n, item) in self.deferred.drain(..) {
386            if n == node {
387                due.push(item);
388            } else {
389                keep.push((n, item));
390            }
391        }
392        self.deferred = keep;
393        for item in due {
394            self.schedule(at, item);
395        }
396    }
397
398    /// Whether `node` is currently stopped, for either reason.
399    pub fn is_stopped(&self, node: NodeId) -> bool {
400        self.nodes.get(&node).map(|n| n.liveness != Liveness::Running).unwrap_or(false)
401    }
402
403    /// Timers and undelivered messages are lost with the incarnation that was waiting for them,
404    /// so a crash takes them — including everything held while the process was suspended.
405    fn discard_pending_of(&mut self, node: NodeId) {
406        self.deferred.retain(|(n, _)| *n != node);
407        let doomed: Vec<(Time, u64)> = self
408            .queue
409            .iter()
410            .filter(|(_, s)| matches!(s, Scheduled::Timer { node: n, .. } if *n == node))
411            .map(|(k, _)| *k)
412            .collect();
413        for k in doomed {
414            self.queue.remove(&k);
415        }
416    }
417
418    /// Split the network into groups. Messages between groups are not delivered.
419    ///
420    /// The special case of [`Sim::sever`] in which the severed pairs are exactly those spanning two
421    /// groups — so reachability *is* transitive here, and a process in a group reaches every other
422    /// member. That is the easy case, and it is not the only one; see `sever`.
423    ///
424    /// Replaces whatever was severed before, so calling this after `sever` discards that severing.
425    pub fn partition(&mut self, groups: &[&[NodeId]]) {
426        let groups: Vec<BTreeSet<NodeId>> =
427            groups.iter().map(|g| g.iter().copied().collect()).collect();
428        let members: Vec<NodeId> = groups.iter().flatten().copied().collect();
429        self.severed.clear();
430        for (i, a) in members.iter().enumerate() {
431            for b in members.iter().skip(i + 1) {
432                if !groups.iter().any(|g| g.contains(a) && g.contains(b)) {
433                    self.severed.insert(pair(*a, *b));
434                }
435            }
436        }
437        if self.config.is_session_based() {
438            self.end_severed_sessions();
439        }
440    }
441
442    /// Cut `a` and `b` off from each other, in both directions, leaving every other pair alone.
443    ///
444    /// This is what a grouping cannot express. Severing one pair of three processes leaves a
445    /// **bridge**: `A` reaches `B` and `B` reaches `C`, but `A` does not reach `C`. All three are
446    /// correct, none of them is wrong about what it can see, and there is no group any of them
447    /// belongs to — which is the case every layer above that depends on processes agreeing about
448    /// who is reachable has never been asked about.
449    ///
450    /// Severing is symmetric. A link that works one way and not the other is a different fault and
451    /// a harder question for the session model, which treats a session as a property of a pair; see
452    /// this change's `design.md`.
453    pub fn sever(&mut self, a: NodeId, b: NodeId) {
454        if a == b {
455            return;
456        }
457        self.severed.insert(pair(a, b));
458        if self.config.is_session_based() {
459            self.end_severed_sessions();
460        }
461    }
462
463    /// Restore connectivity between `a` and `b`, leaving every other severing in place.
464    pub fn reconnect(&mut self, a: NodeId, b: NodeId) {
465        self.severed.remove(&pair(a, b));
466    }
467
468    /// Whether `a` and `b` can currently reach each other.
469    ///
470    /// For asserting the topology a test built rather than assuming it: severing one pair of *four*
471    /// processes is not a bridge, and a test that thinks it is would be testing nothing.
472    pub fn reachable(&self, a: NodeId, b: NodeId) -> bool {
473        self.connected(a, b)
474    }
475
476    /// Restore full connectivity, discarding every severing however it was made.
477    pub fn heal(&mut self) {
478        self.severed.clear();
479    }
480
481    /// Process every event scheduled at or before `until`.
482    pub fn run_until(&mut self, until: Time) {
483        while self.steps < self.config.max_steps {
484            let Some((&key, _)) = self.queue.iter().next() else { break };
485            if key.0 > until {
486                break;
487            }
488            let item = self.queue.remove(&key).expect("key just observed");
489            self.now = key.0;
490            self.steps += 1;
491            self.dispatch(item);
492        }
493        if self.now < until {
494            self.now = until;
495        }
496    }
497
498    /// Process every event scheduled within `d` of now.
499    /// Dispatch everything scheduled for the current instant, and nothing later. The clock does not
500    /// move.
501    ///
502    /// For sequencing a test by *events* rather than by durations: a command is scheduled, not run,
503    /// so `command(...)` followed by `break_session(...)` breaks a session with nothing in flight.
504    /// `step_now()` between them runs the handler, whose sends then sit in the queue at their
505    /// latency — in flight, and the break finds them. The older idiom, `run_for(1 ms)` with a
506    /// comment, depends on the latency being longer than the millisecond.
507    pub fn step_now(&mut self) {
508        let now = self.now;
509        while self.steps < self.config.max_steps {
510            let Some((&key, _)) = self.queue.iter().next() else { break };
511            if key.0 > now {
512                break;
513            }
514            let item = self.queue.remove(&key).expect("key just observed");
515            self.steps += 1;
516            self.dispatch(item);
517        }
518    }
519
520    /// Dispatch the next scheduled event, moving the clock to it. `false` when there is nothing
521    /// left to dispatch, or the step budget is spent.
522    ///
523    /// For a test searching for a state that one event creates and the next may destroy — "exactly
524    /// one process has decided" — stepping by event cannot skip it, where `run_for(1 ms)` can when
525    /// two events fall inside the millisecond.
526    pub fn step(&mut self) -> bool {
527        if self.steps >= self.config.max_steps {
528            return false;
529        }
530        let Some((&key, _)) = self.queue.iter().next() else { return false };
531        let item = self.queue.remove(&key).expect("key just observed");
532        self.now = key.0;
533        self.steps += 1;
534        self.dispatch(item);
535        true
536    }
537
538    pub fn run_for(&mut self, d: Duration) {
539        self.run_until(self.now + d);
540    }
541
542    // ---------------------------------------------------------------- internals
543
544    fn schedule(&mut self, at: Time, item: Scheduled<P>) {
545        let key = (at, self.seq);
546        self.seq += 1;
547        self.queue.insert(key, item);
548    }
549
550    fn dispatch(&mut self, item: Scheduled<P>) {
551        match item {
552            Scheduled::Command { node, cmd, op } => {
553                // Discarded rather than held when the process is not running, and recorded either
554                // way. A command is a call from the layer above, on this process — a stalled
555                // process's layer above is stalled with it, so there is nothing to delay. What was
556                // wrong before was the silence, not the discarding.
557                if let Some(why) = self.why_not_begun(node) {
558                    let at = self.now;
559                    self.record(TraceEvent::NotInvoked { at, node, op, cmd, why });
560                    return;
561                }
562                let at = self.now;
563                self.record(TraceEvent::Invoked { at, node, op, cmd: cmd.clone() });
564                self.run_handler(node, |p, cx| p.on_cmd(cmd, cx));
565            }
566            Scheduled::Reconnect => {
567                self.reconnect_sweep();
568                let at = self.now + self.config.reconnect_interval;
569                self.schedule(at, Scheduled::Reconnect);
570            }
571            Scheduled::ScopeEvent { node, scope } => {
572                if self.suspended(node) {
573                    // A local notification, not a network message: held for the same reason a
574                    // timer is. A stalled process has not stopped existing.
575                    self.deferred.push((node, Scheduled::ScopeEvent { node, scope }));
576                    return;
577                }
578                if self.stopped(node) {
579                    return;
580                }
581                self.run_handler(node, |p, cx| p.on_scope_event(scope, cx));
582            }
583            Scheduled::Timer { node, id } => {
584                if self.suspended(node) {
585                    // Held, not dropped: the process still exists and will want this.
586                    self.deferred.push((node, Scheduled::Timer { node, id }));
587                    return;
588                }
589                if self.stopped(node) {
590                    return;
591                }
592                let at = self.now;
593                self.record(TraceEvent::TimerFired { at, node, id });
594                self.run_handler(node, |p, cx| p.on_timer(id, cx));
595            }
596            Scheduled::Deliver { from, to, msg } => {
597                if self.suspended(to) {
598                    // Held. Dropping it would lose a message inside a session that is still up,
599                    // with no `SessionEnded` raised — the model's own invariant, broken by the
600                    // model. The recipient sees it when it resumes, as it would after a stall.
601                    self.deferred.push((to, Scheduled::Deliver { from, to, msg }));
602                    return;
603                }
604                if self.stopped(to) {
605                    let at = self.now;
606                    self.record(TraceEvent::Dropped {
607                        at,
608                        from,
609                        to,
610                        msg,
611                        reason: DropReason::RecipientCrashed,
612                    });
613                    return;
614                }
615                let msg = match self.check_codec(&msg) {
616                    Ok(m) => m,
617                    Err(e) => panic!("codec check failed for a message from {from} to {to}: {e}"),
618                };
619                let at = self.now;
620                self.record(TraceEvent::Delivered { at, from, to, msg: msg.clone() });
621                self.run_handler(to, |p, cx| p.on_msg(from, msg, cx));
622            }
623        }
624    }
625
626    /// Stopped with its state intact — a stall, from which everything held is replayed.
627    fn suspended(&self, node: NodeId) -> bool {
628        self.nodes.get(&node).map(|n| n.liveness == Liveness::Suspended).unwrap_or(false)
629    }
630
631    /// Stopped having lost everything volatile, or never here at all.
632    ///
633    /// Distinct from [`Sim::stopped`], and the distinction matters: what is owed to a crashed
634    /// process is nothing, and what is owed to a suspended one is everything, later.
635    fn crashed(&self, node: NodeId) -> bool {
636        self.nodes.get(&node).map(|n| n.liveness == Liveness::Crashed).unwrap_or(true)
637    }
638
639    /// Not handling events, for either reason.
640    /// Why an operation given to `node` cannot begin, or `None` if it can.
641    fn why_not_begun(&self, node: NodeId) -> Option<NotBegun> {
642        match self.nodes.get(&node) {
643            None => Some(NotBegun::NotAProcess),
644            Some(n) => match n.liveness {
645                Liveness::Running => None,
646                Liveness::Suspended => Some(NotBegun::Stalled),
647                Liveness::Crashed => Some(NotBegun::Crashed),
648            },
649        }
650    }
651
652    fn stopped(&self, node: NodeId) -> bool {
653        self.nodes.get(&node).map(|n| n.liveness != Liveness::Running).unwrap_or(true)
654    }
655
656    /// Run one handler and interpret everything it emits.
657    fn run_handler(&mut self, node: NodeId, f: impl FnOnce(&mut P, &mut ProtoCx<'_, P>)) {
658        let mut effects = core::mem::take(&mut self.effects);
659        effects.clear();
660
661        // Drawn before the handler runs, so the store need not borrow the generator the context
662        // already holds.
663        // Armed until a write actually happens: a handler that writes nothing is not the one.
664        let doomed = self.doomed.contains(&node);
665        let keep = doomed && self.rng.random::<bool>();
666        let mut writes: Vec<WriteKind> = Vec::new();
667        let mut died = false;
668        let mut notes = core::mem::take(&mut self.notes);
669        notes.clear();
670
671        {
672            let Some(n) = self.nodes.get_mut(&node) else {
673                self.effects = effects;
674                self.notes = notes;
675                return;
676            };
677            let inner = self.storage.entry(node).or_default();
678            let mut store =
679                FaultyStore { inner, writes: &mut writes, doomed, keep, died: &mut died };
680            // The protocol's code is the same either way — it calls `cx.note` regardless — which is
681            // what makes narrating unable to change the run.
682            let mut discard = NoNotes;
683            let listener: &mut dyn NoteSink<P::Note> =
684                if self.record_notes { &mut notes } else { &mut discard };
685            let mut cx = Cx::new(
686                &mut effects,
687                self.now,
688                &mut self.rng,
689                &mut store,
690                &mut self.next_timer,
691                listener,
692            );
693            f(&mut n.protocol, &mut cx);
694        }
695
696        let at = self.now;
697        // Before the writes and the effects: a note marks the decision, and those are what it led
698        // to. A handler that narrated and then wrote reads in that order.
699        for note in notes.drain(..) {
700            self.record(TraceEvent::Said { at, node, note });
701        }
702        self.notes = notes;
703        if !writes.is_empty() {
704            self.doomed.remove(&node);
705        }
706        for kind in writes {
707            self.record(TraceEvent::Wrote { at, node, kind });
708        }
709
710        if died {
711            // Everything the handler went on to do is discarded — a crash loses volatile state
712            // anyway, so nothing decided on the strength of that write can escape.
713            effects.clear();
714            self.effects = effects;
715            self.record(TraceEvent::DiedWriting { at, node });
716            self.crash(node);
717            return;
718        }
719
720        for effect in effects.drain(..) {
721            match effect {
722                Effect::Send { to, msg } => self.transmit(node, to, msg),
723                Effect::Indicate(ind) => {
724                    let at = self.now;
725                    self.record(TraceEvent::Indicated { at, node, ind });
726                }
727                Effect::SetTimer { after, id } => {
728                    self.schedule(self.now + after, Scheduled::Timer { node, id });
729                }
730            }
731        }
732
733        self.effects = effects;
734    }
735
736    /// Apply the network model to one outgoing message.
737    fn transmit(&mut self, from: NodeId, to: NodeId, msg: P::Msg) {
738        let at = self.now;
739
740        // A message a process addressed to itself crosses no wire. A driver is what turns a request
741        // to send into a packet, and a packet addressed to the process it came from is a hand-off
742        // between two roles of one state machine — `multi_paxos_synod` holds both a leader and an
743        // acceptor, and the leader reaching its own acceptor is a call. So it is delivered at this
744        // instant, with no latency and no fault applied, and recorded as its own kind of event
745        // rather than as a send.
746        //
747        // This narrows what can be lost and never widens it: a hand-off cannot be delayed, dropped,
748        // duplicated, reordered, or discarded by a scope ending, and there is no session between a
749        // process and itself to end. It is *delivered* rather than executed here, as every other
750        // effect is, so no handler ever runs inside another handler's effect loop.
751        if from == to {
752            if self.stopped(from) {
753                return;
754            }
755            self.record(TraceEvent::HandedToSelf { at, node: from, msg: msg.clone() });
756            self.schedule(at, Scheduled::Deliver { from, to, msg });
757            return;
758        }
759
760        self.record(TraceEvent::Sent { at, from, to, msg: msg.clone() });
761
762        if self.stopped(from) {
763            self.record(TraceEvent::Dropped {
764                at,
765                from,
766                to,
767                msg,
768                reason: DropReason::SenderCrashed,
769            });
770            return;
771        }
772
773        if !self.connected(from, to) {
774            self.record(TraceEvent::Dropped { at, from, to, msg, reason: DropReason::Partitioned });
775            return;
776        }
777
778        if self.config.is_session_based() {
779            if !self.ensure_session(from, to) {
780                // `connected` was checked above, so a refusal here is either a crashed recipient
781                // or the instant of an ending — not a partition, and the trace must not say so.
782                let reason = if self.crashed(to) {
783                    DropReason::RecipientCrashed
784                } else {
785                    DropReason::NoSession
786                };
787                self.record(TraceEvent::Dropped { at, from, to, msg, reason });
788                return;
789            }
790            let deliver_at = self.session_delivery_time(from, to);
791            self.schedule(deliver_at, Scheduled::Deliver { from, to, msg });
792            return;
793        }
794
795        let synchronous = self.config.is_synchronous();
796
797        if !synchronous && self.config.loss > 0.0 && self.rng.random::<f64>() < self.config.loss {
798            self.record(TraceEvent::Dropped { at, from, to, msg, reason: DropReason::Lost });
799            return;
800        }
801
802        let delay = self.draw_delay(from, to, &msg);
803        self.schedule(at + delay, Scheduled::Deliver { from, to, msg: msg.clone() });
804
805        if !synchronous
806            && self.config.duplication > 0.0
807            && self.rng.random::<f64>() < self.config.duplication
808        {
809            self.record(TraceEvent::Duplicated { at, from, to, msg: msg.clone() });
810            let second = self.draw_delay(from, to, &msg);
811            self.schedule(at + second, Scheduled::Deliver { from, to, msg });
812        }
813    }
814
815    fn draw_delay(&mut self, from: NodeId, to: NodeId, msg: &P::Msg) -> Duration {
816        let lo = self.config.latency_min.as_nanos() as u64;
817        let hi = self.config.latency_max.as_nanos() as u64;
818        let base = if hi > lo { self.rng.random_range(lo..=hi) } else { lo };
819
820        let mut delay = Duration::from_nanos(base);
821        if let Some(bound) = self.config.synchronous {
822            // A reordering spike would exceed the bound, which is the one thing this mode
823            // promises not to do.
824            return delay.min(bound);
825        }
826        if self.config.reorder > 0.0 && self.rng.random::<f64>() < self.config.reorder {
827            let at = self.now;
828            self.record(TraceEvent::Reordered { at, from, to, msg: msg.clone() });
829            delay += self.config.reorder_delay;
830        }
831        delay
832    }
833
834    fn connected(&self, a: NodeId, b: NodeId) -> bool {
835        a == b || !self.severed.contains(&pair(a, b))
836    }
837
838    fn check_codec(&self, msg: &P::Msg) -> Result<P::Msg, CodecError> {
839        match self.codec_check {
840            None => Ok(msg.clone()),
841            Some(f) => f(msg),
842        }
843    }
844}
845
846impl<P> Sim<P>
847where
848    P: Protocol,
849    P::Cmd: Clone,
850    P::Msg: Clone + PartialEq + serde::Serialize + serde::de::DeserializeOwned,
851    P::Ind: Clone,
852    P::Meta: Clone,
853    P::Entry: Clone,
854{
855    /// Round-trip every delivered message through the wire codec.
856    ///
857    /// Off by default: the simulator moves typed values, so a codec defect cannot be mistaken
858    /// for a protocol defect. Turn it on to check that messages actually survive encoding,
859    /// without paying for it on every run.
860    pub fn enable_codec_check(&mut self) {
861        self.codec_check = Some(crate::codec::round_trip);
862    }
863}
864
865impl<P> Sim<P>
866where
867    P: Protocol,
868    P::Msg: Clone + PartialEq + core::fmt::Debug,
869    P::Ind: Clone + core::fmt::Debug,
870    P::Note: core::fmt::Debug,
871    P::Cmd: Clone + core::fmt::Debug,
872    P::Meta: Clone,
873    P::Entry: Clone,
874{
875    /// Emit every recorded event to whatever `tracing` subscriber is installed, as it is recorded.
876    ///
877    /// Off by default, like the codec check: a run pays nothing for an audience it does not have.
878    /// Turning it on does not change the run — the events are the same ones the trace already
879    /// holds, in the same order, and nothing a protocol can observe is affected.
880    ///
881    /// Pair it with [`Sim::record_notes`] to see what the protocols *said* as well as what happened
882    /// to them; without it the rendering shows the run, which is what the trace showed before this
883    /// existed.
884    pub fn enable_tracing(&mut self) {
885        self.render = Some(render::<P>);
886    }
887}
888
889fn pair(a: NodeId, b: NodeId) -> (NodeId, NodeId) {
890    if a <= b { (a, b) } else { (b, a) }
891}
892
893impl<P> Sim<P>
894where
895    P: Protocol,
896    P::Cmd: Clone,
897    P::Msg: Clone + PartialEq,
898    P::Ind: Clone,
899    P::Meta: Clone,
900    P::Entry: Clone,
901{
902    /// The current epoch of the session between `a` and `b`, if one is established.
903    pub fn session_epoch(&self, a: NodeId, b: NodeId) -> Option<u64> {
904        self.sessions.get(&pair(a, b)).copied()
905    }
906
907    /// End the session between `a` and `b`, discarding an unknown suffix of what was in flight.
908    ///
909    /// A new session opens at a higher epoch the next time either sends and the pair is able to
910    /// communicate.
911    pub fn break_session(&mut self, a: NodeId, b: NodeId) {
912        self.end_session(a, b, DropReason::SessionEnded);
913    }
914
915    /// Whether a session currently exists between `a` and `b`.
916    pub fn has_session(&self, a: NodeId, b: NodeId) -> bool {
917        self.sessions.contains_key(&pair(a, b))
918    }
919
920    fn end_session(&mut self, a: NodeId, b: NodeId, reason: DropReason) {
921        let key = pair(a, b);
922        let Some(epoch) = self.sessions.remove(&key) else {
923            return;
924        };
925        let at = self.now;
926
927        // Discard an unknown suffix of what was in flight, in either direction. The cut is drawn
928        // from the run's generator, so it varies with the seed and can be everything or nothing.
929        let mut inflight: Vec<(Time, u64)> = self
930            .queue
931            .iter()
932            .filter(|(_, s)| {
933                matches!(s, Scheduled::Deliver { from, to, .. } if pair(*from, *to) == key)
934            })
935            .map(|(k, _)| *k)
936            .collect();
937        inflight.sort();
938        let keep = self.rng.random_range(0..=inflight.len());
939        for k in inflight.iter().copied().skip(keep) {
940            if let Some(Scheduled::Deliver { from, to, msg }) = self.queue.remove(&k) {
941                self.record(TraceEvent::SuffixLost { at, from, to, msg });
942            }
943        }
944
945        // What survived is flushed now, before the ending is announced. A transport delivers
946        // nothing on a connection after it has surfaced the close, and a scope boundary that
947        // arrivals can trail is not a boundary: the layer above resends on `Established`, so a
948        // straggler from the old epoch would arrive behind the new epoch's traffic under an
949        // identifier nothing distinguishes. Re-scheduled in their existing order, so the session
950        // is FIFO right up to its last message.
951        for k in inflight.into_iter().take(keep) {
952            if let Some(item) = self.queue.remove(&k) {
953                self.schedule(at, item);
954            }
955        }
956
957        // Ordering restarts with the next session, and the next one cannot be this instant.
958        self.last_delivery.remove(&(key.0, key.1));
959        self.last_delivery.remove(&(key.1, key.0));
960        self.ended_at.insert(key, at);
961
962        self.record(TraceEvent::SessionEnded { at, a: key.0, b: key.1, epoch, reason });
963
964        // Both endpoints are told, if they are alive to hear it. The epoch named is the one that
965        // ended: at the moment of failure the next is not a fact, and may never become one.
966        if let Some(f) = self.session_scope {
967            for (node, peer) in [(key.0, key.1), (key.1, key.0)] {
968                if !self.crashed(node) {
969                    let scope = f(SessionEvent::Ended { peer, epoch });
970                    self.schedule(at, Scheduled::ScopeEvent { node, scope });
971                }
972            }
973        }
974    }
975
976    /// End every session involving `node`.
977    fn end_sessions_of(&mut self, node: NodeId) {
978        let peers: Vec<NodeId> = self
979            .sessions
980            .keys()
981            .filter_map(|(a, b)| {
982                if *a == node {
983                    Some(*b)
984                } else if *b == node {
985                    Some(*a)
986                } else {
987                    None
988                }
989            })
990            .collect();
991        for peer in peers {
992            self.end_session(node, peer, DropReason::SessionEnded);
993        }
994    }
995
996    /// End every session that the current partitioning has severed.
997    fn end_severed_sessions(&mut self) {
998        let severed: Vec<(NodeId, NodeId)> =
999            self.sessions.keys().copied().filter(|(a, b)| !self.connected(*a, *b)).collect();
1000        for (a, b) in severed {
1001            self.end_session(a, b, DropReason::Partitioned);
1002        }
1003    }
1004
1005    /// Establish a session if the pair can communicate and none exists.
1006    ///
1007    /// A process needs no session with itself, and giving it one would announce `Established`
1008    /// twice per node — the pair loop visits `(a, b)` and `(b, a)` — so a layer that resends on
1009    /// re-establishment would resend to itself twice per round.
1010    fn ensure_session(&mut self, a: NodeId, b: NodeId) -> bool {
1011        if a == b {
1012            return true;
1013        }
1014        let key = pair(a, b);
1015        if self.sessions.contains_key(&key) {
1016            return true;
1017        }
1018        // Strictly crashed: a suspended process still exists, and the scope event it is owed is
1019        // held for it rather than skipped.
1020        if !self.connected(a, b) || self.crashed(a) || self.crashed(b) {
1021            return false;
1022        }
1023        // Not in the instant the last one ended. Everything that ending owes the two endpoints —
1024        // the flushed prefix, then `Ended` — is scheduled at that instant, and a successor
1025        // established among them would reach a layer as `Established` before the `Ended` it
1026        // replaces. Reconnection takes time; the sweep opens the next one.
1027        if self.ended_at.get(&key) == Some(&self.now) {
1028            return false;
1029        }
1030        // Epochs are per pair and only ever increase, so a re-established session is
1031        // distinguishable from the one it replaces.
1032        let epoch = self.next_epoch.get(&key).copied().unwrap_or(1);
1033        self.next_epoch.insert(key, epoch + 1);
1034        self.sessions.insert(key, epoch);
1035        let at = self.now;
1036        self.record(TraceEvent::SessionOpened { at, a: key.0, b: key.1, epoch });
1037
1038        // Both endpoints are told. This is the actionable event: the peer is reachable, so
1039        // anything sent in response arrives.
1040        if let Some(f) = self.session_scope {
1041            for (node, peer) in [(key.0, key.1), (key.1, key.0)] {
1042                if !self.crashed(node) {
1043                    let scope = f(SessionEvent::Established { peer, epoch });
1044                    self.schedule(at, Scheduled::ScopeEvent { node, scope });
1045                }
1046            }
1047        }
1048        true
1049    }
1050
1051    /// Try to establish every session that is not up.
1052    fn reconnect_sweep(&mut self) {
1053        let nodes: Vec<NodeId> = self.nodes.keys().copied().collect();
1054        for (i, a) in nodes.iter().enumerate() {
1055            for b in nodes.iter().skip(i + 1) {
1056                self.ensure_session(*a, *b);
1057            }
1058        }
1059    }
1060
1061    /// Delivery time under a session: never before something sent earlier the same way.
1062    fn session_delivery_time(&mut self, from: NodeId, to: NodeId) -> Time {
1063        let lo = self.config.latency_min.as_nanos() as u64;
1064        let hi = self.config.latency_max.as_nanos() as u64;
1065        let base = if hi > lo { self.rng.random_range(lo..=hi) } else { lo };
1066        let earliest = self.now + Duration::from_nanos(base);
1067
1068        let key = (from, to);
1069        let at = match self.last_delivery.get(&key) {
1070            Some(prev) if *prev >= earliest => *prev + Duration::from_nanos(1),
1071            _ => earliest,
1072        };
1073        self.last_delivery.insert(key, at);
1074        at
1075    }
1076}
1077
1078impl<P> Sim<P>
1079where
1080    P: Protocol,
1081    P::Cmd: Clone,
1082    P::Msg: Clone + PartialEq,
1083    P::Ind: Clone,
1084    P::Meta: Clone,
1085    P::Entry: Clone,
1086    P::Scope: From<SessionEvent>,
1087{
1088    /// Deliver session events to the protocol as scope events.
1089    ///
1090    /// Opt-in, like the codec check: a protocol that declares no scopes cannot receive one, and
1091    /// the bound lives only on this method so ordinary runs need nothing.
1092    ///
1093    /// **Forgetting this is silent, and it disables everything the session layers do.** Sessions
1094    /// still open, end and lose their suffixes; no layer is ever told, so every resend clause is
1095    /// dead, every `[session]` tag is unearned, and nothing fails to say so. A session-based run
1096    /// of a protocol whose `Scope` is inhabited should call it, and
1097    /// `forgetting_deliver_session_events_silently_disables_the_whole_bridge` is what that costs.
1098    pub fn deliver_session_events(&mut self) {
1099        self.session_scope = Some(P::Scope::from);
1100    }
1101}
1102
1103/// The store a protocol writes through: records what happened, and can kill the process.
1104///
1105/// When `doomed`, the first write applies or does not by a coin drawn before the handler ran, and
1106/// the process is then killed.
1107struct FaultyStore<'a, Me, En> {
1108    inner: &'a mut MemStore<Me, En>,
1109    writes: &'a mut Vec<WriteKind>,
1110    doomed: bool,
1111    keep: bool,
1112    died: &'a mut bool,
1113}
1114
1115impl<Me, En> FaultyStore<'_, Me, En> {
1116    /// Whether the write takes effect. Recorded either way: it was attempted.
1117    fn allow(&mut self, kind: WriteKind) -> bool {
1118        self.writes.push(kind);
1119        if !self.doomed {
1120            return true;
1121        }
1122        if *self.died {
1123            // The process is already gone; the handler is still running only because a synchronous
1124            // call cannot be interrupted. Nothing more it does reaches the disk.
1125            return false;
1126        }
1127        *self.died = true;
1128        self.keep
1129    }
1130}
1131
1132impl<Me, En> Store<Me, En> for FaultyStore<'_, Me, En> {
1133    fn get(&self) -> Option<&Me> {
1134        self.inner.get()
1135    }
1136
1137    fn set(&mut self, meta: Me) {
1138        if self.allow(WriteKind::Set) {
1139            self.inner.set(meta);
1140        }
1141    }
1142
1143    fn append(&mut self, entry: En) -> Position {
1144        if self.allow(WriteKind::Append) {
1145            return self.inner.append(entry);
1146        }
1147        self.inner.end()
1148    }
1149
1150    fn read_from(&self, from: Position) -> Vec<&En> {
1151        self.inner.read_from(from)
1152    }
1153
1154    fn end(&self) -> Position {
1155        self.inner.end()
1156    }
1157}