Skip to main content

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}