Skip to main content

recon_protocols/
multi_paxos_replica.rs

1//! The Multi-Paxos replica: slots become positions in a log.
2//!
3//! **Status: implementation. Space: the bookkeeping is bounded — `decisions` and `performed` by a
4//! retention window, `requests` and `proposals` by `WINDOW` — and the ordered sequence is
5//! deliberately not, because it is the data rather than the bookkeeping.**
6//!
7//! §4.2 is applied, and it is worth being clear which half. The section separates state that is
8//! *unnecessary* — the leader's and acceptor's, collected below a watermark once `f + 1` replicas
9//! have applied past it, which is [`crate::multi_paxos_synod`]'s — from state that is
10//! **unavoidable**, which is this module's:
11//!
12//! > because commands may be decided in multiple slots, replicas each maintain a set of all
13//! > decisions to filter out such duplicates. In practice, it is often sufficient if such
14//! > information is only kept for a certain amount of time, making the probability of duplicate
15//! > execution negligible.
16//!
17//! So the filter is bounded by a **retention window** rather than by the watermark beneath it, and
18//! the reason is worth stating because using the watermark is the obvious wrong move: that
19//! watermark says `f + 1` replicas have *applied* up to a slot, which says nothing whatever about
20//! whether a command decided below it may be decided again above it. A duplicate filter has to
21//! outlive every slot a duplicate could span, and no watermark bounds that.
22//!
23//! **What that weakens, stated rather than discovered.** "A command decided in two slots takes one
24//! position" holds *within the retention window*. A duplicate arriving more than [`RETAIN`] slots
25//! after the first would take a second position. As things stand that is a weakening of something
26//! unreachable — a command is minted at one replica and occupies one slot at a time, so a duplicate
27//! needs the client-retry deployment this port does not have — but it is the honest bound and it is
28//! in the specification.
29//!
30//! **The ordered sequence is exempt and is not collected.** A log grows with what is appended to
31//! it; that is the data, and a log that discarded its entries would not be one. What would bound it
32//! is a snapshot, which is outside the paper and is also the point past which a lagging replica
33//! can no longer be caught up — see the catch-up below.
34//!
35//! # Source
36//!
37//! van Renesse, R. and Altinbuken, D. (2015) 'Paxos Made Moderately Complex', *ACM Computing
38//! Surveys*, 47(3), pp. 1–36 — §2.1 and **Figure 1**, "Pseudocode for a replica", page 42:6.
39//!
40//! **Not** van Renesse, R. (2011), the Cornell technical report of the same title. The two differ
41//! here more than anywhere else in the algorithm: the report's replica keeps a single `slot num`
42//! and its `propose(p)` searches for the lowest slot not already proposed or decided, with no
43//! window at all. The survey splits that into `slot_in` and `slot_out` and adds `WINDOW`, which is
44//! the version transcribed below. A reader checking this code against a page needs to know which
45//! page.
46//!
47//! ```text
48//! process Replica(leaders, initial_state)
49//!   var state := initial_state, slot_in := 1, slot_out := 1;
50//!   var requests := ∅, proposals := ∅, decisions := ∅;
51//!
52//!   function propose()
53//!     while slot_in < slot_out + WINDOW ∧ ∃c : c ∈ requests do
54//!       if ∃op : ⟨slot_in − WINDOW, ⟨·, ·, op⟩⟩ ∈ decisions ∧ isreconfig(op) then
55//!         leaders := op.leaders;
56//!       end if
57//!       if ∄c' : ⟨slot_in, c'⟩ ∈ decisions then
58//!         requests := requests \ {c};
59//!         proposals := proposals ∪ {⟨slot_in, c⟩};
60//!         ∀λ ∈ leaders : send(λ, ⟨propose, slot_in, c⟩);
61//!       end if
62//!       slot_in := slot_in + 1;
63//!     end while
64//!   end function
65//!
66//!   function perform(⟨κ, cid, op⟩)
67//!     if (∃s : s < slot_out ∧ ⟨s, ⟨κ, cid, op⟩⟩ ∈ decisions) ∨ isreconfig(op) then
68//!       slot_out := slot_out + 1;
69//!     else
70//!       ⟨next, result⟩ := op(state);
71//!       atomic
72//!         state := next; slot_out := slot_out + 1;
73//!       end atomic
74//!       send(κ, ⟨response, cid, result⟩);
75//!     end if
76//!   end function
77//!
78//!   for ever
79//!     switch receive()
80//!       case ⟨request, c⟩ :
81//!         requests := requests ∪ {c};
82//!       end case
83//!       case ⟨decision, s, c⟩ :
84//!         decisions := decisions ∪ {⟨s, c⟩};
85//!         while ∃c' : ⟨slot_out, c'⟩ ∈ decisions do
86//!           if ∃c'' : ⟨slot_out, c''⟩ ∈ proposals then
87//!             proposals := proposals \ {⟨slot_out, c''⟩};
88//!             if c'' ≠ c' then
89//!               requests := requests ∪ {c''};
90//!             end if
91//!           end if
92//!           perform(c');
93//!         end while
94//!       end case
95//!     end switch
96//!     propose();
97//!   end for
98//! end process
99//! ```
100//!
101//! # The invariants, quoted
102//!
103//! - **R1**: "There are no two different commands decided for the same slot:
104//!   `∀s, ρ1, ρ2, c1, c2 : ⟨s, c1⟩ ∈ ρ1.decisions ∧ ⟨s, c2⟩ ∈ ρ2.decisions ⇒ c1 = c2`." Held by the
105//!   child, not here: it is the Synod protocol's S1, and this layer relies on it rather than
106//!   enforcing it. What this layer does is not break it — `decisions` never
107//!   replaces an entry.
108//! - **R2**: "All commands up to `slot out` are in the set of decisions:
109//!   `∀ρ, s : 1 ≤ s < ρ.slot out ⇒ ∃c : ⟨s, c⟩ ∈ ρ.decisions`." Held because `slot_out` advances
110//!   only inside `perform`, which the drain loop calls only for a slot that has a decision.
111//! - **R3**: "For all replicas ρ, `ρ.state` is the result of applying the commands
112//!   `⟨s, cs⟩ ∈ ρ.decisions` to `initial state` for all `s` up to `slot out`, in order of slot
113//!   number." Here `state` is the ordered sequence — see the departures — so R3 is the statement
114//!   that the sequence is exactly the decided commands below `slot_out`, in slot order, minus the
115//!   ones `perform` skipped as already applied.
116//! - **R4**: "For each ρ, the variable `ρ.slot out` cannot decrease over time." Structural: the
117//!   only assignment is `+= 1`.
118//! - **R5**: "A replica proposes commands only for slots for which it knows the configuration:
119//!   `∀ρ : ρ.slot in < ρ.slot out + WINDOW`." The `while` guard in `propose`, kept although the
120//!   reason for it is not — see below.
121//!
122//!   **The formula and the sentence beside it do not say quite the same thing, and the code
123//!   follows the sentence.** `propose`'s loop tests `slot_in < slot_out + WINDOW` at the top and
124//!   increments `slot_in` at the bottom, so an exit with requests still queued leaves
125//!   `slot_in = slot_out + WINDOW` exactly — the strict inequality does not hold of the variable
126//!   between the loop ending and `slot_out` next advancing. What does hold, always, is the English:
127//!   every slot ever *proposed for* is below `slot_out + WINDOW` at the moment of proposing, since
128//!   the guard is checked before the slot is used. The suite asserts that form over
129//!   [`MultiPaxosReplica::proposed_slots`] and the loop's own `≤` over the variable, rather than
130//!   asserting a formula the page's own pseudocode breaks.
131//!
132//! # Where each part of Figure 1 went
133//!
134//! | Figure 1 | Here |
135//! |---|---|
136//! | `var state := initial_state` | the ordered sequence; there is no application state machine, because the port is a log |
137//! | `slot_in`, `slot_out` | fields, counting slots |
138//! | `requests`, `proposals`, `decisions` | fields; `decisions` is append-only, as the page has it |
139//! | `function propose()` | `transfer` plus the send loop in `pump` |
140//! | the `isreconfig` branch in `propose()` | absent — reconfiguration is a later change; the `WINDOW` guard around it is kept |
141//! | `function perform()` | `perform`, minus `op(state)` and the client response |
142//! | `case ⟨request, c⟩` | [`Cmd::Append`], the port's own |
143//! | `case ⟨decision, s, c⟩` | the child's [`crate::multi_paxos_synod::Ind::Decision`] |
144//! | `∀λ ∈ leaders : send(λ, ⟨propose, …⟩)` | a call into the child, per §4.4 |
145//! | `send(κ, ⟨response, cid, result⟩)` | absent, with `state` |
146//!
147//! # Departures from the page
148//!
149//! - **`state` is a sequence, and there is no `op(state)`.** The port this satisfies is
150//!   [`crate::total_order_log::TotalOrderLog`]: a log, not a replicated state machine. `perform`
151//!   therefore appends the command to the sequence where the page applies it, and the client
152//!   response goes with the result it would have carried. What survives is the part R1–R4 are about
153//!   — that every replica applies the same commands in the same order — and the part that goes is
154//!   the application on top of it.
155//!
156//! - **`∀λ ∈ leaders : send(λ, ⟨propose, s, c⟩)` is a call into the child.** §4.4: "each machine
157//!   that runs a replica also runs a leader… the replica can send a proposal for a particular slot
158//!   to its local leader". So `leaders` is not a set here, it is one child, and the fan-out to
159//!   remote leaders is the child's forwarding rather than this layer's broadcast. That is what lets
160//!   this module hold no link, no broadcast and no wire of its own; the cost is in
161//!   [`crate::multi_paxos_synod`]'s own documentation, and it is that proposal *delivery* now rests
162//!   on the leader detector.
163//!
164//! - **The already-decided check is indexed rather than scanned.** The page writes
165//!   `∃s : s < slot_out ∧ ⟨s, ⟨κ, cid, op⟩⟩ ∈ decisions`, a scan of every decision below
166//!   `slot_out`. `performed` is exactly that set — `perform` runs once per slot in slot order, so
167//!   what it has seen *is* `{decisions[s] : s < slot_out}` — and membership answers the same
168//!   question in log time. It grows the same unbounded way `decisions` does, so it changes nothing
169//!   about the space statement above.
170//!
171//! - **`requests` is a queue, where the page says "any command".** `∃c : c ∈ requests` picks an
172//!   arbitrary member. FIFO is a refinement rather than a departure in the strict sense, but it is
173//!   worth naming because it is what stops a command being passed over for ever while later ones
174//!   are proposed — the page leaves that to whoever implements the choice.
175//!
176//! - **`WINDOW` is kept and its reason is deferred.** In the source the window exists because a
177//!   reconfiguration decided in slot `s` takes effect at `s + WINDOW`, so a replica may not propose
178//!   past the last slot whose configuration it knows. Reconfiguration is not built here and the
179//!   membership is fixed for the run, so the guard is a pipeline cap and nothing more. It is kept
180//!   rather than dropped so that R5 reads against the page, and this note is here so that a reader
181//!   does not take the cap for the whole reason.
182//!
183//! - **§4.2's periodic report, and the catch-up that pays for collecting.** A replica tells the
184//!   consensus beneath it how far it has applied, periodically and not as a consequence of doing
185//!   work — a replica applying nothing is the one whose position others most need, and a report
186//!   riding its own traffic would fall silent exactly then.
187//!
188//!   Collecting at `f + 1` means `f` correct replicas may be behind, and what would have helped
189//!   them is what was collected. §4.2's own remedy is that "replicas can learn decisions […] from
190//!   one another", so a replica whose proposal is refused as collected asks a peer, and the peer
191//!   answers from *its* `decisions` — the one place that still holds them. A replica further behind
192//!   than [`RETAIN`] cannot be caught up, because nobody holds those decisions any more.
193//!
194//! - **The fourth liveness violation, and its fix.** Liu, Y.A., Chand, S. and Stoller, S.D. (2019)
195//!   'Moderately Complex Paxos Made Simple', PPDP '19, is the cross-check this module's source is
196//!   read against. Its fourth liveness violation is this layer's: if no decision arrives for a
197//!   slot, `slot_out` stops moving, `WINDOW` fills, `slot_in` stops advancing and the replica
198//!   wedges — with nothing on the page to get it out, because Figure 1 acts only on messages that
199//!   arrive. The fix has two halves, and the other one is the leader's: a proposal outstanding
200//!   longer than `repropose_after` is proposed **again for the same slot**,
201//!   and a leader that has seen the slot decided answers with the decision rather than dropping
202//!   the repeat. Re-proposing into a *new* slot instead would leave the old one unfilled for ever,
203//!   which is the wedge itself.
204//!
205//!   The replica does not try to tell why a decision has not come. The proposal may have been
206//!   forwarded to a process that has crashed, the decision may have been lost at a session ending,
207//!   or consensus for the slot may simply still be running. Every case is answered by the same
208//!   message and none is made unsafe by asking: where the slot is open the leader's `∄c'` guard
209//!   drops the repeat, where the proposal was lost the repeat is the first the leader hears of it,
210//!   and where the slot is decided the leader answers. So the threshold can be generous and wrong
211//!   without being unsafe. The one thing it must not be is absent.
212//!
213//! # A position is not a slot
214//!
215//! [`Position`] and [`Slot`] are both `u64` and are **not** interchangeable. Slots count from 1;
216//! `Position::START` is 0. More importantly `perform` skips a command already applied at a lower
217//! slot, so a slot can pass without a position being taken — the two diverge by exactly the number
218//! of commands decided more than once. `Position(slot)` is right until the first duplicate and then
219//! silently wrong, which is why `a_command_decided_in_two_slots_takes_one_position` drives that case
220//! deliberately rather than waiting for a run to produce it.
221//!
222//! # Scope
223//!
224//! This layer bridges nothing. A session ending reaches it from the child and is propagated in its
225//! own indications: its redundancy is the child's, and the child's is the other processes rather
226//! than anything that outlives a session. See `docs/conditional-guarantees.md`.
227//!
228//! Crash-stop, and one step stronger than the child's. `MultiPaxosSynod` keeps nothing durably, so
229//! a returning process has forgotten which ballots it took up; this layer adds that a returning
230//! process has forgotten its *sequence*, and would answer a read with a shorter one than it had
231//! already served. A total order that shortens is not a total order, so a crashed process here is
232//! crashed for good. §4.3 of the source is what changes it, and it is not this change.
233
234use core::time::Duration;
235use recon_core::{Child, NodeId, Position, ProtoCx, Protocol, Time, TimerId};
236use serde::{Deserialize, Serialize};
237use std::collections::{BTreeMap, BTreeSet, VecDeque};
238
239use crate::Timing;
240use crate::link::{Boundary, VolatileLink};
241use crate::multi_paxos_synod::{self as synod, MultiPaxosSynod, Slot, SynodMsg};
242use crate::session_link::SessionLink;
243use crate::total_order_log::{LogInd, TotalOrderLog};
244
245/// How many decided slots a replica keeps for the duplicate filter — §4.2's retention window.
246///
247/// Generously more than [`WINDOW`], because the filter has to outlive every slot a duplicate could
248/// span and the window only bounds how far ahead proposals may run. See
249/// `MultiPaxosReplica::retain`.
250pub const RETAIN: Slot = 64;
251
252/// How far `slot_in` may run ahead of `slot_out` — the source's `WINDOW`.
253///
254/// Large enough that several commanders run at once, which is what makes out-of-order decisions
255/// ordinary rather than incidental, and small enough that a stalled slot is visibly a stall.
256pub const WINDOW: Slot = 8;
257
258/// The source's `c = ⟨κ, cid, op⟩`: who asked, which request of theirs it is, and what it says.
259///
260/// Both identifying halves are load-bearing. `from` is what the port's
261/// [`LogInd::Ordered`] reports, and the page carries it for the same reason — the reply goes back
262/// to `κ`. `cid` is what keeps two appends of the same value from collapsing into one: `perform`
263/// skips a command it has already applied, and without `cid` a client appending `7` twice would see
264/// one entry.
265///
266/// **`cid`'s scope is this incarnation.** It is a counter in volatile state, exactly as the Synod
267/// ballot's round is, and a restarted process re-mints values it has already used. A request
268/// carrying a reused `⟨from, cid⟩` can be taken for one already applied and dropped, which is the
269/// identity rule's worked example: an identifier that crosses the wire outlives the handler that
270/// minted it, so its generator is state with a scope, and this one survives nothing. Making it
271/// durable is part of the fail-recovery change, not this one.
272#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
273pub struct Command<V> {
274    /// `κ` — the process that appended it.
275    pub from: NodeId,
276    /// `cid` — which request of that process's this is. Scope: this incarnation.
277    pub cid: u64,
278    /// `op` — what the command says. A value rather than an operation; see the departures.
279    pub value: V,
280}
281
282/// What the Synod protocol beneath carries for this layer.
283pub type Carried<V> = SynodMsg<Command<V>>;
284
285/// The consensus this layer runs over. A type alias rather than a parameter: the replica is the
286/// half of Multi-Paxos that Figure 1 describes, and the other half is what it is.
287pub type Synod<V, L> = MultiPaxosSynod<Command<V>, L>;
288
289/// Requests from the layer above.
290#[derive(Debug, Clone, PartialEq, Eq)]
291pub enum Cmd<V> {
292    /// `case ⟨request, c⟩` — append `value` to the log.
293    Append(V),
294    /// Read the ordered sequence from `from` onwards. The page has no such request; see the port's
295    /// own documentation for why it exists and what it does not promise.
296    Read { from: Position },
297}
298
299/// Indications to the layer above.
300#[derive(Debug, Clone, PartialEq, Eq)]
301pub enum Ind<V> {
302    /// A command took its place in the agreed sequence — the page's `op(state)`, in the vocabulary
303    /// of a log.
304    Ordered { position: Position, from: NodeId, value: V },
305    /// The answer to a [`Cmd::Read`].
306    Contents { from: Position, entries: Vec<V> },
307    /// The scope with `peer` ended at `epoch`, as the Synod protocol beneath reported it.
308    ///
309    /// Propagated rather than absorbed: this layer holds no redundancy that outlives a session.
310    SessionEnded { peer: NodeId, epoch: u64 },
311    /// A scope with `peer` is in force at `epoch`.
312    SessionEstablished { peer: NodeId, epoch: u64 },
313}
314
315/// A totally ordered log, built on Multi-Paxos: one consensus per slot, under a stable leader that
316/// keeps phase one across all of them.
317#[derive(Debug)]
318pub struct MultiPaxosReplica<V: Clone + Ord, L: VolatileLink<Carried<V>> = SessionLink<Carried<V>>>
319{
320    me: NodeId,
321    /// `slot_in` — the next slot this replica will propose for.
322    slot_in: Slot,
323    /// `slot_out` — the next slot this replica will apply. Never decreases; R4.
324    slot_out: Slot,
325    /// `requests` — appended and not yet given a slot. A queue rather than a set; see the
326    /// departures.
327    requests: VecDeque<Command<V>>,
328    /// `proposals` — given a slot and not yet applied.
329    proposals: BTreeMap<Slot, Command<V>>,
330    /// `decisions` — **append-only**. A slot's command never changes, because R1 says there is only
331    /// one, and an entry is never removed, because R2 reads `decisions` for every slot below
332    /// `slot_out` and R3 for the whole sequence. §4.2's watermark is what may remove one, and it is
333    /// a later change.
334    decisions: BTreeMap<Slot, Command<V>>,
335    /// `state`, as a log: the commands applied, in the order they were applied.
336    sequence: Vec<Command<V>>,
337    /// `{decisions[s] : s < slot_out}`, indexed. See the departures.
338    performed: BTreeSet<Command<V>>,
339    /// `cid`. Volatile, scope this incarnation — see [`Command`].
340    cid: u64,
341    /// `WINDOW`.
342    window: Slot,
343
344    // ---- the fourth liveness violation ----
345    /// When each outstanding proposal was last handed to the child.
346    proposed_at: BTreeMap<Slot, Time>,
347    /// How long a proposal may be outstanding before it is proposed again for the same slot.
348    repropose_after: Duration,
349    /// How often the outstanding set is swept.
350    sweep_every: Duration,
351    /// How often this replica tells the consensus beneath how far it has applied — §4.2's periodic
352    /// update, which is the only thing that lets anything below collect.
353    report_every: Duration,
354    /// When it last did.
355    last_report: Time,
356    /// How many decided slots this replica keeps for the duplicate filter — §4.2's "kept for a
357    /// certain amount of time, making the probability of duplicate execution negligible", measured
358    /// in slots rather than in seconds.
359    ///
360    /// Slots rather than time because a duplicate's risk is a function of how far apart two slots
361    /// deciding one command can be, not of how long the run has lasted: a fast run and a slow one
362    /// with the same slots in flight need the same window, and only the slot count says that.
363    retain: Slot,
364    /// The sweep's handle. Compared before acting, because an expiry is offered to every layer.
365    tick: Option<TimerId>,
366
367    synod: Child<Synod<V, L>>,
368}
369
370impl<V: Clone + Ord> MultiPaxosReplica<V> {
371    /// A replica among `peers`, over the session link the Synod protocol defaults to.
372    ///
373    /// The re-proposal threshold is derived rather than passed: it must exceed the time a decision
374    /// legitimately takes, whose worst case is a leadership change — `detect_after` for Ω to move,
375    /// then phase one, then phase two. `detect_after * 3` clears that with room, and it is far below
376    /// what filling a `WINDOW` of eight slots would take.
377    pub fn new(me: NodeId, peers: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
378        let peers: Vec<NodeId> = peers.into_iter().collect();
379        MultiPaxosReplica {
380            me,
381            slot_in: 1,
382            slot_out: 1,
383            requests: VecDeque::new(),
384            proposals: BTreeMap::new(),
385            decisions: BTreeMap::new(),
386            sequence: Vec::new(),
387            performed: BTreeSet::new(),
388            cid: 0,
389            window: WINDOW,
390            proposed_at: BTreeMap::new(),
391            repropose_after: timing.detect_after * 3,
392            sweep_every: timing.retransmit,
393            report_every: timing.heartbeat,
394            last_report: Time::ZERO,
395            retain: RETAIN,
396            tick: None,
397            synod: Child::new(MultiPaxosSynod::new(me, peers, timing)),
398        }
399    }
400
401    /// The window, for a test that needs to fill one without driving a hundred slots.
402    pub fn with_window(mut self, window: Slot) -> Self {
403        self.window = window;
404        self
405    }
406
407    /// The re-proposal threshold, for a test that needs to reach it inside a settle window.
408    pub fn with_repropose_after(mut self, after: Duration) -> Self {
409        self.repropose_after = after;
410        self
411    }
412
413    /// The retention window, for a test that needs to fill one without deciding sixty-four slots.
414    pub fn with_retain(mut self, retain: Slot) -> Self {
415        self.retain = retain;
416        self
417    }
418}
419
420impl<V: Clone + Ord, L: VolatileLink<Carried<V>>> MultiPaxosReplica<V, L> {
421    /// The ordered sequence as this process holds it.
422    pub fn entries(&self) -> impl Iterator<Item = &V> {
423        self.sequence.iter().map(|c| &c.value)
424    }
425
426    /// How many commands this process has applied — the next [`Position`], as a count.
427    pub fn len(&self) -> usize {
428        self.sequence.len()
429    }
430
431    pub fn is_empty(&self) -> bool {
432        self.sequence.is_empty()
433    }
434
435    /// `slot_out`: the next slot to apply. Diverges from [`MultiPaxosReplica::len`] by exactly the
436    /// number of commands decided in more than one slot.
437    pub fn slot_out(&self) -> Slot {
438        self.slot_out
439    }
440
441    /// `slot_in`: the next slot to propose for. R5 keeps it below `slot_out + WINDOW`.
442    pub fn slot_in(&self) -> Slot {
443        self.slot_in
444    }
445
446    /// How many appends are waiting for a slot.
447    pub fn waiting(&self) -> usize {
448        self.requests.len()
449    }
450
451    /// How many slots this replica has proposed for and not yet applied.
452    pub fn outstanding(&self) -> usize {
453        self.proposals.len()
454    }
455
456    /// The slots this replica has proposed for and not yet applied. R5's substantive form is
457    /// asserted over these — see the module's note on where the page's `<` reads as `≤`.
458    pub fn proposed_slots(&self) -> impl Iterator<Item = Slot> + '_ {
459        self.proposals.keys().copied()
460    }
461
462    /// How many decisions this replica holds. Grows with commands handled, and is never collected —
463    /// the measurement `docs/bounded-space.md` wants of a transcription.
464    pub fn decisions_held(&self) -> usize {
465        self.decisions.len()
466    }
467
468    /// The command decided for `slot`, if this replica has heard.
469    pub fn decision(&self, slot: Slot) -> Option<&Command<V>> {
470        self.decisions.get(&slot)
471    }
472
473    /// How far this replica has told the consensus beneath it that it has applied.
474    pub fn reported_slot_out(&self) -> Slot {
475        self.slot_out
476    }
477
478    /// The Synod protocol beneath, for a test that asks who leads.
479    pub fn synod(&self) -> &Synod<V, L> {
480        &self.synod
481    }
482
483    /// `function propose()`, minus the sending: move what it can from `requests` into `proposals`,
484    /// and hand back what the caller must give the child.
485    ///
486    /// Split in two because the send in Figure 1 goes to a set of leaders and here it goes into a
487    /// child, which needs `&mut self` for the duration — see `pump`.
488    ///
489    /// ```text
490    /// while slot_in < slot_out + WINDOW ∧ ∃c : c ∈ requests do
491    ///   if ∄c' : ⟨slot_in, c'⟩ ∈ decisions then
492    ///     requests := requests \ {c};
493    ///     proposals := proposals ∪ {⟨slot_in, c⟩};
494    ///     ∀λ ∈ leaders : send(λ, ⟨propose, slot_in, c⟩);
495    ///   end if
496    ///   slot_in := slot_in + 1;
497    /// end while
498    /// ```
499    ///
500    /// Note where `slot_in` advances: **outside** the `∄c'` arm, so a slot already decided is
501    /// stepped over without spending a request. And note the `while` guard is R5.
502    fn transfer(&mut self, now: Time, cx: &mut ProtoCx<'_, Self>) -> Vec<(Slot, Command<V>)> {
503        let mut out = Vec::new();
504        while !self.requests.is_empty() {
505            if self.slot_in >= self.slot_out + self.window {
506                // R5 refusing to go further. Nothing at all reaches the trace from this: a replica
507                // holding requests it may not propose for looks exactly like an idle one.
508                cx.note(crate::Note::WindowFull { slot_in: self.slot_in, slot_out: self.slot_out });
509                break;
510            }
511            if !self.decisions.contains_key(&self.slot_in) {
512                let command = self.requests.pop_front().expect("the loop guard");
513                self.proposals.insert(self.slot_in, command.clone());
514                self.proposed_at.insert(self.slot_in, now);
515                out.push((self.slot_in, command));
516            }
517            self.slot_in += 1;
518        }
519        out
520    }
521
522    /// `function perform(⟨κ, cid, op⟩)`, minus `op(state)` and the client response.
523    ///
524    /// ```text
525    /// if (∃s : s < slot_out ∧ ⟨s, ⟨κ, cid, op⟩⟩ ∈ decisions) ∨ isreconfig(op) then
526    ///   slot_out := slot_out + 1;
527    /// else
528    ///   ⟨next, result⟩ := op(state);
529    ///   atomic
530    ///     state := next; slot_out := slot_out + 1;
531    ///   end atomic
532    ///   send(κ, ⟨response, cid, result⟩);
533    /// end if
534    /// ```
535    ///
536    /// The already-decided arm is what makes a position not a slot: the slot passes, the sequence
537    /// does not grow. `isreconfig` is absent with reconfiguration.
538    fn perform(&mut self, command: Command<V>, cx: &mut ProtoCx<'_, Self>) {
539        if self.performed.contains(&command) {
540            self.slot_out += 1;
541            return;
542        }
543        let position = Position(self.sequence.len() as u64);
544        let Command { from, value, .. } = command.clone();
545        self.sequence.push(command.clone());
546        self.performed.insert(command);
547        // The page updates `state` and `slot_out` atomically and only then answers. Here the
548        // indication is what reveals the entry, so the sequence and `slot_out` both move first.
549        self.slot_out += 1;
550        cx.indicate(Ind::Ordered { position, from, value });
551    }
552
553    /// `case ⟨decision, s, c⟩`, in full.
554    ///
555    /// ```text
556    /// decisions := decisions ∪ {⟨s, c⟩};
557    /// while ∃c' : ⟨slot_out, c'⟩ ∈ decisions do
558    ///   if ∃c'' : ⟨slot_out, c''⟩ ∈ proposals then
559    ///     proposals := proposals \ {⟨slot_out, c''⟩};
560    ///     if c'' ≠ c' then
561    ///       requests := requests ∪ {c''};
562    ///     end if
563    ///   end if
564    ///   perform(c');
565    /// end while
566    /// ```
567    ///
568    /// The `while` is why a decision arriving out of order is held rather than dropped: slot 4
569    /// deciding before slot 3 leaves `slot_out` at 3, and slot 3's decision then drains both.
570    fn decided(&mut self, slot: Slot, command: Command<V>, cx: &mut ProtoCx<'_, Self>) {
571        // `decisions := decisions ∪ {⟨s, c⟩}` — a union, so a repeat writes nothing. The child
572        // announces a decision again whenever a later ballot re-commands the slot, and answers a
573        // re-proposal for a decided one deliberately; R1 makes both the same command.
574        self.decisions.entry(slot).or_insert(command);
575        self.proposed_at.remove(&slot);
576        while let Some(decided) = self.decisions.get(&self.slot_out).cloned() {
577            if let Some(mine) = self.proposals.remove(&self.slot_out) {
578                self.proposed_at.remove(&self.slot_out);
579                if mine != decided {
580                    // Somebody else's command took the slot. Ours goes back and is proposed again
581                    // at a later one — nothing at all reaches the trace from this decision, and a
582                    // replica that dropped it instead would lose an append silently.
583                    cx.note(crate::Note::ProposalDisplaced { slot: self.slot_out });
584                    self.requests.push_back(mine);
585                }
586            }
587            self.perform(decided, cx);
588        }
589        self.forget_below();
590    }
591
592    /// Liu et al.'s fourth violation, swept: anything outstanding past the threshold is proposed
593    /// again **for the same slot**.
594    ///
595    /// Undecided slots only. A slot whose decision has arrived is not stalled even if `slot_out`
596    /// has not reached it, because the drain loop will pass over it as soon as the slots below
597    /// decide.
598    fn resweep(&mut self, cx: &mut ProtoCx<'_, Self>) -> Vec<(Slot, Command<V>)> {
599        let now = cx.now();
600        let due: Vec<Slot> = self
601            .proposed_at
602            .iter()
603            .filter(|(slot, at)| {
604                **at + self.repropose_after <= now && !self.decisions.contains_key(slot)
605            })
606            .map(|(slot, _)| *slot)
607            .collect();
608        let mut out = Vec::new();
609        for slot in due {
610            let Some(command) = self.proposals.get(&slot).cloned() else { continue };
611            self.proposed_at.insert(slot, now);
612            // The child's `Propose` is a call when this process leads, so a re-proposal can reach
613            // the trace as nothing whatever. This is what says it happened.
614            cx.note(crate::Note::SlotReproposed { slot });
615            out.push((slot, command));
616        }
617        out
618    }
619
620    fn arm(&mut self, cx: &mut ProtoCx<'_, Self>) {
621        self.tick = Some(cx.set_timer(self.sweep_every));
622    }
623
624    /// §4.2's periodic update: tell the consensus beneath how far this replica has applied, so that
625    /// leaders and acceptors can discard what enough replicas already hold.
626    ///
627    /// **Periodic, and not a consequence of doing work.** A replica applying nothing is exactly the
628    /// one whose position the others most need — a run in which one replica is idle is a run in
629    /// which the watermark is pinned by it — and a report carried on this replica's own traffic
630    /// would fall silent precisely then. It rides the sweep's timer at its own coarser interval,
631    /// so it costs one message per member per `report_every` and nothing per entry: the cost
632    /// identity counts it as its own kind for that reason.
633    fn maybe_report(&mut self, cx: &mut ProtoCx<'_, Self>) {
634        let now = cx.now();
635        if self.last_report + self.report_every > now && now != Time::ZERO {
636            return;
637        }
638        self.last_report = now;
639        let slot_out = self.slot_out;
640        self.through_synod(cx, |s, ccx| s.on_cmd(synod::Cmd::Applied { slot_out }, ccx));
641    }
642
643    /// §4.2's retention window on the duplicate filter, which is the *other* half of that section
644    /// and a different mechanism from the watermark beneath.
645    ///
646    /// The source calls this state unavoidable and bounds it by time rather than by a watermark:
647    /// "it is often sufficient if such information is only kept for a certain amount of time,
648    /// making the probability of duplicate execution negligible".
649    ///
650    /// **Not the collection watermark, and this is the obvious wrong move.** That watermark says
651    /// `f + 1` replicas have *applied* up to a slot, which says nothing whatever about whether a
652    /// command decided below it may be decided again above it. A duplicate filter has to outlive
653    /// every slot a duplicate could span, and no watermark bounds that.
654    ///
655    /// What it costs is stated in the module: the no-duplication guarantee is scoped to the window,
656    /// and a command decided again more than `retain` slots later would take a second position.
657    fn forget_below(&mut self) {
658        let keep_from = self.slot_out.saturating_sub(self.retain);
659        if keep_from == 0 {
660            return;
661        }
662        let dropped: Vec<Command<V>> =
663            self.decisions.range(..keep_from).map(|(_, command)| command.clone()).collect();
664        self.decisions.retain(|slot, _| *slot >= keep_from);
665        for command in dropped {
666            self.performed.remove(&command);
667        }
668    }
669
670    /// Run `f` against the child, then handle everything that falls out of it — including the
671    /// proposals this replica makes in response, which go back into the same child.
672    ///
673    /// Figure 1 calls `propose()` at the bottom of its `for ever` loop, after every message, and
674    /// that is what the tail of this loop is. It iterates rather than recursing because handling a
675    /// decision can produce a proposal and a proposal could in principle produce an indication;
676    /// today it cannot — a decision needs a round trip — and this survives the day it can.
677    fn pump(
678        &mut self,
679        mut pending: Vec<synod::Ind<Command<V>>>,
680        cx: &mut ProtoCx<'_, Self>,
681        mut extra: Vec<synod::Cmd<Command<V>>>,
682    ) {
683        loop {
684            for ind in pending.drain(..) {
685                match ind {
686                    synod::Ind::Decision { slot, command } => self.decided(slot, command, cx),
687                    // §4.2: the slot is decided and the consensus beneath has collected it, so it
688                    // cannot answer. "Replicas can learn decisions … from one another" — ask one.
689                    synod::Ind::Collected { slot } => {
690                        cx.note(crate::Note::CaughtUpFrom { slot });
691                        extra.push(synod::Cmd::CatchUp { from_slot: slot });
692                    }
693                    // The other side of it. What this replica still holds is what the asker needs,
694                    // and it is the only place left that holds it.
695                    // Everything this replica still holds from that slot on, in one answer.
696                    //
697                    // Bounded by the retention window rather than by the proposal window: a
698                    // capped answer leaves the asker behind with nothing to ask again *with*,
699                    // because it asks only when a proposal of its own is refused and it has no
700                    // proposal for a slot somebody else filled. Measured — a `WINDOW`-sized answer
701                    // left the asker one entry short for ever.
702                    //
703                    // **A replica more than the retention window behind cannot be caught up**, and
704                    // nothing here can change that: nobody holds those decisions any more. That is
705                    // where a real deployment takes a snapshot, which is outside the paper.
706                    synod::Ind::CatchUpWanted { peer, from_slot } => {
707                        let teach: Vec<(Slot, Command<V>)> = self
708                            .decisions
709                            .range(from_slot..)
710                            .map(|(slot, command)| (*slot, command.clone()))
711                            .collect();
712                        for (slot, command) in teach {
713                            extra.push(synod::Cmd::Teach { to: peer, slot, command });
714                        }
715                    }
716                    // The child bridges no session ending and neither does this layer: its
717                    // redundancy is the other processes, which a session ending does not restore.
718                    synod::Ind::SessionEnded { peer, epoch } => {
719                        cx.indicate(Ind::SessionEnded { peer, epoch });
720                    }
721                    synod::Ind::SessionEstablished { peer, epoch } => {
722                        cx.indicate(Ind::SessionEstablished { peer, epoch });
723                    }
724                }
725            }
726            // `propose();`
727            let now = cx.now();
728            let mut outgoing: Vec<synod::Cmd<Command<V>>> = self
729                .transfer(now, cx)
730                .into_iter()
731                .map(|(slot, command)| synod::Cmd::Propose { slot, command })
732                .collect();
733            outgoing.append(&mut extra);
734            if outgoing.is_empty() {
735                break;
736            }
737            for cmd in outgoing {
738                let mut inds = self.synod.run(cx, |m| m, |s, ccx| s.on_cmd(cmd, ccx));
739                pending.append(&mut inds);
740                self.synod.reclaim(inds);
741            }
742            if pending.is_empty() {
743                break;
744            }
745        }
746        self.synod.reclaim(pending);
747    }
748
749    /// The composition, in the transforming form: a slot decision becomes a position in a sequence,
750    /// so the child's indications come back for this layer to handle rather than passing through.
751    fn through_synod(
752        &mut self,
753        cx: &mut ProtoCx<'_, Self>,
754        f: impl FnOnce(&mut Synod<V, L>, &mut ProtoCx<'_, Synod<V, L>>),
755    ) {
756        // The wrap is the identity: this layer adds no header, so its wire is the child's.
757        let inds = self.synod.run(cx, |m| m, f);
758        self.pump(inds, cx, Vec::new());
759    }
760}
761
762impl<V: Clone + Ord, L: VolatileLink<Carried<V>>> Protocol for MultiPaxosReplica<V, L> {
763    type Cmd = Cmd<V>;
764    type Ind = Ind<V>;
765    /// The child's. This layer adds no per-hop state, so it adds no wire field.
766    type Msg = <Synod<V, L> as Protocol>::Msg;
767    /// Whatever the link beneath the child is conditional on. This layer bridges none of it.
768    type Scope = <Synod<V, L> as Protocol>::Scope;
769    type Note = crate::Note;
770    /// Keeps nothing durably, which is what makes the crash-stop boundary in the module
771    /// documentation the one that applies. §4.3 is the change that alters it.
772    type Meta = core::convert::Infallible;
773    type Entry = core::convert::Infallible;
774
775    fn on_cmd(&mut self, cmd: Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
776        match cmd {
777            // `case ⟨request, c⟩ : requests := requests ∪ {c};` then `propose()`.
778            Cmd::Append(value) => {
779                let command = Command { from: self.me, cid: self.cid, value };
780                self.cid += 1;
781                self.requests.push_back(command);
782                self.pump(Vec::new(), cx, Vec::new());
783            }
784            // The departure the port records. Served from this process's own sequence, so it may
785            // lag an append that has completed elsewhere.
786            Cmd::Read { from } => {
787                let entries: Vec<V> =
788                    self.sequence.iter().skip(from.0 as usize).map(|c| c.value.clone()).collect();
789                cx.indicate(Ind::Contents { from, entries });
790            }
791        }
792    }
793
794    fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
795        self.through_synod(cx, |s, ccx| s.on_msg(from, msg, ccx));
796    }
797
798    /// An expiry is offered to every layer, so the child is given it and this layer acts only on
799    /// the handle it registered itself.
800    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
801        self.through_synod(cx, |s, ccx| s.on_timer(id, ccx));
802        if self.tick != Some(id) {
803            return;
804        }
805        self.arm(cx);
806        self.maybe_report(cx);
807        let due: Vec<synod::Cmd<Command<V>>> = self
808            .resweep(cx)
809            .into_iter()
810            .map(|(slot, command)| synod::Cmd::Propose { slot, command })
811            .collect();
812        self.pump(Vec::new(), cx, due);
813    }
814
815    /// `⟨ Init ⟩` — arm the sweep, and start the child, whose own `on_init` starts the leader
816    /// detector. A protocol is owed exactly one of `on_init` and `on_recovery` before its first
817    /// event, and the detector beneath is the one that has gone without twice.
818    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
819        self.arm(cx);
820        self.through_synod(cx, |s, ccx| s.on_init(ccx));
821    }
822
823    /// Hand the boundary to the child, which knows what it means. Leaving it to the trait's default
824    /// would take a scope event the driver raised and drop it.
825    fn on_scope_event(&mut self, scope: Self::Scope, cx: &mut ProtoCx<'_, Self>) {
826        self.through_synod(cx, |s, ccx| s.on_scope_event(scope, ccx));
827    }
828}
829
830impl<V: Clone + Ord, L: VolatileLink<Carried<V>>> TotalOrderLog<V> for MultiPaxosReplica<V, L> {
831    fn append(value: V) -> Cmd<V> {
832        Cmd::Append(value)
833    }
834
835    fn read(from: Position) -> Cmd<V> {
836        Cmd::Read { from }
837    }
838
839    fn classify(ind: Ind<V>) -> LogInd<V> {
840        match ind {
841            Ind::Ordered { position, from, value } => LogInd::Ordered { position, from, value },
842            Ind::Contents { from, entries } => LogInd::Contents { from, entries },
843            Ind::SessionEnded { peer, epoch } => LogInd::Boundary(Boundary::Ended { peer, epoch }),
844            Ind::SessionEstablished { peer, epoch } => {
845                LogInd::Boundary(Boundary::Established { peer, epoch })
846            }
847        }
848    }
849}