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}