Skip to main content

recon_protocols/
logged_epoch_consensus.rs

1//! Read/write epoch consensus that survives a restart.
2//!
3//! **Status: implementation. Space: bounded by membership, plus what the stubborn children hold
4//! outstanding — which nothing here retires. See the departure on `Stop`.**
5//!
6//! Cachin, Guerraoui & Rodrigues, Module 5.7 (`LoggedEpochConsensus`) and Algorithm 5.9 ("Logged
7//! Read/Write Epoch Consensus"), quoted from the book:
8//!
9//! ```text
10//! Algorithm 5.9: Logged Read/Write Epoch Consensus
11//! Implements: EpochConsensus, instance lep, with timestamp ets and leader ℓ.
12//! Uses:
13//!     StubbornPointToPointLinks, instance sl;
14//!     StubbornBestEffortBroadcast, instance sbeb;
15//!
16//! upon event ⟨ lep, Init | state ⟩ do
17//!     (valts, val) := state;
18//!     store(valts, val);
19//!     tmpval := ⊥;
20//!     states := [⊥]^N;
21//!     accepted := 0;
22//!
23//! upon event ⟨ lep, Recovery ⟩ do
24//!     retrieve(valts, val);
25//!
26//! upon event ⟨ lep, Propose | v ⟩ do                       // only leader ℓ
27//!     tmpval := v;
28//!     trigger ⟨ sbeb, Broadcast | [READ] ⟩;
29//!
30//! upon event ⟨ sbeb, Deliver | ℓ, [READ] ⟩ do
31//!     trigger ⟨ sl, Send | ℓ, [STATE, valts, val] ⟩;
32//!
33//! upon event ⟨ sl, Deliver | q, [STATE, ts, v] ⟩ do        // only leader ℓ
34//!     states[q] := (ts, v);
35//!
36//! upon #(states) > N/2 do                                  // only leader ℓ
37//!     (ts, v) := highest(states);
38//!     if v ≠ ⊥ then tmpval := v;
39//!     states := [⊥]^N;
40//!     trigger ⟨ sbeb, Broadcast | [WRITE, tmpval] ⟩;
41//!
42//! upon event ⟨ sbeb, Deliver | ℓ, [WRITE, v] ⟩ do
43//!     (valts, val) := (ets, v);
44//!     store(valts, val);
45//!     trigger ⟨ sl, Send | ℓ, [ACCEPT] ⟩;
46//!
47//! upon event ⟨ sl, Deliver | q, [ACCEPT] ⟩ do              // only leader ℓ
48//!     accepted := accepted + 1;
49//!
50//! upon accepted > N/2 do                                   // only leader ℓ
51//!     accepted := 0;
52//!     trigger ⟨ sbeb, Broadcast | [DECIDED, tmpval] ⟩;
53//!
54//! upon event ⟨ sbeb, Deliver | ℓ, [DECIDED, v] ⟩ do
55//!     epochdecision := v;
56//!     store(epochdecision);
57//!     trigger ⟨ lep, Decide | epochdecision ⟩;
58//!
59//! upon event ⟨ lep, Abort ⟩ do
60//!     trigger ⟨ lep, Aborted | (valts, val) ⟩;
61//!     halt;                                                // stop operating when aborted
62//! ```
63//!
64//! The safety argument is [`crate::epoch_consensus`]'s and is not restated here: two majorities
65//! intersect, so a later epoch's read reaches a process that accepted whatever an earlier epoch
66//! decided, and `if v ≠ ⊥ then tmpval := v` makes the later leader adopt it. What this module adds
67//! is that the argument still holds when the processes holding that intersection go down and come
68//! back.
69//!
70//! # Two `store` calls, one metadata value
71//!
72//! The book writes `store(valts, val)` and `store(epochdecision)` as separate calls. `Cx::storage`
73//! offers one rewritten metadata value and an appended sequence, so both land in one [`Durable`]
74//! that is rewritten each time. Nothing accumulates: an epoch accepts at most one value and decides
75//! at most one, so the record is a fixed size and `Entry` is uninhabited.
76//!
77//! # Durable before visible, twice, and both in the handler's own text
78//!
79//! `store(valts, val); trigger ⟨ sl, Send | ℓ, [ACCEPT] ⟩` — the acceptance is a **promise to a
80//! quorum**. A process that told the leader it had accepted `v` at `ets`, and then came back with
81//! no record of it, would answer a later epoch's read with an empty state; the later leader would
82//! find nothing in the intersection and be free to write something else, after `v` had already been
83//! decided. That is `EPC4` failing, and it fails silently.
84//!
85//! `store(epochdecision); trigger ⟨ lep, Decide | v ⟩` — same shape one step later. The layer above
86//! reads `epochdecision` back on recovery ([`LoggedEpochConsensus::epoch_decision`]) and that is how
87//! Algorithm 5.10 knows a process had decided before it went down.
88//!
89//! Both orders are written here, in these handlers, and not left to a driver to arrange by
90//! buffering effects until the handler returns. `Cx` supports eager sinks, so buffering is not
91//! something this code may assume.
92//!
93//! # Departure: repeats are idempotent, and the standing conditions fire once
94//!
95//! [`crate::epoch_consensus`] runs over perfect links, which deliver each message once. This one
96//! runs over stubborn ones, which must not deduplicate — repeating for ever is what reaches a
97//! process that was down when the message was sent. So every handler here sees its message many
98//! times, and the book's counters do not survive that:
99//!
100//! - `accepted := accepted + 1` counts *messages*, and one process's ACCEPT arrives for ever. The
101//!   count would pass `N/2` on its own with a single acceptance in the whole run. It is a set of
102//!   the processes that have accepted, so a repeat adds nothing.
103//! - `upon #(states) > N/2` and `upon accepted > N/2` are standing conditions the book re-arms by
104//!   clearing what they count. Clearing is not enough when the messages come back: `written` and
105//!   `announced` make each fire once, as they already do in [`crate::epoch_consensus`].
106//! - `store` is a rewrite of the same value on a repeat, which is idempotent, but it is still a
107//!   write. `WRITE` is applied only when it changes something, so the write count stays one per
108//!   acceptance and a test can check that rather than take it on trust.
109//! - **A follower answers `READ` once and `WRITE` once.** The book answers every delivery. Over a
110//!   stubborn link the answer is itself retransmitted until this instance ends, so a second answer
111//!   to a redelivered `READ` is a second stubborn transmission carrying the same content — and
112//!   since redeliveries never stop, neither would the transmissions. Measured before this guard: the
113//!   send rate grew linearly in time, 12.6k → 76.6k per 400 ms across five windows, with nothing
114//!   faulty. One answer is enough for the same reason retransmission exists at all: a leader that
115//!   crashed and came back re-proposes, and what reaches its new incarnation is the follower's
116//!   *original* reply, still going. A follower that crashes forgets it answered and answers again,
117//!   which is correct — its link forgot the transmission too.
118//!
119//! # Departure: messages carry the epoch they belong to
120//!
121//! As in [`crate::epoch_consensus`]: instances are addressed `lep.ets` in the book and by nothing at
122//! all on a real wire, so a `WRITE` from epoch 7 arriving after epoch 11 began would be accepted and
123//! recorded at timestamp 11 — an acceptance that never happened. The stamp is in [`Tagged`], inside
124//! the instance, because the epoch is the instance's own identity.
125//!
126//! # Departure: nothing calls `Stop`
127//!
128//! As in [`crate::logged_epoch_change`]. The stubborn children retransmit until retired and nothing
129//! retires them, so space grows with the number of distinct messages an epoch sends rather than
130//! with the membership. Bounded in practice by the epoch ending, which is what `Abort` is for.
131//!
132//! ```text
133//! EPC1 [always]  Validity — a decided value was proposed in this epoch, or was the highest-
134//!                timestamped value some process had already accepted
135//! EPC2 [always]  Uniform agreement — no two processes decide differently in one epoch
136//! EPC3 [always]  Integrity — a process decides at most once
137//! EPC4 [always]  Lock-in — a value decided in an earlier epoch is what a later one decides, and
138//!                **this holds across a crash**: what a process accepted is read back on recovery
139//! EPC5 [always]  Abort behaviour — an abandoned instance reports its state and then is silent
140//! ```
141
142use recon_core::{Child, NodeId, ProtoCx, Protocol, TimerId};
143use serde::{Deserialize, Serialize};
144use std::collections::{BTreeMap, BTreeSet};
145
146use crate::stubborn_broadcast::{self as sbeb, BroadcastId, StubbornBroadcast};
147use crate::stubborn_link::{self as sl, SendId, StubbornLink};
148
149/// `(valts, val)` — what a process has accepted, and when.
150///
151/// `val` is `None` for the book's `⊥`: nothing accepted yet, at timestamp zero.
152#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
153pub struct State<V> {
154    pub valts: u64,
155    pub val: Option<V>,
156}
157
158impl<V> Default for State<V> {
159    fn default() -> Self {
160        State { valts: 0, val: None }
161    }
162}
163
164/// Everything this instance keeps durably, as one rewritten value.
165#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
166pub struct Durable<V> {
167    /// `(valts, val)`.
168    pub state: State<V>,
169    /// `epochdecision`, once this epoch has decided.
170    pub decision: Option<V>,
171}
172
173/// What travels by `sbeb` — the leader speaking to everyone.
174#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
175pub enum Announce<V> {
176    /// `[READ]`.
177    Read,
178    /// `[WRITE, v]`.
179    Write { val: V },
180    /// `[DECIDED, v]`.
181    Decided { val: V },
182}
183
184/// What travels by `sl` — a follower answering the leader.
185#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
186pub enum Reply<V> {
187    /// `[STATE, valts, val]`.
188    StateIs { valts: u64, val: Option<V> },
189    /// `[ACCEPT]`.
190    Accept,
191}
192
193/// A message stamped with the epoch it belongs to.
194///
195/// The stamp lives here rather than in the layer above because the epoch is this instance's own
196/// identity — it stamps what it sends and drops what is not addressed to it.
197#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
198pub struct Tagged<M> {
199    pub ets: u64,
200    pub msg: M,
201}
202
203/// The wire, multiplexing the two children the book names.
204#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
205pub enum Wire<V> {
206    /// `sbeb` — the leader's announcements.
207    Announce(Tagged<Announce<V>>),
208    /// `sl` — the followers' replies, each to one process.
209    Reply(Tagged<Reply<V>>),
210}
211
212/// Requests from the layer above.
213#[derive(Debug, Clone, PartialEq, Eq)]
214pub enum Cmd<V> {
215    /// `⟨ lep, Propose | v ⟩`. Acted on only by this epoch's leader.
216    Propose(V),
217    /// `⟨ lep, Abort ⟩`.
218    Abort,
219}
220
221/// Indications to the layer above.
222#[derive(Debug, Clone, PartialEq, Eq)]
223pub enum Ind<V> {
224    /// `⟨ lep, Decide | v ⟩`. Raised only after the decision is durable.
225    Decide(V),
226    /// `⟨ lep, Aborted | (valts, val) ⟩` — the state the replacement instance begins from.
227    Aborted(State<V>),
228}
229
230/// Abortable consensus within one epoch, whose acceptances survive a restart.
231#[derive(Debug)]
232pub struct LoggedEpochConsensus<V: Clone> {
233    me: NodeId,
234    peers: BTreeSet<NodeId>,
235    /// `ets` — this instance's epoch timestamp.
236    ets: u64,
237    /// `ℓ` — this epoch's leader.
238    leader: NodeId,
239    /// `(valts, val)` and `epochdecision` — durable, and mirrored here.
240    durable: Durable<V>,
241    /// `tmpval` — the value the leader is trying to write. Volatile, as in the book.
242    tmpval: Option<V>,
243    /// `states` — what the leader has read back, by process. A map, so a repeat replaces.
244    states: BTreeMap<NodeId, State<V>>,
245    /// `accepted` — **which** processes have acknowledged, not how many messages said so.
246    accepted: BTreeSet<NodeId>,
247    /// Whether the write has been sent, so a repeat cannot resend it.
248    written: bool,
249    /// Whether the decision has been announced, so a repeat cannot re-announce it.
250    announced: bool,
251    /// Whether the decision has been reported upward, so a repeated `DECIDED` decides once.
252    decided: bool,
253    /// Whether this follower has answered the leader's `READ`. One stubborn reply is enough.
254    state_sent: bool,
255    /// Whether this follower has answered the leader's `WRITE`. Likewise.
256    accept_sent: bool,
257    /// Test-only: answer *every* redelivery, which is what this module did before the two flags
258    /// above were added. See [`LoggedEpochConsensus::with_reply_per_redelivery_defect`].
259    reply_per_redelivery: bool,
260    /// `halt`. Every handler returns immediately once this is set.
261    aborted: bool,
262    /// Names the next stubborn transmission. Volatile, and so is what it keys.
263    next_send: u64,
264    /// Names the next stubborn broadcast. Volatile, and so is what it keys.
265    next_broadcast: u64,
266    sbeb: Child<StubbornBroadcast<Tagged<Announce<V>>>>,
267    sl: Child<StubbornLink<Tagged<Reply<V>>>>,
268}
269
270impl<V: Clone> LoggedEpochConsensus<V> {
271    /// `⟨ lep, Init | state ⟩` — an instance for epoch `ets` led by `leader`, beginning from
272    /// `state`.
273    ///
274    /// The book's `store(valts, val)` in `Init` happens on the first event this instance handles,
275    /// because a constructor has no context to write through. [`Protocol::on_init`] is where it
276    /// lands, and it lands before anything is sent.
277    pub fn new(
278        me: NodeId,
279        peers: impl IntoIterator<Item = NodeId>,
280        ets: u64,
281        leader: NodeId,
282        state: State<V>,
283        retransmit: core::time::Duration,
284    ) -> Self {
285        let mut peers: BTreeSet<NodeId> = peers.into_iter().collect();
286        peers.insert(me);
287        LoggedEpochConsensus {
288            me,
289            peers: peers.clone(),
290            ets,
291            leader,
292            durable: Durable { state, decision: None },
293            tmpval: None,
294            states: BTreeMap::new(),
295            accepted: BTreeSet::new(),
296            written: false,
297            announced: false,
298            decided: false,
299            state_sent: false,
300            accept_sent: false,
301            reply_per_redelivery: false,
302            aborted: false,
303            next_send: 0,
304            next_broadcast: 0,
305            sbeb: Child::new(StubbornBroadcast::new(me, peers.clone(), retransmit)),
306            sl: Child::new(StubbornLink::new(retransmit)),
307        }
308    }
309
310    /// `ets`.
311    pub fn timestamp(&self) -> u64 {
312        self.ets
313    }
314
315    /// Whether this instance has been abandoned.
316    pub fn is_aborted(&self) -> bool {
317        self.aborted
318    }
319
320    /// **Put a fixed defect back.** Answer every redelivered `READ` and `WRITE` on a fresh
321    /// stubborn transmission, as this module did before `state_sent` and `accept_sent` existed.
322    ///
323    /// The consequence is not a wrong decision: it is work that grows with how long the run has
324    /// been going rather than with membership — 12.6k, 28.6k, 44.6k, 60.6k, 76.6k sends in
325    /// successive 400 ms windows, because each answer joins a stubborn set that is never emptied.
326    /// `the_send_rate_does_not_grow_after_the_epoch_has_decided` is the test that holds it fixed.
327    ///
328    /// It exists so that the scenario shrinker can be demonstrated against a defect this project
329    /// actually had rather than against a toy; `shrinking_a_real_defect.rs` is the only caller, and
330    /// a test asserts that reintroducing it does break the bound. Nothing else may call it: a
331    /// process built this way violates the module's own stated space bound on purpose.
332    #[doc(hidden)]
333    pub fn with_reply_per_redelivery_defect(mut self) -> Self {
334        self.reply_per_redelivery = true;
335        self
336    }
337
338    /// `(valts, val)` — what this process has accepted.
339    pub fn state(&self) -> &State<V> {
340        &self.durable.state
341    }
342
343    /// `epochdecision` — what this epoch decided, if this process saw it decide.
344    ///
345    /// Read by the layer above after a recovery: `retrieve(epochdecision) of instance lep.ets` is
346    /// how Algorithm 5.10 learns that a process had decided before it went down.
347    pub fn epoch_decision(&self) -> Option<&V> {
348        self.durable.decision.as_ref()
349    }
350
351    /// `N/2` — the threshold both majorities are measured against.
352    fn majority(&self) -> usize {
353        self.peers.len() / 2
354    }
355
356    fn is_leader(&self) -> bool {
357        self.me == self.leader
358    }
359
360    fn broadcast(&mut self, msg: Announce<V>, cx: &mut ProtoCx<'_, Self>) {
361        let msg = Tagged { ets: self.ets, msg };
362        let id = BroadcastId(self.next_broadcast);
363        self.next_broadcast += 1;
364        self.through_sbeb(cx, |b, ccx| b.on_cmd(sbeb::Cmd::Broadcast { id, msg }, ccx));
365    }
366
367    fn send_to(&mut self, to: NodeId, msg: Reply<V>, cx: &mut ProtoCx<'_, Self>) {
368        let msg = Tagged { ets: self.ets, msg };
369        let id = SendId(self.next_send);
370        self.next_send += 1;
371        self.through_sl(cx, |l, ccx| l.on_cmd(sl::Cmd::Send { id, to, msg }, ccx));
372    }
373
374    /// `highest(states)` — the state with the greatest timestamp among those read.
375    fn highest(&self) -> Option<State<V>> {
376        self.states.values().max_by_key(|s| s.valts).cloned()
377    }
378
379    fn on_announce(&mut self, from: NodeId, msg: Announce<V>, cx: &mut ProtoCx<'_, Self>) {
380        if from != self.leader {
381            return;
382        }
383        match msg {
384            // `upon event ⟨ sbeb, Deliver | ℓ, [READ] ⟩`. Idempotent: the reply says what this
385            // process has accepted, which a repeat does not change.
386            Announce::Read => {
387                if self.state_sent && !self.reply_per_redelivery {
388                    return;
389                }
390                self.state_sent = true;
391                let reply = Reply::StateIs {
392                    valts: self.durable.state.valts,
393                    val: self.durable.state.val.clone(),
394                };
395                self.send_to(from, reply, cx);
396            }
397            // `upon event ⟨ sbeb, Deliver | ℓ, [WRITE, v] ⟩ do (valts, val) := (ets, v);
398            //  store(valts, val); trigger ⟨ sl, Send | ℓ, [ACCEPT] ⟩`
399            //
400            // **The write precedes the acceptance, here in this handler.** The ACCEPT is a promise
401            // to a quorum, and a promise with no record behind it is how `EPC4` fails silently.
402            Announce::Write { val } => {
403                if self.accept_sent && !self.reply_per_redelivery {
404                    return;
405                }
406                if self.durable.state.valts != self.ets {
407                    self.durable.state = State { valts: self.ets, val: Some(val) };
408                    cx.storage().set(self.durable.clone());
409                }
410                self.accept_sent = true;
411                self.send_to(from, Reply::Accept, cx);
412            }
413            // `upon event ⟨ sbeb, Deliver | ℓ, [DECIDED, v] ⟩ do epochdecision := v;
414            //  store(epochdecision); trigger ⟨ lep, Decide | epochdecision ⟩`
415            Announce::Decided { val } => {
416                if self.decided {
417                    return;
418                }
419                self.decided = true;
420                self.durable.decision = Some(val.clone());
421                cx.storage().set(self.durable.clone());
422                cx.indicate(Ind::Decide(val));
423            }
424        }
425    }
426
427    fn on_reply(&mut self, from: NodeId, msg: Reply<V>, cx: &mut ProtoCx<'_, Self>) {
428        if !self.is_leader() {
429            return;
430        }
431        match msg {
432            // `upon event ⟨ sl, Deliver | q, [STATE, ts, v] ⟩ do states[q] := (ts, v)`
433            Reply::StateIs { valts, val } => {
434                self.states.insert(from, State { valts, val });
435                self.maybe_write(cx);
436            }
437            // `upon event ⟨ sl, Deliver | q, [ACCEPT] ⟩ do accepted := accepted + 1`, counted by
438            // process rather than by message. See the module's note on repeats.
439            Reply::Accept => {
440                self.accepted.insert(from);
441                self.maybe_decide(cx);
442            }
443        }
444    }
445
446    /// `upon #(states) > N/2 do … trigger ⟨ sbeb, Broadcast | [WRITE, tmpval] ⟩`.
447    fn maybe_write(&mut self, cx: &mut ProtoCx<'_, Self>) {
448        if self.written || self.states.len() <= self.majority() {
449            return;
450        }
451        // `(ts, v) := highest(states); if v ≠ ⊥ then tmpval := v;` — the line the whole algorithm
452        // turns on, and the one that makes a later epoch adopt what an earlier one may have decided.
453        if let Some(highest) = self.highest()
454            && highest.val.is_some()
455        {
456            self.tmpval = highest.val;
457        }
458        self.states.clear();
459        self.written = true;
460        if let Some(val) = self.tmpval.clone() {
461            self.broadcast(Announce::Write { val }, cx);
462        }
463    }
464
465    /// `upon accepted > N/2 do … trigger ⟨ sbeb, Broadcast | [DECIDED, tmpval] ⟩`.
466    fn maybe_decide(&mut self, cx: &mut ProtoCx<'_, Self>) {
467        if self.announced || self.accepted.len() <= self.majority() {
468            return;
469        }
470        self.accepted.clear();
471        self.announced = true;
472        if let Some(val) = self.tmpval.clone() {
473            self.broadcast(Announce::Decided { val }, cx);
474        }
475    }
476
477    fn through_sbeb(
478        &mut self,
479        cx: &mut ProtoCx<'_, Self>,
480        f: impl FnOnce(
481            &mut StubbornBroadcast<Tagged<Announce<V>>>,
482            &mut ProtoCx<'_, StubbornBroadcast<Tagged<Announce<V>>>>,
483        ),
484    ) {
485        let mut inds = self.sbeb.run(cx, Wire::Announce, f);
486        for sbeb::Ind::Deliver { from, msg } in inds.drain(..) {
487            self.on_announce(from, msg.msg, cx);
488        }
489        self.sbeb.reclaim(inds);
490    }
491
492    fn through_sl(
493        &mut self,
494        cx: &mut ProtoCx<'_, Self>,
495        f: impl FnOnce(
496            &mut StubbornLink<Tagged<Reply<V>>>,
497            &mut ProtoCx<'_, StubbornLink<Tagged<Reply<V>>>>,
498        ),
499    ) {
500        let mut inds = self.sl.run(cx, Wire::Reply, f);
501        for sl::Ind::Deliver { from, msg } in inds.drain(..) {
502            self.on_reply(from, msg.msg, cx);
503        }
504        self.sl.reclaim(inds);
505    }
506}
507
508impl<V: Clone> Protocol for LoggedEpochConsensus<V> {
509    type Cmd = Cmd<V>;
510    type Ind = Ind<V>;
511    type Msg = Wire<V>;
512    type Scope = core::convert::Infallible;
513    type Note = crate::Note;
514    type Meta = Durable<V>;
515    /// An epoch accepts at most one value and decides at most one. Nothing accumulates.
516    type Entry = core::convert::Infallible;
517
518    fn on_cmd(&mut self, cmd: Cmd<V>, cx: &mut ProtoCx<'_, Self>) {
519        if self.aborted {
520            return;
521        }
522        match cmd {
523            // `upon event ⟨ lep, Propose | v ⟩ do tmpval := v; … // only leader ℓ`
524            Cmd::Propose(v) => {
525                if self.is_leader() {
526                    self.tmpval = Some(v);
527                    self.broadcast(Announce::Read, cx);
528                }
529            }
530            // `upon event ⟨ lep, Abort ⟩ do trigger ⟨ lep, Aborted | (valts, val) ⟩; halt;`
531            Cmd::Abort => {
532                self.aborted = true;
533                cx.indicate(Ind::Aborted(self.durable.state.clone()));
534            }
535        }
536    }
537
538    /// `such that ts = ets`, applied at the door.
539    ///
540    /// Unlike [`crate::epoch_consensus`], the link beneath keeps no duplicate set for a foreign
541    /// message to poison — it deduplicates nothing at all. The guard is here for the safety reason
542    /// alone: an acceptance recorded at the wrong timestamp is an acceptance that never happened.
543    fn on_msg(&mut self, from: NodeId, msg: Wire<V>, cx: &mut ProtoCx<'_, Self>) {
544        if self.aborted {
545            return;
546        }
547        match msg {
548            Wire::Announce(m) if m.ets == self.ets => {
549                self.through_sbeb(cx, |b, ccx| b.on_msg(from, m, ccx))
550            }
551            Wire::Reply(m) if m.ets == self.ets => {
552                self.through_sl(cx, |l, ccx| l.on_msg(from, m, ccx))
553            }
554            _ => {}
555        }
556    }
557
558    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
559        if self.aborted {
560            return;
561        }
562        self.through_sbeb(cx, |b, ccx| b.on_timer(id, ccx));
563        self.through_sl(cx, |l, ccx| l.on_timer(id, ccx));
564    }
565
566    /// `upon event ⟨ lep, Init | state ⟩ do (valts, val) := state; store(valts, val); …`
567    ///
568    /// The state came in through the constructor; this is where it becomes durable, before this
569    /// instance answers anything.
570    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
571        cx.storage().set(self.durable.clone());
572    }
573
574    /// `upon event ⟨ lep, Recovery ⟩ do retrieve(valts, val)`.
575    ///
576    /// `epochdecision` comes back with it, because they share one metadata value. Nothing is
577    /// re-indicated: a process that decided before it went down told the layer above at the time,
578    /// and it is that layer's own record — not a second `Decide` from here — that restores it. See
579    /// [`LoggedEpochConsensus::epoch_decision`].
580    fn on_recovery(&mut self, cx: &mut ProtoCx<'_, Self>) {
581        if let Some(durable) = cx.storage().get().cloned() {
582            self.decided = durable.decision.is_some();
583            self.durable = durable;
584        }
585    }
586}