recon_protocols/multi_paxos_synod.rs
1//! The Synod protocol of Multi-Paxos: ballots, acceptors, scouts, commanders and leaders.
2//!
3//! **Status: implementation. Space: bounded by membership and by the distance between the
4//! collection watermark and the frontier — see the section on it below, and the condition it
5//! carries.**
6//!
7//! van Renesse, R. and Altinbuken, D. (2015) 'Paxos Made Moderately Complex', *ACM Computing
8//! Surveys*, 47(3), pp. 1–36, §2. Figures 4, 6 and 7 are quoted above the code that implements
9//! them. **The edition matters**: the 2011 Cornell technical report of the same title numbers its
10//! figures differently and its acceptor differs materially from this one, so a reader checking the
11//! code against the page needs to know which page. See `CLAUDE.md`, Reference material.
12//!
13//! The cross-check is Liu, Y.A., Chand, S. and Stoller, S.D. (2019) 'Moderately Complex Paxos Made
14//! Simple', PPDP '19 — the same algorithm specified in DistAlgo with machine-checked TLA+ safety
15//! proofs. It reports four liveness violations in this specification **when messages can be lost**,
16//! which is this repository's setting rather than a hypothetical one, and one further issue in the
17//! acceptor. Three of the four and the acceptor issue are fixed here; each is cited where it
18//! departs from the survey's pseudocode.
19//!
20//! # What it guarantees, and what it does not
21//!
22//! ```text
23//! S1 [always] At most one proposal is ever chosen for a slot
24//! S2 [always] A chosen proposal is one some process proposed
25//! S3 [Ω settles] Progress — every proposed slot eventually has a proposal chosen
26//! ```
27//!
28//! S1 and S2 hold whatever the schedule, whatever the ballots in flight and whatever has crashed.
29//! S3 does not, and the source is blunt about it: §3 opens by observing that duelling leaders can
30//! preempt one another for ever and that the Synod protocol guarantees nothing about progress
31//! "even in the absence of any failure whatsoever".
32//!
33//! # Roles, and where each one lives
34//!
35//! §4.4 says the roles are co-located on one machine in practice, and that is the shape built here:
36//! one protocol per process holding **both** the acceptor and the leader. Neither has a vocabulary
37//! the other does not, and both send and receive on the same wire, so composing them as children
38//! would buy two wrap functions and a message split that carries nothing.
39//!
40//! Scouts and commanders are the leader's own bookkeeping rather than child protocols, for the
41//! question `link.rs` asks of everything: does this thing have a vocabulary of its own? It does
42//! not. What a scout is, concretely, is a `waitfor` set and a union of pvalues; what a commander
43//! is, is a `waitfor` set and the pvalue it is responsible for. The source models them as threads
44//! because it is written in a language where a thread is the cheapest way to say "wait for a
45//! majority". Here they are the leader's own `scout` and `commanders` fields, and the mapping
46//! from each figure's `for ever / switch receive` arm to where it went is stated above the
47//! handler that took it.
48//!
49//! The source's constraints on their number are in the types rather than in a comment: at most one
50//! scout, "and only for its own ballots", is an `Option`; at most one commander per `⟨ballot,
51//! slot⟩` (Invariant C1) is a map keyed by slot, cleared when the ballot changes.
52//!
53//! # Departure: a `p2b` names its slot
54//!
55//! Figure 4 answers a `p2a` with `⟨p2b, self(), ballot_num⟩` — no slot. It does not need one:
56//! the reply goes to the commander thread that sent the request, and the thread's identity is what
57//! routes it. Folding commanders into the leader removes that identity, and a leader running
58//! commanders for several slots at once cannot tell from `⟨p2b, α, b⟩` which of them answered. So
59//! the slot travels in the message.
60//!
61//! This is the structural decision above paying for itself, and it is the whole of the price: no
62//! guarantee changes, because the slot was already determined by the request the reply answers.
63//! `⟨p1b, α, b, r⟩` needs no such addition — a leader runs at most one scout, so its ballot routes
64//! it.
65//!
66//! # Departure: a reply naming a ballot *below* the attempt's is stale, and is discarded
67//!
68//! The same cause, and it is the half that bites. Figures 6(a) and 6(b) branch on `b' = b` and
69//! treat everything else as a preemption, which is sound for a thread: a reply to a request the
70//! thread did not send cannot reach it, because the thread that sent that request has exited and
71//! the reply is addressed to it. A scout that is a *field* has no such address, so a late answer
72//! to the **previous** attempt arrives at the current one.
73//!
74//! Taking the figure literally there is a liveness bug, and it was measured before this guard
75//! existed. A leader is preempted, climbs, and starts a scout for the new ballot; the second
76//! acceptor's answer to the *old* ballot then arrives, fails `b' = b`, and takes the `else` arm —
77//! which discards the running scout and reports a preemption the leader correctly ignores, because
78//! the ballot in it is below the one it now holds. The leader is left with no scout, not active,
79//! and nothing in the sweep to restart it. It stops for ever, while Ω goes on trusting it.
80//!
81//! So a reply is classified against the attempt it reaches: **above** its ballot is a preemption,
82//! **equal** is an answer, and **below** is an answer to an attempt that has already exited and is
83//! discarded. Under Figure 6's own model the third case cannot arise, which is why the figure does
84//! not name it.
85//!
86//! # Departure: the acceptor accepts under `b ≥ ballot_num`, and adopts what it accepts
87//!
88//! Figure 4 accepts a pvalue only when `b = ballot_num`, and Figure 6(a)'s commander treats any
89//! reply naming a different ballot as a preemption. The two compose only because §2.3 asserts that
90//! every `p2b` carries `b' ≥ b`, and **that assertion rests on the acceptor having seen phase 1
91//! before phase 2** — which the link beneath this module does not guarantee.
92//!
93//! The case is concrete rather than theoretical. An acceptor's `p1a` dies at a session ending; the
94//! scout completes with a majority that excludes it; the retransmitted `p2a` then reaches an
95//! acceptor whose `ballot_num` is below `b`. Under the survey's text that acceptor refuses,
96//! answers with its own lower ballot, and the commander reads the mismatch as a preemption and
97//! exits — while the leader ignores the preemption, because the ballot in it is *below* its own.
98//! The slot then has no commander, so the retry sweep no longer covers it, and nothing moves until
99//! the phase-two escalation notices a whole timeout later.
100//!
101//! So the acceptor here takes the 2011 report's condition — `if b ≥ ballot_num then ballot_num :=
102//! b; accepted := accepted ∪ {⟨b,s,c⟩}` — which is also the fix Liu et al. give for what they name
103//! the useless-replies issue in this acceptor. Every `p2b` then names a ballot at least as high as
104//! the request's, the commander's else-arm is a genuine preemption again, and the property §2.3
105//! asserts without support holds here by construction.
106//!
107//! **Safety is unaffected**, and this is the one place to be sure of it. Invariant A1 — an acceptor
108//! adopts strictly increasing ballots — still holds: `ballot_num` only ever moves up. A2 becomes
109//! "accepts `⟨b,s,c⟩` only where `b = ballot_num` **after** the adoption in the same transition",
110//! which is what the argument in §2.4 actually uses: what matters is that an acceptor which has
111//! adopted `b` never afterwards accepts below `b`, and adopting on the way in makes that *more*
112//! true rather than less. A4 is unaffected because it is enforced by the leader (C1), not the
113//! acceptor.
114//!
115//! # Departure: liveness comes from Ω, not from the source's pinging
116//!
117//! §3 has a preempted leader monitor the preempting one by pinging it on a regular basis, backing
118//! off with an AIMD timeout, and says outright that "this concept is called failure detection".
119//! This module composes [`EventualLeaderDetector`] instead and acts on `Trust`: it starts a scout
120//! for its next ballot when trusted, and stays passive otherwise. A preemption moves `ballot_num`
121//! past the ballot that beat it but starts no scout unless this process is still trusted, which is
122//! what stops the duel §3 opens by describing.
123//!
124//! What differs: Ω names **one** leader, where the source's scheme lets any correct leader win a
125//! race and merely makes the loser wait longer each time. Ω is the stronger assumption and the
126//! cheaper mechanism, and it is conditional on everything its own chain is conditional on —
127//! stated at each link in `docs/conditional-guarantees.md` rather than collapsed into one claim
128//! here. What it costs is that a correct process Ω does not trust will not lead however long it
129//! waits.
130//!
131//! # Colocation: the replica is on this machine, so two messages are this layer's
132//!
133//! §4.4 describes the deployment this module is built for: "each machine that runs a replica also
134//! runs a leader… the replica can send a proposal for a particular slot to its local leader, say λ,
135//! rather than broadcasting the request to all leaders. If λ is passive, monitoring another leader
136//! λ′, it forwards the proposal to λ′. If λ is active, it will start a commander."
137//!
138//! Taking that shape puts two messages here that a separated deployment would put above, and one of
139//! them is on this layer's own figure:
140//!
141//! - **A commander announces its decision to every process.** Figure 6(a)'s last line,
142//! `∀ρ ∈ replicas : send(ρ, ⟨decision, s, c⟩)`. Every process raises `Ind::Decision`, not only
143//! the one whose commander counted the majority, because a log above cannot be built from a
144//! decision one process holds. The addressees are the processes of the run: colocation makes the
145//! replica set and the acceptor set one here.
146//! - **A leader that cannot act on a proposal forwards it** to the process Ω trusts. Not on any
147//! figure — §4.4's sentence above. Without it a proposal made at a process Ω never trusts is
148//! never acted on at all.
149//!
150//! **What this costs, and it is not free: the delivery of a proposal now rests on the detector.** A
151//! proposal goes to one process where it used to reach every leader, so a request handed to a
152//! process whose detector names a leader that has crashed is *lost*. Nothing here recovers it — the
153//! layer above must ask again, and `multi_paxos_replica`'s re-proposal timeout is what does. Before
154//! colocation an inaccurate detector cost nothing for delivery; it now costs a request, and the
155//! suite drives that case rather than assuming it away.
156//!
157//! **A `Propose` that arrived is never forwarded again.** Two processes whose detectors disagree
158//! would otherwise pass one back and forth for as long as they disagree. The message carries no hop
159//! count and needs none: an arrived proposal is handled locally or dropped, so the path is at most
160//! two hops by construction, and the asker's own timeout is what recovers a drop.
161//!
162//! # Departure: a leader answers a re-proposal for a decided slot
163//!
164//! The leader-side half of the fourth liveness fix Liu et al. describe, and the half without which
165//! the replica's re-proposal is a loop rather than a recovery. A decision is announced **once**, by
166//! the commander that counted the majority, which then exits. A process the announcement never
167//! reached has no other way back: nothing here retransmits a decision once its commander is gone,
168//! the retry sweep walks the `waitfor` sets of *live* commanders, and the session link does not
169//! resend across an ending.
170//!
171//! So a leader keeps the slots it has seen decided, and answers a `Propose` for one with
172//! `⟨decision, s, c⟩` sent to the **asker alone** — one message, where the announcement was a
173//! fan-out. Liu et al. put it directly: a leader "can then work on deciding for that slot if a
174//! decision for it has not been made; otherwise, it can send back the decision for that slot".
175//! Without the answer, a re-proposal for a decided slot meets the `∄c'` guard, is dropped, and the
176//! asker re-proposes for ever.
177//!
178//! **The command in the answer comes from the decision, never from `proposals`.** The first draft
179//! took it from `proposals[slot]`, arguing that for a decided slot that is the decided command — it
180//! was the commanded value in the ballot that decided, and any later adoption's `pmax` writes it
181//! back. That argument holds for the leader that decided and for any leader that adopted afterwards,
182//! and it is **false** for the third case: a leader that commanded something else for the slot, was
183//! preempted, and never adopted again. Nothing rewrites that leader's `proposals`; the announcement
184//! of the real decision still marks the slot decided; and a forwarded re-proposal would then be
185//! answered with a command that was never chosen, which a replica whose detector names that process
186//! would apply. `the_answer_names_the_decided_command_not_the_answerers_own_stale_proposal` is the
187//! schedule. So `decided` keeps the command beside the slot, filled from the same two places every
188//! process learns a decision — its own commander counting a majority, and another's announcement —
189//! and R1 is what makes either source the right one. A consequence worth having: any process that
190//! knows the decision can answer, not only the leader that made it.
191//!
192//! **A leader that has yielded forwards even for a slot it remembered.** The same root, on the
193//! liveness side. A trusted process that has not yet adopted remembers a proposal for `adopted` to
194//! command; if it is preempted first and Ω has moved on, it yields with the entry still in
195//! `proposals`, and the `∄c'` guard would then drop every later proposal its own replica makes for
196//! that slot — never forwarded, never commanded by anyone, the one slot nobody else will propose
197//! for, and the re-proposal that exists for exactly this wedge defeated by the process's own memory.
198//! So the guard applies only where this process can act, and a passive, untrusted process forwards
199//! and forgets: what it remembered is dead weight, and anything a majority accepted comes back
200//! through `pmax` if it ever leads again.
201//! `a_proposal_remembered_by_a_leader_that_then_yields_is_forwarded_when_asked_again` is that
202//! schedule.
203//!
204//! # Departure: three liveness fixes from the cross-check
205//!
206//! Each is reachable here because this link loses messages, and each is Liu et al.'s.
207//!
208//! | Where | What is lost | What happens | What this leader does |
209//! |---|---|---|---|
210//! | Phase 1 | `p1a` | no `p1b` majority and no preemption ever arrives, so the leader waits for ever | restarts phase one after a timeout |
211//! | Phase 2 | `p2b` | no decision for that slot, and if it happens at every leader the layer above stalls too | resends `p2a` for that slot once a round trip has passed |
212//! | Phase 2 | `preempt` | a majority has moved to a higher ballot, so `p2a` can never reach one — the leader sends for ever and decides nothing | restarts **phase one** after a timeout |
213//!
214//! The third is the one a naive design gets wrong. Resending `p2a` cannot help once a majority
215//! holds a higher ballot; the leader has to go back to phase 1. A single retransmission sweep over
216//! unanswered requests recovers from the first two and loops for ever on the third.
217//!
218//! # Retransmission: how often it is *asked*, and how often it is *done*
219//!
220//! One timer, at `Timing::retransmit`, runs the sweep. It used to decide both questions, and that
221//! was the mistake: every suite here configures `retransmit` at half the simulator's delivery
222//! bound, so a request went out again before an answer to it could possibly have arrived. Measured
223//! over ten entries at five processes, phase two cost **3.6× what the algorithm needs** — two
224//! thirds of it the protocol talking over itself.
225//!
226//! So the sweep decides how often the question is asked and `resend_after`
227//! decides the answer. **`Timing::retransmit` below the delivery bound is not a tuning choice, it
228//! is a mistake** — the same shape as `detect_after`'s own note about exceeding the bound by a
229//! margin, and with the same remedy: state what the parameter has to exceed. Here the sweep may be
230//! as fine as you like, because it is no longer what sets the rate.
231//!
232//! A tick-driven resend is a **stubborn link's** idiom in the first place — resend because the
233//! network may have dropped it — and this module runs over a session link, which drops nothing
234//! while a session holds. The only loss is at a session ending, and the establishment that follows
235//! is when a resend can succeed. That event is what `resend_to` acts on, which
236//! makes the sweep a backstop rather than the mechanism. This is the first module in the repository
237//! to act on `SessionEstablished` rather than merely propagate it; `docs/conditional-guarantees.md`
238//! records what that obliges.
239//!
240//! **Both restarts rerun phase one under the same ballot.** For a lost `p1a` the rerun is
241//! idempotent: an acceptor that already adopted the ballot answers again, and the scout recollects.
242//! For a lost preemption the rerun is how the leader learns what it missed — acceptors that have
243//! moved answer `p1b` naming their higher ballot, the scout reports that as a preemption, and only
244//! then does `ballot_num` move. Minting a higher ballot on every timeout would discard phase-two
245//! work already accepted under the current ballot and learn nothing a rerun does not.
246//!
247//! The fourth violation is in the replica: if no decision arrives for a slot, every replica stops
248//! applying from that slot, `slot_out` stops moving, `WINDOW` fills and the system wedges. Its fix
249//! has two halves. The replica re-proposes after a timeout, which is `multi_paxos_replica`'s; and
250//! the leader answers a re-proposal for a decided slot, which is this module's and is the departure
251//! above.
252//!
253//! # Space
254//!
255//! **Bounded, and this is an implementation.** An acceptor keeps one pvalue per slot between the
256//! collection watermark and the frontier; a leader keeps one proposal and one decision record over
257//! the same range; and what each member last said it had applied is one entry per member. None of
258//! them grows with the slots handled, which is what `docs/bounded-space.md` asks of an
259//! implementation.
260//!
261//! Both of the source's reductions are applied. §4 opens by saying "the described protocol is not
262//! practical" and gives them: **§4.1**, one pvalue per slot rather than one per `⟨ballot, slot⟩`,
263//! is the section below; **§4.2**, collecting below a watermark, is the one after it.
264//!
265//! **The bound is conditional, and the condition is the source's own.** Collection needs `f + 1` of
266//! the `2 f + 1` members reporting their progress, so a run in which `f` have crashed collects
267//! nothing and grows again — "no garbage collection could be done", as §4.2 puts it. That is
268//! specified behaviour rather than a defect, and
269//! `collection_stalls_when_too_few_members_remain_to_report` is there so nobody later reads a stall
270//! as one.
271//!
272//! The set of decided slots grows the same way, and the announcement is work per decision rather
273//! than per tick: one fan-out when a commander completes, plus one directed answer per re-proposal
274//! for a slot already decided. Both are bounded by membership for a given slot and unbounded in
275//! slots, which is the same statement as everything else here.
276//!
277//! What *is* bounded is the retry sweep, and it is worth stating precisely because the obvious
278//! claim is wrong. The sweep visits the outstanding `waitfor` sets: each is bounded by membership,
279//! but their number is not — it grows with the slots proposed and not yet decided. So the sweep
280//! costs membership times the slots in flight. A decision retires its commander and leaves the
281//! sweep, so once the work completes the sweep is empty, which is the window
282//! `tests/common::assert_send_rate_flat!` measures. Nothing here resends history.
283//!
284//! # §4.1: an acceptor keeps only the highest-ballot pvalue per slot
285//!
286//! The first of the source's reductions, applied. §4.1 gives the reason in one sentence:
287//!
288//! > First, note that although a leader obtains for each slot a set of all accepted pvalues from a
289//! > majority of acceptors, it only needs to know if this set is empty or not, and if not, what the
290//! > maximum pvalue is. Thus, a large step toward practicality is that acceptors only maintain the
291//! > most recently accepted pvalue for each slot (`⊥` if no pvalue has been accepted) and return
292//! > only these pvalues in a `p1b` message to the scout. This gives the leader all information
293//! > needed to enforce Invariant C2.
294//!
295//! So `accepted` is keyed by slot, the scout reduces per slot as it collects, and `p1b` carries one
296//! entry per slot. What this removes is growth in the number of *ballots* a run has seen; what it
297//! leaves is growth in slots, which is §4.2's and not this section's.
298//!
299//! **The comparison is the leader's, not the acceptor's**, and the sentence quoted above splits it
300//! that way: an acceptor keeps "the most recently accepted pvalue", and the leader is what "needs
301//! to know … what the maximum pvalue is". An acceptor writes over its record with no comparison,
302//! because its own promise has already ordered the writes — a stored pvalue's ballot became
303//! `ballot_num` when it was stored, `ballot_num` never falls, and the `b ≥ ballot_num` arm admits
304//! nothing below it. The one case where an arriving ballot is not strictly above the stored one is
305//! equality, where A4 makes the command the same. A guard there would be a branch nothing can take.
306//!
307//! The scout is where the maximum is genuinely taken: two acceptors answering one phase one can
308//! report different ballots for one slot, and nothing orders their answers. That is `keep_max`, and
309//! it is load-bearing — a scout keeping the last arrival instead of the highest would hand `pmax` a
310//! command a lower ballot proposed, which is a split slot. **Nothing in the suite caught that**
311//! until this change measured it: the property used to be structural, carried by the old
312//! `⟨ballot, slot⟩` key's iteration order rather than by any code, and a structural property is one
313//! no test has to name. Moving the reduction to collection time made it code, and code gets a test.
314//!
315//! **Invariant A4 stops being structural, and does not stop being true.** The old key was
316//! `⟨ballot, slot⟩`, which made "at most one command per ballot and slot" the map's own property.
317//! A4 was never the acceptor's to enforce: Invariant C1 gives it — at most one commander per
318//! `⟨ballot, slot⟩` — and the leader is what holds C1. The old key bought a second, redundant
319//! enforcement and cost a dimension of growth.
320//!
321//! ## The record of a choice may be overwritten while the choice stands
322//!
323//! The paper raises this against its own reduction, and it is worth stating in full because a
324//! reader who assumes otherwise would take a correct run for a broken one:
325//!
326//! > This optimization leads to a worrisome effect. We know that when a majority of acceptors have
327//! > accepted the same pvalue `⟨b, s, c⟩`, then proposal `c` is chosen for slot `s`. Consider now
328//! > the following scenario. […] Acceptors `α₁` and `α₂` accept `⟨⟨0, λ⟩, 1, c⟩`, and thus proposal
329//! > `c` is chosen for slot 1 by ballot `⟨0, λ⟩`. However, leader `λ` crashes before learning this.
330//! > Now leader `λ′` gets acceptors `α₂` and `α₃` to adopt ballot `⟨0, λ′⟩`. After determining the
331//! > maximum pvalue among the responses, leader `λ′` has to select proposal `c`. Now suppose that
332//! > acceptor `α₂` accepts `⟨⟨0, λ′⟩, 1, c⟩`. At this point, there is no majority of acceptors that
333//! > store the same most recently accepted pvalue, and in fact no proof that ballot `⟨0, λ⟩` even
334//! > chose proposal `c`, as that part of the history has been overwritten.
335//!
336//! And the answer: "However, the leader of any ballot `b` after `⟨0, λ⟩` can only select
337//! `⟨b, 1, c⟩`. This is by Invariant C2 and because acceptors `α₁` and `α₂` both accepted
338//! `⟨⟨0, λ⟩, 1, c⟩` and together form a majority."
339//!
340//! What the reduction discards is **evidence**, not agreement. The fact outlives the record because
341//! every later ballot had to read the maximum from a majority, and any two majorities intersect —
342//! so the choice is carried forward by a chain of adoptions rather than by anything still stored.
343//! `agreement_survives_the_record_of_it_being_overwritten` drives exactly the paper's scenario and
344//! asserts the non-vacuity half from acceptor state: no majority still holds the chosen proposal at
345//! the instant the assertion is made.
346//!
347//! One consequence for the suite, and it is why the checker was built the way it was: **a checker
348//! reading acceptor state would now be wrong.** `tests/multi_paxos_synod.rs` feeds its checker from
349//! the trace — what was sent, what was indicated — so the reduction does not reach it.
350//!
351//! # §4.2: collecting what enough replicas already hold
352//!
353//! The second reduction, and what makes this an implementation. §4.2 separates two things, and only
354//! the second is this module's:
355//!
356//! > In our description of Paxos, much information is kept about slots that have already been
357//! > decided. Some of this is unavoidable. For example, because commands may be decided in multiple
358//! > slots, replicas each maintain a set of all decisions to filter out such duplicates. […] But
359//! > some state, and indeed some work that results from having that state, is unnecessary. The
360//! > leader maintains state for each slot in which it has a proposal. […] Similarly, even if the
361//! > state reduction of Section 4.1 is implemented, each acceptor maintains state for each slot.
362//! > However, once at least `f + 1` replicas have learned about the decision of some slot, it is no
363//! > longer necessary for leaders and acceptors to maintain this state — replicas can learn
364//! > decisions, and the application state that results from those decisions, from one another.
365//!
366//! So each replica reports its `slot_out` periodically, every process keeps what it was told, and
367//! everything below the highest slot `f + 1` members have applied to is discarded: the acceptor's
368//! pvalues, the leader's proposals, the record of which slots are decided. The replica's own
369//! duplicate filter is the *unavoidable* half and is bounded by a retention window instead — see
370//! [`crate::multi_paxos_replica`], which explains why no watermark bounds it.
371//!
372//! ## The hazard the section names, and the clause that answers it
373//!
374//! > However, we must prevent other leaders from mistakenly concluding that the acceptors have not
375//! > accepted any pvalues for the garbage-collected slots. To achieve this, the state of an
376//! > acceptor can be extended with a new variable that contains a slot number: all pvalues lower
377//! > than that slot number have been garbage collected. This slot number must be included in `p1b`
378//! > messages so that leaders can skip the lower numbered slots.
379//!
380//! Before collection, an acceptor reporting nothing for a slot meant nothing had been accepted for
381//! it. Afterwards it means one of two things, and `collected` is the only thing that tells them
382//! apart. A leader that read a collected slot as free would put a second command up for a slot
383//! already decided — a split slot, reached from the opposite direction to `synod-ignore-pmax`.
384//!
385//! **The skip is in two places, and the second is not on the page.** `adopted` drops proposals
386//! below the watermark it was told about, which covers what a leader already held; `on_propose`
387//! refuses a slot below the watermark, which covers one arriving *afterwards*. Only the first was
388//! written at first, and `agreement_holds_across_a_collection` split a slot on the second: the
389//! successor adopted cleanly and was then asked for a different command for a collected slot, and
390//! commanded it. The page does not need the second clause because there a `propose` comes from a
391//! replica whose `slot_in` never falls below its own `slot_out`; nothing in a *port* guarantees
392//! that, and a layer that assumes its caller is well behaved is not one.
393//!
394//! `synod-skip-collected` is the mutation registered against both.
395//!
396//! ## What collecting costs, and the transfer that pays for it
397//!
398//! Collecting at `f + 1` means `f` correct replicas may be behind — and everything that could have
399//! helped them is precisely what was discarded. This layer no longer holds the decision, and the
400//! leader's answer to a re-proposal needs that record. A correct process that missed one decision
401//! would be **stranded for ever**, which is what `a_replica_stranded_by_a_collection_catches_up_
402//! from_a_peer` drove before the transfer existed: one replica at `slot_out = 1` with the leader
403//! holding zero decisions.
404//!
405//! The section's own justification is the remedy — "replicas can learn decisions […] from one
406//! another" — so this layer carries the ask and the answer without holding either. A refused
407//! proposal raises [`Ind::Collected`]; the layer above asks with [`Cmd::CatchUp`]; the peer's layer
408//! above answers with [`Cmd::Teach`], from *its* decisions, and the answers arrive as the ordinary
409//! `⟨decision, s, c⟩` the replica already handles. A replica further behind than the retention
410//! window cannot be caught up at all, because nobody holds those decisions any more; that is where
411//! a deployment takes a snapshot, which is outside the paper.
412//!
413//! # The boundary this module does not cross
414//!
415//! This is **crash-stop**, which is the source's own model rather than a scope dodge. A crashed
416//! state machine there "will make no more transitions and thus its current state is fixed
417//! indefinitely", and a process that comes back off disk "is not theoretically considered
418//! crashed—it is simply slow for a while". There is no third case, and a process that returns
419//! having forgotten what it knew is the first one acting when it is not permitted to.
420//!
421//! The simulator can produce that case, and Ω will trust such a process again, so the consequence
422//! has to be stated rather than assumed away. The round counter is volatile: **its scope is this
423//! incarnation.** A process that restarts re-mints a ballot it has already used; an acceptor still
424//! holding that ballot accepts a second, different proposal under it; and two proposals accepted at
425//! one ballot and slot is exactly what Invariant A4 forbids and what the argument that two
426//! majorities agree depends on. A **durable ballot counter**, read back in `on_recovery`, is what
427//! makes leading again after a restart legal, and it belongs with the rest of §4.3 (*Keeping State
428//! on Disk*) in the fail-recovery change.
429//!
430//! A returning process could instead lead under a **new identity**, which is sound for the leader
431//! role — ballot uniqueness is a proposer-side obligation and the proposer set need not be fixed —
432//! and unsound for the acceptor role, because majorities of two different acceptor sets need not
433//! intersect. The source's own answer to that is reconfiguration (§5's Cheap Paxos "reconfigures
434//! the system replacing the suspected acceptor with a fresh one"), decided in a slot like any other
435//! command. There are no slots to decide it in until the replica exists, so the membership here is
436//! **fixed for the run**.
437//!
438//! # The figures
439//!
440//! ```text
441//! process Acceptor()
442//! var ballot_num := ⊥, accepted := ∅;
443//!
444//! for ever
445//! switch receive()
446//! case ⟨p1a, λ, b⟩ :
447//! if b > ballot_num then
448//! ballot_num := b;
449//! end if
450//! send(λ, ⟨p1b, self(), ballot_num, accepted⟩);
451//! end case
452//! case ⟨p2a, λ, ⟨b, s, c⟩⟩ :
453//! if b = ballot_num then
454//! accepted := accepted ∪ {⟨b, s, c⟩};
455//! end if
456//! send(λ, ⟨p2b, self(), ballot_num⟩);
457//! end case
458//! end switch
459//! end for
460//! end process
461//! ```
462//! Figure 4. Pseudocode for an acceptor. The `p2a` arm departs — see above.
463//!
464//! ```text
465//! process Commander(λ, acceptors, replicas, ⟨b, s, c⟩)
466//! var waitfor := acceptors;
467//!
468//! ∀α ∈ acceptors : send(α, ⟨p2a, self(), ⟨b, s, c⟩⟩);
469//! for ever
470//! switch receive()
471//! case ⟨p2b, α, b'⟩ :
472//! if b' = b then
473//! waitfor := waitfor − {α};
474//! if |waitfor| < |acceptors|/2 then
475//! ∀ρ ∈ replicas :
476//! send(ρ, ⟨decision, s, c⟩);
477//! exit();
478//! end if
479//! else
480//! send(λ, ⟨preempted, b'⟩);
481//! exit();
482//! end if
483//! end case
484//! end switch
485//! end for
486//! end process
487//!
488//! process Scout(λ, acceptors, b)
489//! var waitfor := acceptors, pvalues := ∅;
490//!
491//! ∀α ∈ acceptors : send(α, ⟨p1a, self(), b⟩);
492//! for ever
493//! switch receive()
494//! case ⟨p1b, α, b', r⟩ :
495//! if b' = b then
496//! pvalues := pvalues ∪ r;
497//! waitfor := waitfor − {α};
498//! if |waitfor| < |acceptors|/2 then
499//! send(λ, ⟨adopted, b, pvalues⟩);
500//! exit();
501//! end if
502//! else
503//! send(λ, ⟨preempted, b'⟩);
504//! exit();
505//! end if
506//! end case
507//! end switch
508//! end for
509//! end process
510//! ```
511//! Figure 6. (a) a commander, (b) a scout.
512//!
513//! ```text
514//! process Leader(acceptors, replicas)
515//! var ballot_num = (0, self()), active = false, proposals = ∅;
516//!
517//! spawn(Scout(self(), acceptors, ballot_num));
518//! for ever
519//! switch receive()
520//! case ⟨propose, s, c⟩ :
521//! if ∄c' : ⟨s, c'⟩ ∈ proposals then
522//! proposals := proposals ∪ {⟨s, c⟩};
523//! if active then
524//! spawn(Commander(self(), acceptors, replicas, ⟨ballot_num, s, c⟩));
525//! end if
526//! end if
527//! end case
528//! case ⟨adopted, ballot_num, pvals⟩ :
529//! proposals := proposals ◁ pmax(pvals);
530//! ∀⟨s, c⟩ ∈ proposals :
531//! spawn(Commander(self(), acceptors, replicas, ⟨ballot_num, s, c⟩));
532//! active := true;
533//! end case
534//! case ⟨preempted, ⟨r', λ'⟩⟩ :
535//! if ⟨r', λ'⟩ > ballot_num then
536//! active := false;
537//! ballot_num := (r' + 1, self());
538//! spawn(Scout(self(), acceptors, ballot_num));
539//! end if
540//! end case
541//! end switch
542//! end for
543//! end process
544//! ```
545//! Figure 7. Pseudocode skeleton for a leader. The `spawn(Scout(…))` calls depart — see above.
546
547use core::marker::PhantomData;
548use core::time::Duration;
549use recon_core::{Child, NodeId, ProtoCx, Protocol, Time, TimerId};
550use serde::{Deserialize, Serialize};
551use std::collections::{BTreeMap, BTreeSet};
552
553use crate::eventual_leader_detector::{self as eld, EventualLeaderDetector};
554use crate::link::{Boundary, LinkInd, VolatileLink};
555
556use crate::perfect_failure_detector::Heartbeat;
557use crate::session_link::SessionLink;
558use crate::{Note, Timing};
559
560/// A slot number. Orthogonal to a ballot number, as Figure 5 has it: one ballot can decide many
561/// slots, and one slot may be targeted by many ballots.
562pub type Slot = u64;
563
564/// A ballot number: `⟨round, leader⟩`, ordered lexicographically.
565///
566/// §2.2 makes ballot numbers "lexicographically ordered pairs of an integer and its leader
567/// identifier (consequently, leader identifiers need to be totally ordered)". Two consequences the
568/// derived ordering gives for free: any two ballots are comparable, and a ballot names its leader,
569/// so it is trivial to see who owns one.
570///
571/// The source's `⊥` is `Option::<Ballot>::None`, whose derived ordering already places it below
572/// every `Some` — which is exactly "ordered before any normal ballot number".
573///
574/// **The round counter is volatile, and its scope is this incarnation.** See the module
575/// documentation on the boundary this module does not cross.
576#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
577pub struct Ballot {
578 /// The integer half. Increases when this process is preempted.
579 pub round: u64,
580 /// The leader half, which makes two processes' ballots distinct however far their rounds drift.
581 pub leader: NodeId,
582}
583
584impl Ballot {
585 /// `(0, self())` — Figure 7's initial ballot number.
586 pub fn initial(leader: NodeId) -> Self {
587 Ballot { round: 0, leader }
588 }
589
590 /// `(r' + 1, self())` — the ballot to take up after being preempted by `beat_by`.
591 pub fn above(beat_by: Ballot, me: NodeId) -> Self {
592 Ballot { round: beat_by.round + 1, leader: me }
593 }
594}
595
596impl core::fmt::Display for Ballot {
597 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
598 write!(f, "({},{})", self.round, self.leader)
599 }
600}
601
602/// `p = ⟨b, s, c⟩` — a ballot number, a slot number and a command.
603#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
604pub struct Pvalue<C> {
605 pub ballot: Ballot,
606 pub slot: Slot,
607 pub command: C,
608}
609
610/// The pvalues held for a set of slots: **one per slot**, the one carrying the highest ballot.
611///
612/// The book writes a set and unions into it, and this module first held one — a `BTreeMap` keyed by
613/// `⟨ballot, slot⟩`, which made Invariant A4 the map's own property. §4.1 is why it no longer does;
614/// see `MultiPaxosSynod::accepted`. A4 is unaffected, because it was never the acceptor's to
615/// enforce: Invariant C1 — at most one commander per `⟨ballot, slot⟩` — is what gives it, and the
616/// leader is what holds C1. What the old key bought was a second, redundant enforcement at the
617/// acceptor; what it cost was a whole dimension of growth.
618pub type Pvalues<C> = BTreeMap<Slot, (Ballot, C)>;
619
620/// Union `pvalue` into `into`, keeping the higher ballot for the slot.
621///
622/// **The scout's, and only the scout's.** §4.1 splits the work in two: an acceptor keeps "the most
623/// recently accepted pvalue for each slot", and the leader takes the maximum across the majority
624/// that answers it. The acceptor needs no comparison — see `MultiPaxosSynod::on_p2a`, where its
625/// own promise already makes the latest acceptance the highest. The scout does, because two
626/// acceptors can report different ballots for one slot and nothing orders their answers: they
627/// arrive as the network delivers them.
628///
629/// This is where the safety argument's `pmax` really happens, so it is where a mistake is a split
630/// slot rather than a tidiness question. `a_scout_keeps_the_highest_ballot_reported_for_a_slot_not
631/// _the_last_one_to_arrive` is the test, and it was written because a mutation showed the whole
632/// suite green with this reduced to a plain insert.
633fn keep_max<C>(into: &mut Pvalues<C>, ballot: Ballot, slot: Slot, command: C) {
634 match into.get(&slot) {
635 Some((held, _)) if *held >= ballot => {}
636 _ => {
637 into.insert(slot, (ballot, command));
638 }
639 }
640}
641
642/// What this layer puts on the wire, beneath the link.
643///
644/// Figures 4 and 6: the acceptor's two requests, its two replies, and the commander's `decision`.
645/// `adopted` and `preempted` are not here — each is a thread reporting to the leader that owns it,
646/// which is a function call once the thread is a field.
647///
648/// [`SynodMsg::Propose`] is not on any figure. It is §4.4's colocation: a replica hands its
649/// proposal to the leader on its own machine, and a leader that cannot act on it passes it to the
650/// one that can. See the module documentation.
651#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
652pub enum SynodMsg<C> {
653 /// `⟨p1a, λ, b⟩` — phase one, from a scout.
654 P1a { ballot: Ballot },
655 /// `⟨p1b, α, ballot_num, accepted⟩` — an acceptor's answer, carrying **one pvalue per slot** it
656 /// has accepted for rather than everything it has ever accepted. §4.1; see the module
657 /// documentation. This message is what grew fastest in the book's version, because it grew with
658 /// the ballots the run had seen as well as with the slots.
659 P1b { ballot: Ballot, accepted: Vec<Pvalue<C>>, collected: Slot },
660 /// `⟨p2a, λ, ⟨b, s, c⟩⟩` — phase two, from a commander.
661 P2a { pvalue: Pvalue<C> },
662 /// `⟨p2b, α, ballot_num⟩`, **plus the slot it answers for**. Figure 4 has no slot, because the
663 /// reply goes to a thread whose identity supplies it; a leader holding its commanders as fields
664 /// needs it in the message. See the module documentation.
665 P2b { ballot: Ballot, slot: Slot },
666 /// `⟨decision, s, c⟩` — Figure 6(a)'s last line, sent by a commander that has counted a
667 /// majority, so that every process learns what was chosen rather than only the one that
668 /// counted.
669 ///
670 /// Also the answer to a [`SynodMsg::Propose`] for a slot already decided, sent to the asker
671 /// alone. A decision is announced once and its commander then exits, so a process the
672 /// announcement never reached has no other way back; asking is the way, and this is the answer.
673 Decision { slot: Slot, command: C },
674 /// A request for the decisions from `from_slot` onwards, from a process that has fallen behind
675 /// and whose consensus layer has collected them. §4.2's replica-to-replica transfer; see
676 /// [`Cmd::CatchUp`].
677 CatchUp { from_slot: Slot },
678 /// `⟨applied, ρ, s⟩` — how far the replica on the sending machine has applied. Not on any
679 /// figure; §4.2's periodic update, which is what makes collection possible at all.
680 Applied { slot_out: Slot },
681 /// A proposal for a slot, from a replica on this machine or from a leader that could not act on
682 /// it. Not on any figure — §4.4's colocation; see the module documentation.
683 ///
684 /// **A `Propose` that arrived is never forwarded again.** Two processes whose detectors
685 /// disagree would otherwise pass one back and forth for as long as they disagree.
686 Propose { slot: Slot, command: C },
687}
688
689/// The wire, multiplexing the leader detector and the Synod protocol itself.
690#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
691pub enum Wire<M> {
692 /// The leader detector's heartbeats.
693 Detector(Heartbeat),
694 /// The Synod protocol's traffic, as the link beneath wraps it.
695 Synod(M),
696}
697
698/// Requests from the layer above.
699#[derive(Debug, Clone, PartialEq, Eq)]
700pub enum Cmd<C> {
701 /// `⟨propose, s, c⟩` — propose `command` for `slot`.
702 Propose { slot: Slot, command: C },
703 /// How far the replica on this machine has applied — §4.2's "each replica periodically updates
704 /// leaders and acceptors about its `slot out` variable".
705 ///
706 /// A command rather than a message, because §4.4 puts the replica on this machine: it tells its
707 /// own leader and acceptor by a call, and this layer is what tells everybody else's.
708 Applied { slot_out: Slot },
709 /// Ask a peer that is ahead for the decisions from `from_slot` onwards.
710 ///
711 /// §4.2's other half: "replicas can learn decisions, and the application state that results
712 /// from those decisions, from one another". Once a slot is collected this layer cannot answer
713 /// for it, and a replica that missed it has nowhere else to ask.
714 CatchUp { from_slot: Slot },
715 /// Answer a peer's catch-up with one decision the replica above still holds.
716 ///
717 /// The decision comes from the *replica's* store, not this layer's: this layer collected it,
718 /// which is what provoked the request.
719 Teach { to: NodeId, slot: Slot, command: C },
720}
721
722/// Indications to the layer above.
723#[derive(Debug, Clone, PartialEq, Eq)]
724pub enum Ind<C> {
725 /// `send(ρ, ⟨decision, s, c⟩)` — a majority of acceptors accepted this pvalue, so `command` is
726 /// chosen for `slot`. Figure 6(a) addresses this to the replicas; there is no replica here, so
727 /// it is raised to whatever is above.
728 Decision { slot: Slot, command: C },
729 /// A proposal was refused because the slot is below this process's collection watermark: it is
730 /// decided, `f + 1` members hold the decision, and this layer no longer does.
731 ///
732 /// **The layer above must catch up from a peer**, which is what [`Cmd::CatchUp`] asks for.
733 /// Without it a correct process that missed one decision is stranded for ever, because the only
734 /// other way back — a leader answering a re-proposal — needs the record this layer discarded.
735 Collected { slot: Slot },
736 /// `peer` has fallen behind and wants the decisions from `from_slot` onwards. The layer above
737 /// answers with [`Cmd::Teach`] for what it still holds.
738 CatchUpWanted { peer: NodeId, from_slot: Slot },
739 /// The scope with `peer` ended at `epoch`. **Propagated, not absorbed**: this layer holds no
740 /// redundancy that outlives a session — what it knows of a ballot is in memory a crash takes,
741 /// and its redundancy is the other processes rather than a resend across an ending. A layer
742 /// above may need to know that an answer it was waiting for will not arrive.
743 SessionEnded { peer: NodeId, epoch: u64 },
744 /// A scope with `peer` is in force at `epoch`.
745 SessionEstablished { peer: NodeId, epoch: u64 },
746}
747
748/// Figure 6(b)'s scout, as a field rather than a thread.
749///
750/// `var waitfor := acceptors, pvalues := ∅` is the whole of its state. Its `∀α ∈ acceptors :
751/// send(α, ⟨p1a, self(), b⟩)` prologue is in [`MultiPaxosSynod::start_scout`]; its `⟨p1b, α, b',
752/// r⟩` arm is in [`MultiPaxosSynod::on_p1b`], whose two branches are the figure's `then` and
753/// `else`; and `exit()` is the `Option` becoming `None`.
754#[derive(Debug)]
755struct Scout<C> {
756 /// `b` — the ballot this scout is running phase one for.
757 ballot: Ballot,
758 /// `waitfor` — the acceptors that have not yet answered.
759 waitfor: BTreeSet<NodeId>,
760 /// `pvalues` — the union of what the answers carried.
761 pvalues: Pvalues<C>,
762 /// The highest slot any answering acceptor said it had collected below. §4.2: a leader must
763 /// "skip the lower numbered slots", or it reads a collected slot as one nothing was accepted
764 /// for. The **highest**, because a slot collected at any acceptor of the answering majority is
765 /// one whose decision `f + 1` replicas hold.
766 collected: Slot,
767 /// When phase one began, for the escalation. Not the figure's: the figure waits for ever.
768 started: Time,
769 /// When its requests last went out, so a resend is timed against the delivery bound rather than
770 /// against how often the sweep happens to run. See `resend_after`.
771 last_sent: Time,
772}
773
774/// Figure 6(a)'s commander, as a field rather than a thread.
775///
776/// `var waitfor := acceptors`, plus the pvalue it is responsible for, which the figure passes as a
777/// constructor argument. Its prologue is in [`MultiPaxosSynod::start_commander`]; its `⟨p2b, α,
778/// b'⟩` arm is in [`MultiPaxosSynod::on_p2b`]; `exit()` is removal from the map.
779#[derive(Debug)]
780struct Commander<C> {
781 /// `b` — the ballot this commander is running phase two under.
782 ballot: Ballot,
783 /// `c` — the command it is trying to get chosen. Held rather than looked up in `proposals`,
784 /// because it is the figure's own constructor argument and because a resend needs it.
785 command: C,
786 /// `waitfor` — the acceptors that have not yet answered.
787 waitfor: BTreeSet<NodeId>,
788 /// When phase two began for this slot, for the escalation.
789 started: Time,
790 /// When its requests last went out. See `resend_after`.
791 last_sent: Time,
792}
793
794/// `pmax(pvalues) ≡ {⟨s, c⟩ | ∃b : ⟨b, s, c⟩ ∈ pvalues ∧ ∀b', c' : ⟨b', s, c'⟩ ∈ pvalues ⇒ b' ≤ b}`
795///
796/// For each slot, the command of the pvalue with the maximum ballot number. Invariant A4 makes that
797/// command unique — there cannot be two different commands for the same ballot and slot — which is
798/// why taking one is well defined rather than a choice.
799///
800/// **After §4.1 this is the identity on the map's commands**, because the collection that built the
801/// map already kept only the maximum per slot — see [`keep_max`]. It stays because what it names is
802/// the algorithm's step, Figure 7's `proposals := proposals ◁ pmax(pvals)`, and a reader checking
803/// the code against the page needs to find it. Where the maximum is now *taken* is the one thing
804/// that moved.
805fn pmax<C: Clone>(pvalues: &Pvalues<C>) -> BTreeMap<Slot, C> {
806 pvalues.iter().map(|(slot, (_, command))| (*slot, command.clone())).collect()
807}
808
809/// `|waitfor| < |acceptors|/2` — the figures' majority test.
810///
811/// Written as a multiplication because the figures' division is over the reals. In Rust
812/// `waitfor.len() < acceptors.len() / 2` truncates: with five acceptors it reads `< 2`, which
813/// demands four answers rather than three, so a run losing two acceptors would never adopt or
814/// decide. Safety would survive that and liveness would not, silently.
815fn is_majority(waitfor: usize, acceptors: usize) -> bool {
816 waitfor * 2 < acceptors
817}
818
819/// The Synod protocol: acceptor and leader in one process, as §4.4 co-locates them.
820///
821/// `L` is the link beneath, a parameter rather than a fixed type, and it defaults to
822/// [`SessionLink`] rather than to the perfect link. That is deliberate: the real-world set's first
823/// obligation is running over a session link, and this is the first module here that can meet it
824/// before joining rather than after. A boundary is classified and propagated; this layer bridges
825/// nothing.
826#[derive(Debug)]
827pub struct MultiPaxosSynod<C: Clone, L: VolatileLink<SynodMsg<C>> = SessionLink<SynodMsg<C>>> {
828 me: NodeId,
829 /// The acceptor set, over which majorities are counted. **Fixed for the run** — see the module
830 /// documentation on why changing it without deciding the change is unsound.
831 acceptors: BTreeSet<NodeId>,
832
833 // ---- the acceptor, Figure 4 ----
834 /// `α.ballot_num`, initially `⊥`.
835 ballot_num: Option<Ballot>,
836 /// The slot below which this acceptor has discarded its pvalues — §4.2's "new variable that
837 /// contains a slot number: all pvalues lower than that slot number have been garbage
838 /// collected". Travels in every `p1b`, and only ever rises.
839 collected: Slot,
840 /// How far each member last said it had applied — §4.2's periodic update. One entry per member,
841 /// so bounded by membership by construction.
842 reported: BTreeMap<NodeId, Slot>,
843 /// `α.accepted`, initially `∅` — **one pvalue per slot**, the one carrying the highest ballot.
844 ///
845 /// §4.1, quoted in the module documentation. The book keeps every pvalue ever accepted; a
846 /// leader reads only the maximum per slot, so everything below it is read by nothing. Still
847 /// grows with the slots handled, which is why this is still a transcription; what it no longer
848 /// grows with is the ballots the run has seen.
849 accepted: Pvalues<C>,
850
851 // ---- the leader, Figure 7 ----
852 /// `λ.ballot_num`, initially `(0, self())`.
853 leader_ballot: Ballot,
854 /// `λ.active`.
855 active: bool,
856 /// `λ.proposals` — at most one entry per slot, which the map makes structural.
857 proposals: BTreeMap<Slot, C>,
858 /// The scout, if one is running. `Option` rather than a set, because a leader "starts at most
859 /// one of these for any ballot `b`, and only for its own ballots".
860 scout: Option<Scout<C>>,
861 /// The commanders, keyed by slot **within the current ballot** — Invariant C1. Cleared when the
862 /// ballot changes, which is what makes the key a slot rather than a `⟨ballot, slot⟩` pair.
863 commanders: BTreeMap<Slot, Commander<C>>,
864 /// Every decision this process has learned, so that a `Propose` for a decided slot can be
865 /// answered with the decision instead of ignored. Not on any figure; see the module
866 /// documentation on why the answer is what makes a re-proposal a recovery rather than a loop,
867 /// and on why the command lives here rather than being read back out of `proposals`.
868 ///
869 /// Grows with slots decided. That is the same growth `proposals` already has, so it changes
870 /// nothing about this module's stated bound, and §4.2's watermark collects both.
871 decided: BTreeMap<Slot, C>,
872 /// Whether Ω currently trusts this process. The departure from §3: a leader is passive by
873 /// default and competes only while trusted.
874 ///
875 /// The **identity** rather than a boolean, because §4.4's forward needs to know *who* to hand a
876 /// proposal to, not merely that it is not this process. `None` until Ω first speaks.
877 trusted: Option<NodeId>,
878
879 // ---- retries ----
880 /// The periodic sweep's handle. Compared before acting, since an expiry is offered to every
881 /// layer and this one composes two children that register their own.
882 tick: Option<TimerId>,
883 /// How often the sweep runs, and so how often an unanswered request is resent.
884 retry: Duration,
885 /// How long an attempt may go without adopting, deciding or being preempted before it restarts
886 /// phase one. Taken from the detector's own timeout rather than adding a knob: it is already
887 /// configured as "long enough that silence means something is wrong", which is exactly the
888 /// question being asked.
889 escalate_after: Duration,
890
891 omega: Child<EventualLeaderDetector>,
892 link: Child<L>,
893 _command: PhantomData<fn() -> C>,
894}
895
896impl<C: Clone, L: VolatileLink<SynodMsg<C>>> MultiPaxosSynod<C, L> {
897 /// The Synod protocol among `acceptors`, over a link the caller supplies.
898 pub fn with_link(
899 me: NodeId,
900 acceptors: impl IntoIterator<Item = NodeId>,
901 timing: Timing,
902 link: L,
903 ) -> Self {
904 let Timing { retransmit, heartbeat, detect_after } = timing;
905 let mut acceptors: BTreeSet<NodeId> = acceptors.into_iter().collect();
906 acceptors.insert(me);
907 MultiPaxosSynod {
908 me,
909 omega: Child::new(EventualLeaderDetector::new(
910 me,
911 acceptors.clone(),
912 heartbeat,
913 detect_after,
914 )),
915 acceptors,
916 ballot_num: None,
917 collected: 0,
918 reported: BTreeMap::new(),
919 accepted: BTreeMap::new(),
920 leader_ballot: Ballot::initial(me),
921 active: false,
922 proposals: BTreeMap::new(),
923 scout: None,
924 commanders: BTreeMap::new(),
925 decided: BTreeMap::new(),
926 trusted: None,
927 tick: None,
928 retry: retransmit,
929 escalate_after: detect_after,
930 link: Child::new(link),
931 _command: PhantomData,
932 }
933 }
934
935 /// `α.ballot_num` — the ballot this process has adopted as an acceptor, or `⊥`.
936 pub fn adopted_ballot(&self) -> Option<Ballot> {
937 self.ballot_num
938 }
939
940 /// `λ.ballot_num` — the ballot this process is leading with.
941 pub fn leader_ballot(&self) -> Ballot {
942 self.leader_ballot
943 }
944
945 /// `λ.active` — whether phase one has completed for the current ballot.
946 pub fn is_active(&self) -> bool {
947 self.active
948 }
949
950 /// Whether Ω currently trusts this process.
951 pub fn is_trusted(&self) -> bool {
952 self.trusted == Some(self.me)
953 }
954
955 /// Who Ω currently trusts, if it has spoken.
956 pub fn trusted_leader(&self) -> Option<NodeId> {
957 self.trusted
958 }
959
960 /// The slots this leader has seen decided. Grows with slots decided, exactly as `proposals`
961 /// does; §4.2 collects both.
962 pub fn decided_count(&self) -> usize {
963 self.decided.len()
964 }
965
966 /// How many pvalues this acceptor holds — after §4.1, the number of **slots** it has accepted
967 /// for, not the number of `⟨ballot, slot⟩` pairs. Grows with the slots handled, which is the
968 /// measurement `docs/bounded-space.md` wants for a transcription.
969 pub fn accepted_count(&self) -> usize {
970 self.accepted.len()
971 }
972
973 /// The slots this leader currently has a commander for.
974 pub fn commanded_slots(&self) -> impl Iterator<Item = Slot> + '_ {
975 self.commanders.keys().copied()
976 }
977
978 /// The pvalue this acceptor holds for `slot`, if any — after §4.1, at most one, carrying the
979 /// highest ballot it has accepted for that slot.
980 pub fn accepted_for(&self, slot: Slot) -> Option<(Ballot, &C)> {
981 self.accepted.get(&slot).map(|(ballot, command)| (*ballot, command))
982 }
983
984 /// The slot below which this process has discarded its state — §4.2's watermark, as it has
985 /// been acted on. Only ever rises.
986 pub fn collected_below(&self) -> Slot {
987 self.collected
988 }
989
990 /// How many members have said how far they have applied. Bounded by membership; the
991 /// measurement `docs/bounded-space.md` wants beside the two that are now bounded.
992 pub fn reports_held(&self) -> usize {
993 self.reported.len()
994 }
995
996 /// Whether a scout is running phase one.
997 pub fn is_scouting(&self) -> bool {
998 self.scout.is_some()
999 }
1000}
1001
1002impl<C: Clone> MultiPaxosSynod<C, SessionLink<SynodMsg<C>>> {
1003 /// The Synod protocol among `acceptors`, over the session link this module defaults to.
1004 pub fn new(me: NodeId, acceptors: impl IntoIterator<Item = NodeId>, timing: Timing) -> Self {
1005 Self::with_link(me, acceptors, timing, SessionLink::new())
1006 }
1007}
1008
1009impl<C, L> MultiPaxosSynod<C, L>
1010where
1011 C: Clone,
1012 L: VolatileLink<SynodMsg<C>>,
1013{
1014 // ---------------------------------------------------------------- the acceptor, Figure 4
1015
1016 /// `case ⟨p1a, λ, b⟩ : if b > ballot_num then ballot_num := b; end if;`
1017 /// `send(λ, ⟨p1b, self(), ballot_num, accepted⟩);`
1018 ///
1019 /// The reply carries `ballot_num` **after** any adoption, so a scout whose ballot was refused is
1020 /// told which ballot beat it rather than merely that one did.
1021 fn on_p1a(&mut self, from: NodeId, ballot: Ballot, cx: &mut ProtoCx<'_, Self>) {
1022 if Some(ballot) > self.ballot_num {
1023 self.ballot_num = Some(ballot);
1024 }
1025 // Always `Some` here: either it just adopted, or it already held something at least as
1026 // high, and a real ballot is above `⊥`.
1027 let held = self.ballot_num.expect("an acceptor answering p1a has adopted something");
1028 // §4.1: "return only these pvalues in a p1b message to the scout". One per slot, so this
1029 // message grows with the slots this acceptor has accepted for and not with the ballots the
1030 // run has seen.
1031 let accepted = self
1032 .accepted
1033 .iter()
1034 .map(|(slot, (ballot, command))| Pvalue {
1035 ballot: *ballot,
1036 slot: *slot,
1037 command: command.clone(),
1038 })
1039 .collect();
1040 // §4.2: "This slot number must be included in p1b messages so that leaders can skip the
1041 // lower numbered slots." Without it a leader reads absence as "nothing was accepted", which
1042 // after collection is one of two possible meanings and the wrong one.
1043 let collected = self.collected;
1044 self.transmit(from, SynodMsg::P1b { ballot: held, accepted, collected }, cx);
1045 }
1046
1047 /// `case ⟨p2a, λ, ⟨b, s, c⟩⟩ : if b ≥ ballot_num then ballot_num := b; accepted := accepted ∪
1048 /// {⟨b, s, c⟩}; end if; send(λ, ⟨p2b, self(), ballot_num⟩);`
1049 ///
1050 /// **Departed** from Figure 4's `b = ballot_num`, which does not adopt. The module documentation
1051 /// gives the failure this repairs, the edition the condition comes from, and why A1 and A2
1052 /// survive it.
1053 fn on_p2a(&mut self, from: NodeId, pvalue: Pvalue<C>, cx: &mut ProtoCx<'_, Self>) {
1054 let Pvalue { ballot, slot, command } = pvalue;
1055 // `synod-accept-below-promise` is a mutation the safety-evidence guard compiles: an
1056 // acceptor that accepts under a ballot it has already superseded, and answers as though it
1057 // had not. It breaks A2 and nothing else, so a suite that stays green under it is not
1058 // reading agreement. See `scripts/check-safety-tests.sh`.
1059 let sabotaged = cfg!(feature = "synod-accept-below-promise");
1060 let admissible = sabotaged || Some(ballot) >= self.ballot_num;
1061 if admissible {
1062 // Still monotonic even under the mutation: A1 is a separate claim and stays true, so
1063 // exactly one invariant is removed at a time.
1064 self.ballot_num = Some(self.ballot_num.map_or(ballot, |held| held.max(ballot)));
1065 // §4.1: "acceptors only maintain the most recently accepted pvalue for each slot".
1066 // The latest acceptance, written over whatever was there — and no comparison, because
1067 // the promise has already made the latest the highest. Any stored pvalue's ballot
1068 // became `ballot_num` when it was stored and `ballot_num` never falls, so
1069 // `stored ≤ ballot_num ≤ b` for every `b` this arm admits. The one case where `stored`
1070 // is not strictly below `b` is `stored = ballot_num = b`, and A4 makes that the same
1071 // command. A guard here would be a branch nothing can take; the comparison that does
1072 // the work is the leader's, in `keep_max`.
1073 self.accepted.insert(slot, (ballot, command));
1074 }
1075 let held = if sabotaged {
1076 ballot
1077 } else {
1078 self.ballot_num.expect("an acceptor answering p2a has adopted something")
1079 };
1080 self.transmit(from, SynodMsg::P2b { ballot: held, slot }, cx);
1081 }
1082
1083 // ---------------------------------------------------------------- the scout, Figure 6(b)
1084
1085 /// `∀α ∈ acceptors : send(α, ⟨p1a, self(), b⟩)`, and the state the figure's `var` line declares.
1086 ///
1087 /// Called for a new ballot and again to restart phase one under the same one. Restarting
1088 /// discards a partial `waitfor` and `pvalues` deliberately: an acceptor that already answered
1089 /// answers again, and re-collecting from scratch is what makes the restart idempotent.
1090 fn start_scout(&mut self, ballot: Ballot, cx: &mut ProtoCx<'_, Self>) {
1091 self.scout = Some(Scout {
1092 ballot,
1093 waitfor: self.acceptors.clone(),
1094 pvalues: BTreeMap::new(),
1095 collected: 0,
1096 started: cx.now(),
1097 last_sent: cx.now(),
1098 });
1099 for a in self.acceptors.clone() {
1100 self.transmit(a, SynodMsg::P1a { ballot }, cx);
1101 }
1102 }
1103
1104 /// `case ⟨p1b, α, b', r⟩` — the figure's `then` branch collects, and its `else` reports a
1105 /// preemption to the leader. Both are here, because the leader is this same object.
1106 fn on_p1b(
1107 &mut self,
1108 from: NodeId,
1109 ballot: Ballot,
1110 accepted: Vec<Pvalue<C>>,
1111 collected: Slot,
1112 cx: &mut ProtoCx<'_, Self>,
1113 ) {
1114 let Some(scout) = self.scout.as_mut() else { return };
1115 if ballot < scout.ballot {
1116 // A late answer to a scout that has already exited. Figure 6(b) never sees one: the
1117 // reply is addressed to a thread that no longer exists. See the module's departure on
1118 // what a scout being a field costs.
1119 return;
1120 }
1121 if ballot != scout.ballot {
1122 // `else send(λ, ⟨preempted, b'⟩); exit();`
1123 self.scout = None;
1124 self.preempted(ballot, cx);
1125 return;
1126 }
1127 // `pvalues := pvalues ∪ r; waitfor := waitfor − {α};`
1128 //
1129 // The union reduces per slot as it collects, for the same reason §4.1 gives the acceptor: a
1130 // majority's answers still hold one ballot each for a slot, and keeping them all would put
1131 // the growth §4.1 took off the acceptor straight back onto the leader.
1132 for Pvalue { ballot, slot, command } in accepted {
1133 keep_max(&mut scout.pvalues, ballot, slot, command);
1134 }
1135 scout.collected = scout.collected.max(collected);
1136 scout.waitfor.remove(&from);
1137 if is_majority(scout.waitfor.len(), self.acceptors.len()) {
1138 // `send(λ, ⟨adopted, b, pvalues⟩); exit();`
1139 let scout = self.scout.take().expect("borrowed above");
1140 self.adopted(scout.ballot, scout.pvalues, scout.collected, cx);
1141 }
1142 }
1143
1144 // ---------------------------------------------------------------- the commander, Figure 6(a)
1145
1146 /// `∀α ∈ acceptors : send(α, ⟨p2a, self(), ⟨b, s, c⟩⟩)`, and the figure's `var` line.
1147 fn start_commander(&mut self, slot: Slot, command: C, cx: &mut ProtoCx<'_, Self>) {
1148 let ballot = self.leader_ballot;
1149 // Invariant C1: at most one commander for ⟨b, s⟩. The map key enforces it within a ballot,
1150 // and the ballot changing clears the map.
1151 self.commanders.insert(
1152 slot,
1153 Commander {
1154 ballot,
1155 command: command.clone(),
1156 waitfor: self.acceptors.clone(),
1157 started: cx.now(),
1158 last_sent: cx.now(),
1159 },
1160 );
1161 for a in self.acceptors.clone() {
1162 let pvalue = Pvalue { ballot, slot, command: command.clone() };
1163 self.transmit(a, SynodMsg::P2a { pvalue }, cx);
1164 }
1165 }
1166
1167 /// `case ⟨p2b, α, b'⟩` — collect toward a majority, or report the preemption.
1168 ///
1169 /// The slot comes from the message rather than from the thread's identity; see the module
1170 /// documentation.
1171 fn on_p2b(&mut self, from: NodeId, ballot: Ballot, slot: Slot, cx: &mut ProtoCx<'_, Self>) {
1172 let Some(commander) = self.commanders.get_mut(&slot) else { return };
1173 if ballot < commander.ballot {
1174 // A late answer to a commander that has already exited — see `on_p1b`.
1175 return;
1176 }
1177 if ballot != commander.ballot {
1178 // `else send(λ, ⟨preempted, b'⟩); exit();`
1179 self.commanders.remove(&slot);
1180 self.preempted(ballot, cx);
1181 return;
1182 }
1183 commander.waitfor.remove(&from);
1184 if is_majority(commander.waitfor.len(), self.acceptors.len()) {
1185 // `∀ρ ∈ replicas : send(ρ, ⟨decision, s, c⟩); exit();`
1186 let commander = self.commanders.remove(&slot).expect("borrowed above");
1187 self.decide(slot, commander.command, cx);
1188 }
1189 }
1190
1191 /// `∀ρ ∈ replicas : send(ρ, ⟨decision, s, c⟩)`, and the same thing raised here.
1192 ///
1193 /// Every process learns what was chosen, not only the one whose commander counted the
1194 /// majority — which is what a log above this layer needs and what Figure 6(a) says. The
1195 /// addressees are the processes of the run: §4.4 colocates a replica with every acceptor, so
1196 /// the two sets are one here.
1197 fn decide(&mut self, slot: Slot, command: C, cx: &mut ProtoCx<'_, Self>) {
1198 self.decided.insert(slot, command.clone());
1199 for peer in self.acceptors.clone() {
1200 if peer != self.me {
1201 let command = command.clone();
1202 self.transmit(peer, SynodMsg::Decision { slot, command }, cx);
1203 }
1204 }
1205 cx.indicate(Ind::Decision { slot, command });
1206 }
1207
1208 /// `case ⟨decision, s, c⟩` at a process that did not decide it — the receiving half of the line
1209 /// above, and the answer to a re-proposal for a slot already decided.
1210 ///
1211 /// Raised again if it arrives again: a later ballot can re-command a decided slot, and an
1212 /// answer repeats one deliberately. Invariant A5 makes the command the same — the test is
1213 /// `a_slot_decided_twice_is_announced_twice_and_names_one_command` — so **the layer above must
1214 /// be idempotent per slot**, and it is `multi_paxos_replica`'s `decisions` that makes it so.
1215 ///
1216 /// Suppressing the repeat here instead is not free, and the cost is worse than the duplicate.
1217 /// It needs a set of already-announced slots at every *receiving* process, where `decided` is
1218 /// kept only by leaders, and it grows the same unbounded way. Worse, it would make the recovery
1219 /// above depend on this layer's memory: a replica re-proposes for a slot it has not applied,
1220 /// and the leader's answer is a repeat by construction, so a receiver that dropped repeats
1221 /// would drop the very message that unwedges it. This layer's set is volatile and the replica's
1222 /// is what the log is built from; only one of them can decide what "already delivered" means,
1223 /// and it is not this one.
1224 /// **A commander still running for the slot is left running**, deliberately. Cancelling it is
1225 /// safe — Invariant A5 makes its command the decided one — and saves the `p2a` it resends until
1226 /// its ballot ends. It was written that way first, and `scripts/check-safety-tests.sh` caught
1227 /// what it costs: two of the six tests registered against `synod-ignore-pmax` stopped going
1228 /// red, because a second decision under a later ballot is precisely what the cancellation
1229 /// suppresses. Under the correct clause there is no second decision to suppress; under the
1230 /// mutation there is, and it is the evidence. Silencing a contradiction is not the same as not
1231 /// having one.
1232 fn on_decision(&mut self, slot: Slot, command: C, cx: &mut ProtoCx<'_, Self>) {
1233 // A union, as the page has it: R1 makes a second arrival the same command, so the first
1234 // one stays.
1235 self.decided.entry(slot).or_insert_with(|| command.clone());
1236 cx.indicate(Ind::Decision { slot, command });
1237 }
1238
1239 // ---------------------------------------------------------------- the leader, Figure 7
1240
1241 /// `case ⟨adopted, ballot_num, pvals⟩`.
1242 ///
1243 /// `proposals := proposals ◁ pmax(pvals)` — the update operator replaces this leader's proposal
1244 /// for a slot with the highest-ballot pvalue any acceptor in the majority reported, and keeps
1245 /// this leader's own proposal for a slot nobody reported. **This is the step the whole safety
1246 /// argument rests on**: a proposal already accepted by a majority under a lower ballot is what
1247 /// the new leader proposes, whatever it set out to propose.
1248 fn adopted(
1249 &mut self,
1250 ballot: Ballot,
1251 pvalues: Pvalues<C>,
1252 collected: Slot,
1253 cx: &mut ProtoCx<'_, Self>,
1254 ) {
1255 // "If an `adopted` message arrives for an old ballot number, it is ignored."
1256 if ballot != self.leader_ballot {
1257 return;
1258 }
1259 // `synod-ignore-pmax` is the guard's other mutation: a leader that proposes what it set
1260 // out to propose rather than what the majority already accepted. It removes the one step
1261 // the safety argument rests on and leaves everything else working.
1262 if !cfg!(feature = "synod-ignore-pmax") {
1263 for (slot, command) in pmax(&pvalues) {
1264 self.proposals.insert(slot, command);
1265 }
1266 }
1267 self.active = true;
1268 // §4.2: "This slot number must be included in `p1b` messages so that leaders can skip the
1269 // lower numbered slots." **This is the clause the collection's safety rests on.** An
1270 // acceptor that reports nothing for a slot below `collected` has not told us the slot is
1271 // free — it has told us it no longer remembers, and `f + 1` replicas do. Proposing there
1272 // would put a second command up for a slot already decided.
1273 //
1274 // `synod-skip-collected` is the mutation the safety guard compiles against this line.
1275 if !cfg!(feature = "synod-skip-collected") {
1276 self.proposals.retain(|slot, _| *slot >= collected);
1277 self.collected = self.collected.max(collected);
1278 }
1279 // `∀⟨s, c⟩ ∈ proposals : spawn(Commander(…))`
1280 for (slot, command) in self.proposals.clone() {
1281 self.start_commander(slot, command, cx);
1282 }
1283 }
1284
1285 /// `case ⟨preempted, ⟨r', λ'⟩⟩ : if ⟨r', λ'⟩ > ballot_num then active := false; ballot_num :=
1286 /// (r' + 1, self()); spawn(Scout(…)); end if;`
1287 ///
1288 /// **Departed** in one clause: the scout starts only if Ω still trusts this process. Figure 7
1289 /// starts one unconditionally, which is the duel §3 opens by describing.
1290 fn preempted(&mut self, beat_by: Ballot, cx: &mut ProtoCx<'_, Self>) {
1291 if beat_by <= self.leader_ballot {
1292 return;
1293 }
1294 self.active = false;
1295 self.leader_ballot = Ballot::above(beat_by, self.me);
1296 // The ballot changed, so every commander was running under the old one. C1 keys on the
1297 // ballot; the map keys on the slot, and this is what keeps the two in step.
1298 self.commanders.clear();
1299 self.scout = None;
1300 if self.is_trusted() {
1301 let ballot = self.leader_ballot;
1302 self.start_scout(ballot, cx);
1303 } else {
1304 // Nothing whatever reaches the trace from here: a leader that stops competing sends no
1305 // message, sets no timer and raises no indication. That is the shape of silence this
1306 // repository narrates.
1307 cx.note(Note::LeadershipYielded { to: beat_by.leader, round: beat_by.round });
1308 }
1309 }
1310
1311 /// `upon event ⟨ Ω, Trust | p ⟩` — the departure §3 is replaced by.
1312 ///
1313 /// Figure 7 spawns a scout for its initial ballot at startup, unconditionally. Here a process
1314 /// scouts when trusted and stays passive otherwise, so a process Ω does not trust starts no
1315 /// ballot and the duel does not begin.
1316 fn on_trust(&mut self, leader: NodeId, cx: &mut ProtoCx<'_, Self>) {
1317 self.trusted = Some(leader);
1318 if self.is_trusted() && self.scout.is_none() && !self.active {
1319 let ballot = self.leader_ballot;
1320 self.start_scout(ballot, cx);
1321 }
1322 }
1323
1324 /// `case ⟨propose, s, c⟩ : if ∄c' : ⟨s, c'⟩ ∈ proposals then …`
1325 fn on_propose(
1326 &mut self,
1327 from: Option<NodeId>,
1328 slot: Slot,
1329 command: C,
1330 cx: &mut ProtoCx<'_, Self>,
1331 ) {
1332 // The answer, and the reason the asker's re-proposal is a recovery rather than a loop: a
1333 // decision is announced once and its commander then exits, so a process the announcement
1334 // never reached has no other way back. To the asker alone, where the announcement was a
1335 // fan-out. Liu et al.: a leader "can then work on deciding for that slot if a decision for
1336 // it has not been made; otherwise, it can send back the decision for that slot".
1337 //
1338 // From `decided`, never from `proposals`: the module documentation has the schedule in
1339 // which the two differ. Any process that knows the decision may answer.
1340 if let Some(asker) = from
1341 && let Some(decided) = self.decided.get(&slot).cloned()
1342 {
1343 self.transmit(asker, SynodMsg::Decision { slot, command: decided }, cx);
1344 return;
1345 }
1346 // §4.2's skip, and it has to be here as well as in `adopted`. Filtering at adoption covers
1347 // the proposals a leader *already held*; it does nothing about one arriving afterwards, and
1348 // an active leader would command it. `agreement_holds_across_a_collection` split a slot on
1349 // exactly that: everything for slots 1–4 was collected everywhere, the successor adopted
1350 // cleanly, and was then asked for a different command for slot 1 — and commanded it.
1351 //
1352 // A slot below the watermark is decided and `f + 1` replicas hold the decision. There is
1353 // nothing to propose for and nothing this layer could safely put there.
1354 if !cfg!(feature = "synod-skip-collected") && slot < self.collected {
1355 cx.note(Note::ProposalIgnored { slot });
1356 // And say so, rather than dropping it silently. The asker is a replica that missed a
1357 // decision this layer has collected, and the only thing that can help it now is a peer.
1358 cx.indicate(Ind::Collected { slot });
1359 return;
1360 }
1361 // §4.4: "If λ is active, it will start a commander." An adopted ballot stands until
1362 // something preempts it, so an active leader commands even where Ω has moved on — standing
1363 // down means starting no new ballots, not abandoning one a majority already adopted.
1364 if self.active {
1365 if self.proposals.contains_key(&slot) {
1366 // No effect at all: dropped because this leader already has one for the slot,
1367 // which is what enforces C1 against a second commander. The attempt already in
1368 // flight is what fills it.
1369 cx.note(Note::ProposalIgnored { slot });
1370 return;
1371 }
1372 self.proposals.insert(slot, command.clone());
1373 self.start_commander(slot, command, cx);
1374 return;
1375 }
1376 if self.is_trusted() {
1377 // Trusted but not yet adopted: remembered, and `adopted` commands it. Figure 7's
1378 // `if active then` arm doing nothing. Remembered once; a repeat changes nothing.
1379 if self.proposals.contains_key(&slot) {
1380 cx.note(Note::ProposalIgnored { slot });
1381 return;
1382 }
1383 self.proposals.insert(slot, command);
1384 return;
1385 }
1386 // §4.4: "If λ is passive, monitoring another leader λ′, it forwards the proposal to λ′."
1387 // Remembering here would be the same as dropping — this process will not lead — and so is
1388 // anything it remembered *before* it stopped leading: the `∄c'` guard is C1's, and C1 is
1389 // about commanders, of which a passive process has none. What it forgets here comes back
1390 // through `pmax` if a majority accepted it and this process ever leads again.
1391 //
1392 // **A `Propose` that arrived is never forwarded again.** Two processes whose detectors
1393 // disagree would pass one back and forth for as long as they disagree, so the path is at
1394 // most two hops by construction and the asker's own timeout is what recovers a drop.
1395 match (from, self.trusted) {
1396 (None, Some(leader)) if leader != self.me => {
1397 self.proposals.remove(&slot);
1398 self.transmit(leader, SynodMsg::Propose { slot, command }, cx);
1399 }
1400 _ => cx.note(Note::ProposalIgnored { slot }),
1401 }
1402 }
1403
1404 // ---------------------------------------------------------------- retries and escalation
1405
1406 /// How long an attempt may go unanswered before its requests are sent again.
1407 ///
1408 /// **Not the sweep interval, and the two must not be confused.** The sweep decides how often
1409 /// this question is *asked*; this decides the answer. Conflating them is what made a resend
1410 /// happen every `Timing::retransmit`, which every suite configures below the delivery bound —
1411 /// so a request went out again before an answer to it could possibly have arrived, and two
1412 /// thirds of phase two was the protocol talking over itself.
1413 ///
1414 /// It must **exceed one round trip**, or it resends what is merely in flight, and it must stay
1415 /// **below `escalate_after`**, or an attempt is escalated before a retransmission has been
1416 /// tried and a phase restarts for a message that was never lost. Half of `escalate_after`
1417 /// satisfies both by construction and gives exactly one resend before escalating, which is the
1418 /// shape worth having: ask once more, then conclude something is wrong.
1419 ///
1420 /// The lower bound is a constraint on the *configuration* rather than on this expression:
1421 /// `escalate_after` must itself be more than twice the delivery bound, which is weaker than
1422 /// what the detector already needs of `detect_after` and which
1423 /// `the_resend_threshold_sits_between_a_round_trip_and_the_escalation` pins.
1424 fn resend_after(&self) -> Duration {
1425 self.escalate_after / 2
1426 }
1427
1428 /// The periodic sweep: resend what is outstanding, and escalate an attempt that has waited too
1429 /// long. See the module's table of the three liveness fixes.
1430 fn sweep(&mut self, cx: &mut ProtoCx<'_, Self>) {
1431 let now = cx.now();
1432 let resend_after = self.resend_after();
1433 if let Some(scout) = self.scout.as_ref() {
1434 let ballot = scout.ballot;
1435 if now - scout.started >= self.escalate_after {
1436 // Phase 1, lost `p1a`: neither adopted nor preempted. Restart, same ballot.
1437 self.start_scout(ballot, cx);
1438 } else if now - scout.last_sent >= self.resend_after() {
1439 // Long enough that an answer would have arrived. Anything sooner would be
1440 // resending what is still in flight.
1441 let waiting = scout.waitfor.clone();
1442 if let Some(scout) = self.scout.as_mut() {
1443 scout.last_sent = now;
1444 }
1445 for a in waiting {
1446 self.transmit(a, SynodMsg::P1a { ballot }, cx);
1447 }
1448 }
1449 return;
1450 }
1451 if !self.active {
1452 return;
1453 }
1454 let stalled = self
1455 .commanders
1456 .values()
1457 .any(|commander| now - commander.started >= self.escalate_after);
1458 if stalled {
1459 // Phase 2, lost `preempt`: a majority may hold a higher ballot, in which case no `p2a`
1460 // can ever reach one and resending is a loop. Go back to phase one, same ballot — the
1461 // `p1b` answers are how the higher ballot is learnt.
1462 self.active = false;
1463 self.commanders.clear();
1464 let ballot = self.leader_ballot;
1465 self.start_scout(ballot, cx);
1466 return;
1467 }
1468 // Phase 2, lost `p2b`: resend the pvalue to whoever has not answered — but only for a
1469 // commander that has waited longer than an answer could take.
1470 let outstanding: Vec<(Slot, Ballot, C, Vec<NodeId>)> = self
1471 .commanders
1472 .iter_mut()
1473 .filter(|(_, c)| now - c.last_sent >= resend_after)
1474 .map(|(slot, c)| {
1475 c.last_sent = now;
1476 (*slot, c.ballot, c.command.clone(), c.waitfor.iter().copied().collect())
1477 })
1478 .collect();
1479 for (slot, ballot, command, waitfor) in outstanding {
1480 for a in waitfor {
1481 let pvalue = Pvalue { ballot, slot, command: command.clone() };
1482 self.transmit(a, SynodMsg::P2a { pvalue }, cx);
1483 }
1484 }
1485 }
1486
1487 /// Send again, to one peer, what that peer has not answered — because a session with it has
1488 /// just been established.
1489 ///
1490 /// **The event this repository documents and nothing acted on.** `session_link.rs` calls an
1491 /// establishment "the moment on which anything that must be resent can be", and until this
1492 /// every module over a session link propagated it and did nothing else. A session ending is the
1493 /// only way this stack loses a message, so the establishment that follows is the only moment a
1494 /// resend can succeed; waiting out the threshold instead makes recovery slower than the
1495 /// information already available, and it is why that threshold can afford to be generous.
1496 ///
1497 /// To the peer the event names, and to nobody else. A fan-out on every establishment would cost
1498 /// membership squared as a cluster reconnects, for a peer that is owed nothing.
1499 ///
1500 /// `last_sent` is deliberately **not** reset. It belongs to the attempt rather than to a peer,
1501 /// so moving it here would delay the sweep's resends to peers whose sessions never broke. The
1502 /// cost is that one peer may be asked twice in quick succession after an ending, which is
1503 /// bounded by how often sessions end and is the cheaper mistake.
1504 fn resend_to(&mut self, peer: NodeId, cx: &mut ProtoCx<'_, Self>) {
1505 if let Some(scout) = self.scout.as_ref() {
1506 if scout.waitfor.contains(&peer) {
1507 let ballot = scout.ballot;
1508 self.transmit(peer, SynodMsg::P1a { ballot }, cx);
1509 }
1510 // A leader in phase one has no commanders: `preempted` clears them and `adopted` is
1511 // what starts them. Nothing below applies.
1512 return;
1513 }
1514 if !self.active {
1515 return;
1516 }
1517 let owed: Vec<Pvalue<C>> = self
1518 .commanders
1519 .iter()
1520 .filter(|(_, c)| c.waitfor.contains(&peer))
1521 .map(|(slot, c)| Pvalue { ballot: c.ballot, slot: *slot, command: c.command.clone() })
1522 .collect();
1523 for pvalue in owed {
1524 self.transmit(peer, SynodMsg::P2a { pvalue }, cx);
1525 }
1526 }
1527
1528 /// `case ⟨applied, ρ, s⟩` — §4.2's periodic update, from a replica's machine.
1529 ///
1530 /// Recording it is all that happens here; the collection is [`MultiPaxosSynod::collect`], run
1531 /// once afterwards because a new report is the only thing that can move the watermark.
1532 fn on_applied(&mut self, from: NodeId, slot_out: Slot, cx: &mut ProtoCx<'_, Self>) {
1533 let held = self.reported.entry(from).or_insert(slot_out);
1534 *held = (*held).max(slot_out);
1535 self.collect(cx);
1536 }
1537
1538 /// `case ⟨catchup, ρ, s⟩` — a peer has fallen behind. This layer holds nothing to answer with,
1539 /// so it asks the layer above, which does.
1540 fn on_catch_up(&mut self, from: NodeId, from_slot: Slot, cx: &mut ProtoCx<'_, Self>) {
1541 cx.indicate(Ind::CatchUpWanted { peer: from, from_slot });
1542 }
1543
1544 /// The highest slot at least `f + 1` members have applied every decision up to.
1545 ///
1546 /// §4.2: "Once a leader or acceptor learns that at least `f + 1` replicas have received all
1547 /// decisions up to some slot `s`, all information about lower numbered slots can be garbage
1548 /// collected."
1549 ///
1550 /// `f + 1` of `2 f + 1` is a majority, and it is the *same* majority the rest of this module
1551 /// counts, because §4.4 puts a replica on every acceptor's machine — the replica set and the
1552 /// acceptor set are one here, so there is one membership and one quorum size.
1553 ///
1554 /// **With fewer than `f + 1` members reporting, this stays where it is and nothing is
1555 /// collected.** The source says so directly: "if there are fewer than `2 f + 1` replicas, the
1556 /// crash of `f` replicas would leave fewer than `f + 1` replicas to send periodic updates and
1557 /// no garbage collection could be done". A run that stops collecting after `f` crashes is
1558 /// behaving as specified, not failing.
1559 fn watermark(&self) -> Slot {
1560 let quorum = self.acceptors.len() / 2 + 1;
1561 let mut applied: Vec<Slot> = self.reported.values().copied().collect();
1562 applied.sort_unstable_by(|a, b| b.cmp(a));
1563 applied.get(quorum - 1).copied().unwrap_or(0)
1564 }
1565
1566 /// Discard everything below the watermark — §4.2, and the whole of what makes this module an
1567 /// implementation rather than a transcription.
1568 ///
1569 /// The three things the audit in `docs/bounded-space.md` lists: the acceptor's pvalues, the
1570 /// leader's proposals, and the record of which slots are decided. Discarding them destroys no
1571 /// information, it *moves* it — "replicas can learn decisions, and the application state that
1572 /// results from those decisions, from one another" — and what stops a later leader reading the
1573 /// gap as emptiness is `collected` travelling in `p1b`.
1574 fn collect(&mut self, cx: &mut ProtoCx<'_, Self>) {
1575 let wm = self.watermark();
1576 if wm <= self.collected {
1577 return;
1578 }
1579 self.collected = wm;
1580 self.accepted.retain(|slot, _| *slot >= wm);
1581 self.proposals.retain(|slot, _| *slot >= wm);
1582 self.decided.retain(|slot, _| *slot >= wm);
1583 cx.note(Note::CollectedBelow { slot: wm });
1584 }
1585
1586 /// Arm the sweep, or re-arm it after it fired.
1587 fn arm(&mut self, cx: &mut ProtoCx<'_, Self>) {
1588 self.tick = Some(cx.set_timer(self.retry));
1589 }
1590
1591 // ---------------------------------------------------------------- composition
1592
1593 /// Put one message on the wire, through the link beneath.
1594 fn transmit(&mut self, to: NodeId, msg: SynodMsg<C>, cx: &mut ProtoCx<'_, Self>) {
1595 self.through_link(cx, |l, ccx| l.on_cmd(L::send(to, msg), ccx));
1596 }
1597
1598 fn through_omega(
1599 &mut self,
1600 cx: &mut ProtoCx<'_, Self>,
1601 f: impl FnOnce(&mut EventualLeaderDetector, &mut ProtoCx<'_, EventualLeaderDetector>),
1602 ) {
1603 let mut inds = self.omega.run(cx, Wire::Detector, f);
1604 for eld::Ind::Trust { leader } in inds.drain(..) {
1605 self.on_trust(leader, cx);
1606 }
1607 self.omega.reclaim(inds);
1608 }
1609
1610 fn through_link(
1611 &mut self,
1612 cx: &mut ProtoCx<'_, Self>,
1613 f: impl FnOnce(&mut L, &mut ProtoCx<'_, L>),
1614 ) {
1615 let mut inds = self.link.run(cx, Wire::Synod, f);
1616 for ind in inds.drain(..) {
1617 match L::classify(ind) {
1618 LinkInd::Deliver { from, msg } => self.on_synod_msg(from, msg, cx),
1619 // Propagated, never absorbed: nothing here bridges a session ending.
1620 LinkInd::Boundary(Boundary::Ended { peer, epoch }) => {
1621 cx.indicate(Ind::SessionEnded { peer, epoch })
1622 }
1623 LinkInd::Boundary(Boundary::Established { peer, epoch }) => {
1624 cx.indicate(Ind::SessionEstablished { peer, epoch });
1625 self.resend_to(peer, cx);
1626 }
1627 }
1628 }
1629 self.link.reclaim(inds);
1630 }
1631
1632 /// The four `switch receive()` arms of Figures 4 and 6, dispatched.
1633 fn on_synod_msg(&mut self, from: NodeId, msg: SynodMsg<C>, cx: &mut ProtoCx<'_, Self>) {
1634 match msg {
1635 SynodMsg::P1a { ballot } => self.on_p1a(from, ballot, cx),
1636 SynodMsg::P1b { ballot, accepted, collected } => {
1637 self.on_p1b(from, ballot, accepted, collected, cx)
1638 }
1639 SynodMsg::Applied { slot_out } => self.on_applied(from, slot_out, cx),
1640 SynodMsg::CatchUp { from_slot } => self.on_catch_up(from, from_slot, cx),
1641 SynodMsg::P2a { pvalue } => self.on_p2a(from, pvalue, cx),
1642 SynodMsg::P2b { ballot, slot } => self.on_p2b(from, ballot, slot, cx),
1643 SynodMsg::Decision { slot, command } => self.on_decision(slot, command, cx),
1644 SynodMsg::Propose { slot, command } => self.on_propose(Some(from), slot, command, cx),
1645 }
1646 }
1647}
1648
1649impl<C, L> Protocol for MultiPaxosSynod<C, L>
1650where
1651 C: Clone,
1652 L: VolatileLink<SynodMsg<C>>,
1653{
1654 type Cmd = Cmd<C>;
1655 type Ind = Ind<C>;
1656 type Msg = Wire<L::Msg>;
1657 /// Whatever the link's guarantees are conditional on. This layer adds no condition of its own
1658 /// and bridges none of the link's.
1659 type Scope = L::Scope;
1660 type Note = crate::Note;
1661 /// Keeps nothing durably, which is what makes this crash-stop. §4.3 is the change that alters
1662 /// it, and the module documentation says what a durable ballot counter would buy.
1663 type Meta = core::convert::Infallible;
1664 type Entry = core::convert::Infallible;
1665
1666 /// A request from the layer above carries no sender, which is what distinguishes it from a
1667 /// forwarded one: only a proposal that has *not* travelled may be forwarded.
1668 fn on_cmd(&mut self, cmd: Cmd<C>, cx: &mut ProtoCx<'_, Self>) {
1669 match cmd {
1670 Cmd::Propose { slot, command } => self.on_propose(None, slot, command, cx),
1671 // §4.2's periodic update. Recorded here as this machine's own — the replica is on it —
1672 // and passed to everybody else's leader and acceptor, which is what colocation makes
1673 // this layer's job rather than the replica's.
1674 // Ask **one** peer, the one furthest ahead of us. A fan-out would have every member
1675 // answer the same request, which costs membership times the gap for information one
1676 // process could have supplied.
1677 Cmd::CatchUp { from_slot } => {
1678 let ahead = self
1679 .reported
1680 .iter()
1681 .filter(|(peer, applied)| **peer != self.me && **applied > from_slot)
1682 .max_by_key(|(_, applied)| **applied)
1683 .map(|(peer, _)| *peer);
1684 match ahead {
1685 Some(peer) => self.transmit(peer, SynodMsg::CatchUp { from_slot }, cx),
1686 // Nobody has said they are ahead. Nothing to do, and nothing lost: the layer
1687 // above asks again, and a report will arrive.
1688 None => cx.note(Note::ProposalIgnored { slot: from_slot }),
1689 }
1690 }
1691 Cmd::Teach { to, slot, command } => {
1692 self.transmit(to, SynodMsg::Decision { slot, command }, cx);
1693 }
1694 Cmd::Applied { slot_out } => {
1695 self.on_applied(self.me, slot_out, cx);
1696 for peer in self.acceptors.clone() {
1697 if peer != self.me {
1698 self.transmit(peer, SynodMsg::Applied { slot_out }, cx);
1699 }
1700 }
1701 }
1702 }
1703 }
1704
1705 fn on_msg(&mut self, from: NodeId, msg: Self::Msg, cx: &mut ProtoCx<'_, Self>) {
1706 match msg {
1707 Wire::Detector(h) => self.through_omega(cx, |o, ccx| o.on_msg(from, h, ccx)),
1708 Wire::Synod(m) => self.through_link(cx, |l, ccx| l.on_msg(from, m, ccx)),
1709 }
1710 }
1711
1712 /// An expiry is offered to every layer, so both children are given it and this layer acts only
1713 /// on the handle it registered.
1714 fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
1715 self.through_omega(cx, |o, ccx| o.on_timer(id, ccx));
1716 self.through_link(cx, |l, ccx| l.on_timer(id, ccx));
1717 if self.tick != Some(id) {
1718 return;
1719 }
1720 self.arm(cx);
1721 self.sweep(cx);
1722 }
1723
1724 /// `⟨ Init ⟩` — start the detector, whose first `Trust` may immediately make this process scout,
1725 /// and arm the sweep.
1726 fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
1727 self.arm(cx);
1728 self.through_omega(cx, |o, ccx| o.on_init(ccx));
1729 }
1730
1731 /// Hand the boundary down to the link, which is the layer that knows what it means. Leaving
1732 /// this to the trait's default would take a scope event the driver raised and drop it — the
1733 /// cardinal sin of `docs/conditional-guarantees.md`.
1734 fn on_scope_event(&mut self, scope: L::Scope, cx: &mut ProtoCx<'_, Self>) {
1735 self.through_link(cx, |l, ccx| l.on_scope_event(scope, ccx));
1736 }
1737}