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}