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}