Skip to main content

recon_protocols/
epoch_consensus.rs

1//! Read/write epoch consensus — the quorum core, and where Paxos's safety argument lives.
2//!
3//! **Status: implementation. Space: bounded by membership.**
4//!
5//! Cachin, Guerraoui & Rodrigues, Module 5.4 and Algorithm 5.6 ("Read/Write Epoch Consensus"),
6//! quoted from the book:
7//!
8//! ```text
9//! Algorithm 5.6: Read/Write Epoch Consensus
10//! Implements: EpochConsensus, instance ep, with timestamp ets and leader ℓ.
11//! Uses:
12//!     PerfectPointToPointLinks, instance pl;
13//!     BestEffortBroadcast, instance beb.
14//!
15//! upon event ⟨ ep, Init | state ⟩ do
16//!     (valts, val) := state; tmpval := ⊥; states := [⊥]^N; accepted := 0;
17//!
18//! upon event ⟨ ep, Propose | v ⟩ do                       // only leader ℓ
19//!     tmpval := v;
20//!     trigger ⟨ beb, Broadcast | [READ] ⟩;
21//!
22//! upon event ⟨ beb, Deliver | ℓ, [READ] ⟩ do
23//!     trigger ⟨ pl, Send | ℓ, [STATE, valts, val] ⟩;
24//!
25//! upon event ⟨ pl, Deliver | q, [STATE, ts, v] ⟩ do       // only leader ℓ
26//!     states[q] := (ts, v);
27//!
28//! upon #(states) > N/2 do                                 // only leader ℓ
29//!     (ts, v) := highest(states);
30//!     if v ≠ ⊥ then tmpval := v;
31//!     states := [⊥]^N;
32//!     trigger ⟨ beb, Broadcast | [WRITE, tmpval] ⟩;
33//!
34//! upon event ⟨ beb, Deliver | ℓ, [WRITE, v] ⟩ do
35//!     (valts, val) := (ets, v);
36//!     trigger ⟨ pl, Send | ℓ, [ACCEPT] ⟩;
37//!
38//! upon event ⟨ pl, Deliver | q, [ACCEPT] ⟩ do             // only leader ℓ
39//!     accepted := accepted + 1;
40//!
41//! upon accepted > N/2 do                                  // only leader ℓ
42//!     accepted := 0;
43//!     trigger ⟨ beb, Broadcast | [DECIDED, tmpval] ⟩;
44//!
45//! upon event ⟨ beb, Deliver | ℓ, [DECIDED, v] ⟩ do
46//!     trigger ⟨ ep, Decide | v ⟩;
47//!
48//! upon event ⟨ ep, Abort ⟩ do
49//!     trigger ⟨ ep, Aborted | (valts, val) ⟩;
50//!     halt;                                               // stop operating when aborted
51//! ```
52//!
53//! # Why two majorities are the whole algorithm
54//!
55//! The leader reads from a majority and writes to a majority, and any two majorities of `Π`
56//! intersect. So if some epoch decided `v` — meaning a majority accepted it — then every later
57//! epoch's read reaches at least one process that accepted `v`, and `highest(states)` returns it.
58//! `if v ≠ ⊥ then tmpval := v` is the line that makes the later leader adopt it instead of its own
59//! proposal. That is the entire reason two epochs cannot decide differently, and every other part
60//! of Paxos exists to arrange for it.
61//!
62//! **`highest` means highest *timestamp*, not highest value.** It picks the state written in the
63//! most recent epoch, which is the one that may already have been decided.
64//!
65//! # `halt` is a safety property, not tidiness
66//!
67//! `upon event ⟨ ep, Abort ⟩ … halt;  // stop operating when aborted`. An instance that kept
68//! answering after being abandoned would be a second leader for its epoch under another name: it
69//! could still collect a quorum and decide, while the epoch that replaced it decided something
70//! else. The flag [`EpochConsensus::is_aborted`] reports is checked at the top of every handler, and
71//! the suite delivers a message to an aborted instance and asserts that nothing at all comes out.
72//!
73//! # Departure: directed replies travel by directed broadcast
74//!
75//! As in [`crate::epoch_change`], the book's `pl` child is absorbed into the broadcast's directed
76//! [`beb::Cmd::SendTo`] — one addressed process, one wire message, which is what `pl, Send` does.
77//! `STATE` and `ACCEPT` both travel that way. One child and one wire variant fewer, and the
78//! guarantee is unchanged.
79//!
80//! ```text
81//! EPC1 [always]  Validity — a decided value was proposed in this epoch, or was the highest-
82//!                timestamped value some process had already accepted
83//! EPC2 [always]  Uniform agreement — no two processes decide differently in one epoch
84//! EPC3 [always]  Integrity — a process decides at most once
85//! EPC4 [always]  Lock-in — a value decided in an earlier epoch is what a later one decides
86//! EPC5 [always]  Abort behaviour — an abandoned instance reports its state and then is silent
87//! ```
88
89use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
90use serde::{Deserialize, Serialize};
91use std::collections::{BTreeMap, BTreeSet};
92
93use crate::best_effort_broadcast::{self as beb, BestEffortBroadcast};
94use crate::perfect_link as pl;
95
96/// `(valts, val)` — what a process has accepted, and when.
97///
98/// `val` is `None` for the book's `⊥`: nothing accepted yet, at timestamp zero.
99#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
100pub struct State<V> {
101    pub valts: u64,
102    pub val: Option<V>,
103}
104
105impl<V> Default for State<V> {
106    /// The state a first epoch begins from: nothing accepted.
107    fn default() -> Self {
108        State { valts: 0, val: None }
109    }
110}
111
112/// What this layer puts on the wire, beneath the broadcast.
113#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
114pub enum EpochMsg<V> {
115    /// `[READ]` — the leader asking what everyone holds.
116    Read,
117    /// `[STATE, valts, val]` — a follower's answer, addressed to the leader.
118    StateIs { valts: u64, val: Option<V> },
119    /// `[WRITE, v]` — the leader telling everyone what to accept.
120    Write { val: V },
121    /// `[ACCEPT]` — a follower's acknowledgement, addressed to the leader.
122    Accept,
123    /// `[DECIDED, v]` — the leader announcing the decision.
124    Decided { val: V },
125}
126
127/// Requests from the layer above.
128#[derive(Debug, Clone, PartialEq, Eq)]
129pub enum Cmd<V> {
130    /// `⟨ ep, Propose | v ⟩`. Acted on only by this epoch's leader.
131    Propose(V),
132    /// `⟨ ep, Abort ⟩`. Answered by [`Ind::Aborted`], after which the instance is silent.
133    Abort,
134}
135
136/// Indications to the layer above.
137#[derive(Debug, Clone, PartialEq, Eq)]
138pub enum Ind<V> {
139    /// `⟨ ep, Decide | v ⟩`.
140    Decide(V),
141    /// `⟨ ep, Aborted | (valts, val) ⟩` — the state this instance held when it was abandoned.
142    Aborted(State<V>),
143}
144
145/// An epoch message, stamped with the instance it belongs to.
146///
147/// The book writes `ep.ts` and guards every handler with `such that ts = ets`, so instances are
148/// addressed by timestamp and a message for one never reaches another. Nothing in this codebase's
149/// wire does that for free, and the consequence of omitting it is a **safety** failure rather than a
150/// lost message: a `WRITE` from epoch 7 arriving after epoch 11 began would be accepted and recorded
151/// at timestamp 11, inventing an acceptance that never happened.
152///
153/// The stamp lives here rather than in the layer above because the epoch is this instance's own
154/// identity — it stamps what it sends and drops what is not addressed to it, so a parent cannot
155/// forget to.
156#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
157pub struct Tagged<V> {
158    pub ets: u64,
159    pub msg: EpochMsg<V>,
160}
161
162/// What the broadcast beneath puts on the wire for this layer's messages.
163pub type BebMsg<V> = pl::Wire<Tagged<V>>;
164
165/// Abortable consensus within one epoch.
166#[derive(Debug)]
167pub struct EpochConsensus<V: Clone> {
168    me: NodeId,
169    peers: BTreeSet<NodeId>,
170    /// `ets` — this instance's epoch timestamp.
171    ets: u64,
172    /// `ℓ` — this epoch's leader.
173    leader: NodeId,
174    /// `(valts, val)`.
175    state: State<V>,
176    /// `tmpval` — the value the leader is trying to write.
177    tmpval: Option<V>,
178    /// `states` — what the leader has read back, by process.
179    states: BTreeMap<NodeId, State<V>>,
180    /// `accepted` — how many have acknowledged the write.
181    accepted: usize,
182    /// Whether the write has already been sent, so a second majority of `STATE` cannot resend it.
183    written: bool,
184    /// Whether the decision has been announced, so a second majority of `ACCEPT` cannot re-announce.
185    announced: bool,
186    /// `halt`. Every handler returns immediately once this is set.
187    aborted: bool,
188    beb: Child<BestEffortBroadcast<Tagged<V>>>,
189}
190
191impl<V: Clone> EpochConsensus<V> {
192    /// `⟨ ep, Init | state ⟩` — an instance for epoch `ets` led by `leader`, beginning from `state`.
193    pub fn new(
194        me: NodeId,
195        peers: impl IntoIterator<Item = NodeId>,
196        ets: u64,
197        leader: NodeId,
198        state: State<V>,
199        retransmit: core::time::Duration,
200    ) -> Self {
201        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
202        peers.insert(me);
203        EpochConsensus {
204            me,
205            peers: peers.clone(),
206            ets,
207            leader,
208            state,
209            tmpval: None,
210            states: BTreeMap::new(),
211            accepted: 0,
212            written: false,
213            announced: false,
214            aborted: false,
215            beb: Child::new(BestEffortBroadcast::new(me, peers, retransmit)),
216        }
217    }
218
219    /// This epoch's timestamp.
220    pub fn timestamp(&self) -> u64 {
221        self.ets
222    }
223
224    /// Whether this instance has been abandoned and is therefore silent.
225    pub fn is_aborted(&self) -> bool {
226        self.aborted
227    }
228
229    /// What this process has accepted, and when.
230    pub fn state(&self) -> &State<V> {
231        &self.state
232    }
233
234    /// `N/2` — the threshold both majorities are measured against.
235    fn majority(&self) -> usize {
236        self.peers.len() / 2
237    }
238
239    fn is_leader(&self) -> bool {
240        self.me == self.leader
241    }
242
243    fn broadcast(&mut self, msg: EpochMsg<V>, cx: &mut ProtoCx<'_, Self>) {
244        let tagged = Tagged { ets: self.ets, msg };
245        self.through_beb(cx, |b, ccx| b.on_cmd(beb::Cmd::Broadcast(tagged), ccx));
246    }
247
248    fn send_to(&mut self, to: NodeId, msg: EpochMsg<V>, cx: &mut ProtoCx<'_, Self>) {
249        let tagged = Tagged { ets: self.ets, msg };
250        self.through_beb(cx, |b, ccx| b.on_cmd(beb::Cmd::SendTo { to, msg: tagged }, ccx));
251    }
252
253    /// `highest(states)` — the state with the greatest timestamp among those read.
254    fn highest(&self) -> Option<State<V>> {
255        self.states.values().max_by_key(|s| s.valts).cloned()
256    }
257
258    fn on_epoch_msg(&mut self, from: NodeId, msg: EpochMsg<V>, cx: &mut ProtoCx<'_, Self>) {
259        // `halt` — an abandoned instance answers nothing, which is what stops it from becoming a
260        // second leader for an epoch that has moved on.
261        if self.aborted {
262            return;
263        }
264        match msg {
265            // `upon event ⟨ beb, Deliver | ℓ, [READ] ⟩`
266            EpochMsg::Read if from == self.leader => {
267                let reply =
268                    EpochMsg::StateIs { valts: self.state.valts, val: self.state.val.clone() };
269                self.send_to(from, reply, cx);
270            }
271            // `upon event ⟨ pl, Deliver | q, [STATE, ts, v] ⟩`  // only leader
272            EpochMsg::StateIs { valts, val } if self.is_leader() => {
273                self.states.insert(from, State { valts, val });
274                self.maybe_write(cx);
275            }
276            // `upon event ⟨ beb, Deliver | ℓ, [WRITE, v] ⟩`
277            EpochMsg::Write { val } if from == self.leader => {
278                self.state = State { valts: self.ets, val: Some(val) };
279                self.send_to(from, EpochMsg::Accept, cx);
280            }
281            // `upon event ⟨ pl, Deliver | q, [ACCEPT] ⟩`  // only leader
282            EpochMsg::Accept if self.is_leader() => {
283                self.accepted += 1;
284                self.maybe_decide(cx);
285            }
286            // `upon event ⟨ beb, Deliver | ℓ, [DECIDED, v] ⟩`
287            EpochMsg::Decided { val } if from == self.leader => {
288                cx.indicate(Ind::Decide(val));
289            }
290            // A message from someone who is not this epoch's leader, or a leader-only message at a
291            // follower. Neither is addressed to this process's role; the book's guards drop them.
292            _ => {}
293        }
294    }
295
296    /// `upon #(states) > N/2 do … trigger ⟨ beb, Broadcast | [WRITE, tmpval] ⟩`.
297    fn maybe_write(&mut self, cx: &mut ProtoCx<'_, Self>) {
298        if self.written || self.states.len() <= self.majority() {
299            return;
300        }
301        // `(ts, v) := highest(states); if v ≠ ⊥ then tmpval := v;`
302        //
303        // This is the line the whole algorithm turns on. A value already accepted in a higher
304        // epoch displaces this leader's own proposal, which is what stops two epochs deciding
305        // differently. The book prints `≠`; an OCR of the page renders it as `=`, which would
306        // invert the meaning and break safety outright.
307        if let Some(highest) = self.highest()
308            && highest.val.is_some()
309        {
310            self.tmpval = highest.val;
311        }
312        self.states.clear();
313        self.written = true;
314        if let Some(val) = self.tmpval.clone() {
315            self.broadcast(EpochMsg::Write { val }, cx);
316        }
317    }
318
319    /// `upon accepted > N/2 do … trigger ⟨ beb, Broadcast | [DECIDED, tmpval] ⟩`.
320    fn maybe_decide(&mut self, cx: &mut ProtoCx<'_, Self>) {
321        if self.announced || self.accepted <= self.majority() {
322            return;
323        }
324        self.accepted = 0;
325        self.announced = true;
326        if let Some(val) = self.tmpval.clone() {
327            self.broadcast(EpochMsg::Decided { val }, cx);
328        }
329    }
330
331    fn through_beb(
332        &mut self,
333        cx: &mut ProtoCx<'_, Self>,
334        f: impl FnOnce(
335            &mut BestEffortBroadcast<Tagged<V>>,
336            &mut ProtoCx<'_, BestEffortBroadcast<Tagged<V>>>,
337        ),
338    ) {
339        let mut inds = self.beb.run(cx, core::convert::identity, f);
340        for ind in inds.drain(..) {
341            match ind {
342                // `such that ts = ets` — traffic for another instance is not this one's business,
343                // and reading it would invent an acceptance at the wrong timestamp. Unreachable
344                // while `on_msg` guards the door, and kept because this is where the book puts it.
345                beb::Ind::Deliver { from, msg } if msg.ets == self.ets => {
346                    self.on_epoch_msg(from, msg.msg, cx)
347                }
348                beb::Ind::Deliver { .. } => {}
349                beb::Ind::SessionEnded { .. } | beb::Ind::SessionEstablished { .. } => {}
350            }
351        }
352        self.beb.reclaim(inds);
353    }
354}
355
356impl<V: Clone> Protocol for EpochConsensus<V> {
357    type Cmd = Cmd<V>;
358    type Ind = Ind<V>;
359    type Msg = BebMsg<V>;
360    type Scope = core::convert::Infallible;
361    type Note = crate::Note;
362    /// Keeps nothing durably. `logged_epoch_consensus` is the variant that does.
363    type Meta = core::convert::Infallible;
364    type Entry = core::convert::Infallible;
365
366    fn on_cmd(&mut self, cmd: Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
367        if self.aborted {
368            return;
369        }
370        match cmd {
371            // `upon event ⟨ ep, Propose | v ⟩ do tmpval := v; … // only leader ℓ`
372            Cmd::Propose(v) => {
373                if self.is_leader() {
374                    self.tmpval = Some(v);
375                    self.broadcast(EpochMsg::Read, cx);
376                }
377            }
378            // `upon event ⟨ ep, Abort ⟩ do trigger ⟨ ep, Aborted | (valts, val) ⟩; halt;`
379            Cmd::Abort => {
380                self.aborted = true;
381                cx.indicate(Ind::Aborted(self.state.clone()));
382            }
383        }
384    }
385
386    /// `such that ts = ets`, applied **at the door** rather than after the link beneath.
387    ///
388    /// The guard has to be here, not only where the delivery is handled, and the reason is the
389    /// perfect link's duplicate-detection set. Each epoch gets a new instance, so each epoch gets a
390    /// new link, and a new link restarts its sequence numbers at one — while the *receiver's* set is
391    /// cleared at a different moment, when its own epoch changes. Hand a foreign-epoch message to
392    /// the link and it records `(src, 1)` as delivered; the real epoch-`ets` message with sequence
393    /// one is then discarded as a duplicate, silently, and that process never answers the leader
394    /// again. Three of five processes stalled this way before the guard moved up here.
395    ///
396    /// This is `CLAUDE.md`'s "identity is as durable as the state it keys" seen from the other side:
397    /// the identifier's scope is one epoch, so nothing outside that epoch may enter the set that
398    /// keys on it.
399    fn on_msg(&mut self, from: NodeId, msg: BebMsg<V>, cx: &mut ProtoCx<'_, Self>) {
400        if self.aborted || msg.payload.ets != self.ets {
401            return;
402        }
403        self.through_beb(cx, |b, ccx| b.on_msg(from, msg, ccx));
404    }
405
406    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
407        if self.aborted {
408            return;
409        }
410        self.through_beb(cx, |b, ccx| b.on_timer(id, ccx));
411    }
412}