Skip to main content

recon_protocols/
perfect_failure_detector.rs

1//! Perfect failure detection.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 2.6 and Algorithm 2.5 ("Exclude on Timeout").
4//!
5//! **Status: deployable where synchrony is real. Space: bounded by membership.** All state is one
6//! entry per peer. A system without a known delivery bound wants an eventually perfect detector
7//! instead — see the timing note below.
8//!
9//! ```text
10//! PFD1: Strong completeness: Eventually, every process that crashes is permanently
11//!       detected by every correct process.
12//! PFD2: Strong accuracy:     If a process p is detected by any process, then p has crashed.
13//! ```
14//!
15//! **This protocol's guarantees are conditional on a timing assumption, and it is the first in
16//! this repository that has one.** Every abstraction below is correct in an asynchronous model: it
17//! assumes nothing about how long a message takes. Perfect detection is impossible there — a live
18//! process whose messages are merely slow is indistinguishable from a dead one, so any detector
19//! either accuses the living or never accuses the dead.
20//!
21//! What makes it possible is a *synchronous* system: message delivery between correct processes
22//! within a known bound Δ. PFD2 holds only while that bound does. Run this detector on a lossy or
23//! unbounded network and it will accuse correct processes — that is the assumption failing, not
24//! the implementation.
25//!
26//! ```text
27//! upon event ⟨ Timeout ⟩ do
28//!     forall p ∈ Π do
29//!         if p ∉ alive ∧ p ∉ detected then
30//!             detected := detected ∪ {p};
31//!             trigger ⟨ P, Crash | p ⟩;
32//!     forall p ∈ Π do
33//!         trigger ⟨ pl, Send | p, [HEARTBEATREQUEST] ⟩;
34//!     alive := ∅;
35//!     starttimer(Δ);
36//! ```
37//!
38//! Four departures from the page:
39//!
40//! - The book exchanges a request and a reply each round. This sends one unsolicited heartbeat
41//!   per round instead: with the same bound it distinguishes the same failures in half the
42//!   messages, and the round-trip's only role was to let the requester choose when to ask.
43//! - **The heartbeat period and the detection timeout are separate.** The book uses one delay Δ
44//!   for both, which makes a single missed round fatal: a process stalled for an instant spanning
45//!   its own send is accused, though it is alive and the network kept its promise. Beating every
46//!   `period` and accusing only after `timeout` of silence tolerates a stall shorter than
47//!   `timeout − period − Δ`. Accuracy still requires `timeout > period + Δ`.
48//! - `⟨P, Init⟩` *is* the init event, and this protocol has no commands at all. The first timer is
49//!   armed on the first tick request from the layer above.
50//! - **Heartbeats go on the wire directly, not through `pl`.** Algorithm 2.5 sends over perfect
51//!   links; this protocol is its own bottom layer and has no child. Under the synchronous
52//!   configuration the two are equivalent — the sim's synchronous mode zeroes the loss knob and
53//!   enforces the bound, so nothing is lost and a perfect link would add only deduplication that
54//!   an idempotent `last_heard` write does not need. Outside it they differ, and in the direction
55//!   that matters: a lost heartbeat is indistinguishable from a silent process, so loss forges
56//!   the very evidence this protocol accuses on. `accuracy_is_lost_when_the_timing_assumption_is`
57//!   is that difference, made to happen.
58//!
59//! # A stall has two sides, and the second one is worse
60//!
61//! The departure above covers the process a stall happens *to*: it misses a send and its peers
62//! accuse it, which the `timeout − period − Δ` margin is what tolerates. The other side is the
63//! stalled process itself, and it is not tolerated at all.
64//!
65//! A process descheduled for longer than `timeout` comes back holding a tick that is due and a
66//! `last_heard` for every peer that is now older than the timeout. The measurement is not wrong —
67//! it genuinely heard nothing for that long — so it accuses every one of them, and detection is
68//! permanent. `a_stalled_process_accuses_its_peers_when_it_comes_back` pins that, and the pair of
69//! suspension tests either side of it shows the margin still holds from the inside for a stall
70//! shorter than the timeout.
71//!
72//! This is the timing assumption failing rather than the implementation, and it fails in the one
73//! place the assumption is easiest to forget: Δ bounds the *network*, and a synchronous system
74//! bounds process scheduling too. A detector that discounts its own stall — noticing that far
75//! more than `period` elapsed between consecutive ticks, and treating that round as unmeasured —
76//! is a real technique and a real departure from the page, so it is a change with a proposal
77//! rather than something to add here. `Sim::resume` documents the same asymmetry from the
78//! simulator's side.
79
80use core::time::Duration;
81use recon_core::{NodeId, ProtoCx, Protocol, Time, TimerId};
82
83use crate::detector::{Detector, DetectorInd};
84use serde::{Deserialize, Serialize};
85use std::collections::{BTreeMap, BTreeSet};
86
87/// What a process says to show it is alive.
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89pub struct Heartbeat;
90
91/// Requests from the layer above: none.
92///
93/// Detection begins at initialisation, as Module 2.6 has it, so there is nothing to ask for. An
94/// uninhabited type says so in a way the compiler checks.
95pub type Cmd = core::convert::Infallible;
96
97/// Indications to the layer above.
98#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum Ind {
100    /// `node` has crashed. Raised exactly once per process, and never retracted.
101    Crash { node: NodeId },
102}
103
104/// Detects crashes by heartbeat timeout.
105#[derive(Debug)]
106pub struct PerfectFailureDetector {
107    me: NodeId,
108    peers: BTreeSet<NodeId>,
109    /// When each peer was last heard from.
110    last_heard: BTreeMap<NodeId, Time>,
111    /// Already reported. Detection is permanent, so these are never revisited.
112    detected: BTreeSet<NodeId>,
113    /// How often this process announces itself.
114    period: Duration,
115    /// How long a peer may be silent before it is declared crashed.
116    timeout: Duration,
117    /// The tick outstanding, if any. A handle rather than a flag: an expiry this detector has
118    /// superseded is then recognisable, where `true` said only that something was pending.
119    tick: Option<TimerId>,
120}
121
122impl PerfectFailureDetector {
123    /// Detect among `peers`, announcing this process every `period` and declaring a peer crashed
124    /// after `timeout` of silence.
125    ///
126    /// For strong accuracy, `timeout` must exceed `period` plus the network's delivery bound —
127    /// otherwise a live process's heartbeat can arrive after it has already been accused.
128    /// Configure the bound from the simulator's own `Sim::delivery_bound` rather than guessing
129    /// it. Named rather than linked: `recon-sim` is a dev-dependency of this crate, so the path
130    /// does not exist for a reader of these docs — and it is the right way round, since a protocol
131    /// that depended on the simulator would be the thing constraint 2 forbids.
132    pub fn new(
133        me: NodeId,
134        peers: impl IntoIterator<Item = NodeId>,
135        period: Duration,
136        timeout: Duration,
137    ) -> Self {
138        debug_assert!(
139            timeout > period,
140            "a timeout no longer than the heartbeat period accuses live processes"
141        );
142        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
143        peers.remove(&me);
144        PerfectFailureDetector {
145            me,
146            peers,
147            last_heard: BTreeMap::new(),
148            detected: BTreeSet::new(),
149            period,
150            timeout,
151            tick: None,
152        }
153    }
154
155    /// How often this process announces itself.
156    pub fn period(&self) -> Duration {
157        self.period
158    }
159
160    /// The processes currently believed crashed.
161    pub fn detected(&self) -> impl Iterator<Item = NodeId> + '_ {
162        self.detected.iter().copied()
163    }
164
165    /// Whether `node` has been detected as crashed.
166    pub fn has_detected(&self, node: NodeId) -> bool {
167        self.detected.contains(&node)
168    }
169
170    /// The processes not yet detected as crashed, this one included.
171    pub fn correct(&self) -> impl Iterator<Item = NodeId> + '_ {
172        core::iter::once(self.me)
173            .chain(self.peers.iter().copied().filter(|p| !self.detected.contains(p)))
174    }
175
176    /// The silence a process is allowed before being declared crashed.
177    pub fn timeout(&self) -> Duration {
178        self.timeout
179    }
180
181    fn beat(&mut self, cx: &mut ProtoCx<'_, Self>) {
182        for &p in &self.peers {
183            cx.send(p, Heartbeat);
184        }
185    }
186
187    /// Treat every peer as heard from now, so the first round does not accuse anyone before a
188    /// heartbeat has had time to arrive.
189    fn assume_alive(&mut self, now: Time) {
190        for &p in &self.peers {
191            self.last_heard.insert(p, now);
192        }
193    }
194}
195
196/// `P` satisfies the detector port, and never withdraws a suspicion.
197///
198/// `PFD2` is *strong* accuracy — a detected process has crashed — so there is nothing to take back.
199/// The port's second arm is unreachable over this detector rather than merely unused, which
200/// `tests/detector_port.rs` pins.
201impl Detector for PerfectFailureDetector {
202    fn classify(Ind::Crash { node }: Ind) -> DetectorInd {
203        DetectorInd::Suspect { node }
204    }
205}
206
207impl Protocol for PerfectFailureDetector {
208    type Cmd = Cmd;
209    type Ind = Ind;
210    type Msg = Heartbeat;
211    /// No scope conditions: this protocol's guarantees do not lapse.
212    type Scope = core::convert::Infallible;
213    type Note = crate::Note;
214    /// Keeps nothing durably: a crash loses everything this protocol knows.
215    type Meta = core::convert::Infallible;
216    type Entry = core::convert::Infallible;
217
218    /// `⟨ P, Init ⟩ do alive := Π; detected := ∅; starttimer(Δ)`.
219    ///
220    /// The book's own trigger, now that there is one. It used to be a `Start` command because
221    /// there was no init event to hang the first timer on; there is, so there is no command.
222    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
223        if self.tick.is_some() {
224            return;
225        }
226        self.assume_alive(cx.now());
227        self.beat(cx);
228        self.tick = Some(cx.set_timer(self.period));
229    }
230
231    fn on_cmd(&mut self, cmd: Cmd, _cx: &mut ProtoCx<'_, Self>) {
232        match cmd {}
233    }
234
235    fn on_msg(&mut self, from: NodeId, Heartbeat: Heartbeat, cx: &mut ProtoCx<'_, Self>) {
236        // A process already detected stays detected: PFD1 demands permanence, and under the
237        // timing assumption a heartbeat from an accused process cannot arrive.
238        if self.peers.contains(&from) {
239            self.last_heard.insert(from, cx.now());
240        }
241    }
242
243    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
244        // Handed down from every layer above along with everybody else's, so only the one this
245        // detector is waiting on does anything. An expiry it has superseded accuses nobody.
246        if self.tick != Some(id) {
247            return;
248        }
249        let now = cx.now();
250        let silent: Vec<NodeId> = self
251            .peers
252            .iter()
253            .copied()
254            .filter(|p| !self.detected.contains(p))
255            .filter(|p| {
256                let last = self.last_heard.get(p).copied().unwrap_or(Time::ZERO);
257                now.saturating_since(last) > self.timeout
258            })
259            .collect();
260        for p in silent {
261            self.detected.insert(p);
262            cx.indicate(Ind::Crash { node: p });
263        }
264
265        self.beat(cx);
266        self.tick = Some(cx.set_timer(self.period));
267    }
268}