Skip to main content

recon_protocols/
logged_link.rs

1//! Logged perfect point-to-point links.
2//!
3//! Cachin, Guerraoui & Rodrigues, Module 2.4 and Algorithm 2.3 ("Log Delivered").
4//!
5//! **Status: transcription. Space: unbounded — and on disk. Write cost: one append per message.**
6//! `delivered` grows with every distinct message log-delivered and nothing retires an entry, as
7//! [`crate::perfect_link`] does, except that here the growth is a file rather than a heap. See
8//! `docs/bounded-space.md`; the fix is a delivered *cursor* rather than a delivered *set*, and it
9//! is a change with a proposal.
10//!
11//! # The indication carries the log, not the message
12//!
13//! This is the whole of what the fail-recovery model changes about an interface, and the reason
14//! this protocol sits at the bottom of the stack where it can be seen.
15//!
16//! A crash-stop protocol notifies the layer above by triggering `⟨ Deliver | m ⟩` once. A
17//! crash-recovery protocol cannot. It may crash immediately afterwards, and then neither it nor
18//! the layer above nor anyone else will ever know the indication happened — the message is lost
19//! in a notification that no longer exists. So the module writes the message into a set in stable
20//! storage, and the indication says only that *the set may have changed*:
21//!
22//! ```text
23//! upon event ⟨ lpl, Init ⟩ do
24//!     delivered := ∅;
25//!     store(delivered);
26//!
27//! upon event ⟨ lpl, Recovery ⟩ do
28//!     retrieve(delivered);
29//!     trigger ⟨ lpl, Deliver | delivered ⟩;
30//!
31//! upon event ⟨ lpl, Send | q, m ⟩ do
32//!     trigger ⟨ sl, Send | q, m ⟩;
33//!
34//! upon event ⟨ sl, Deliver | p, m ⟩ do
35//!     if not exists (p′, m′) ∈ delivered such that m′ = m then
36//!         delivered := delivered ∪ {(p, m)};
37//!         store(delivered);
38//!         trigger ⟨ lpl, Deliver | delivered ⟩;
39//! ```
40//!
41//! The layer above reads the set rather than receiving a message, and must be idempotent: the
42//! same set arrives again after every restart.
43//!
44//! # Reliable delivery is weaker here, and necessarily
45//!
46//! Module 2.3 promises delivery if a *correct* process sends to a correct process. Module 2.4
47//! promises it only if a process that **never crashes** does. The difference is not fussiness: a
48//! sender that crashes immediately after being asked to send may have no record that it was ever
49//! asked, and in the crash-recovery model a process that crashes and recovers is still correct.
50//! There is nothing left in the system to retransmit.
51//!
52//! # What the durable record buys
53//!
54//! [`crate::perfect_link`] keeps its `delivered` set in memory, so a restart forgets it and the
55//! sender's next retransmission is delivered a second time — `no_duplication_does_not_survive_the_recipient_restarting`
56//! records exactly that. Here the record survives, so LPL2 holds across incarnations rather than
57//! within one. That is the entire purchase, and it is what stable storage is for.
58//!
59//! # Departures from the page
60//!
61//! - None worth the name for initialisation: `⟨ Init ⟩` and `⟨ Recovery ⟩` are both here, exactly
62//!   one fires, and `Init` performs the book's initial store — of the empty counter, this
63//!   implementation's metadata, rather than of the empty set. That write is what makes the branch
64//!   real rather than emergent: after it, storage holds something, so every later restart
65//!   recovers. The constructor does volatile setup only, because it runs in both cases and cannot
66//!   emit effects.
67//! - `delivered` is keyed by sender and a per-sender sequence number rather than by message
68//!   content, for the reason [`crate::perfect_link`] gives: identical content sent twice is two
69//!   messages and must be delivered twice.
70//! - The send counter is therefore durable, held in the metadata and restored on recovery. See
71//!   the note below; the obligation comes with the departure that created it.
72//! - No `Stop`: the stubborn link beneath retransmits for ever, which in this model is a feature —
73//!   it is how a process that was down when a message was sent receives it after recovering.
74//!
75//! # The send counter is the metadata
76//!
77//! Content-keyed deduplication, as the book has it, needs no counter: a payload names itself, and
78//! a restart changes nothing. Identifier-keyed deduplication owes the counter the same durability
79//! as the set it keys — and this is the one it cannot recompute, because the log holds messages
80//! *received*, saying nothing about what this process sent.
81//!
82//! Left volatile, a restarted sender re-mints `(me, 1), (me, 2), …`. A recipient whose durable
83//! log already holds those identifiers discards the new payloads as duplicates, permanently, while
84//! the stubborn link beneath retransmits them for ever. LPL1 does not formally break — it promises
85//! delivery only from a process that never crashes — but the loss is silent, and the simulator
86//! restarts processes freely.
87//!
88//! So the counter is written before the message it names goes out, which is what the metadata is
89//! for: one small value, rewritten. A torn write discards that handler's sends along with it, so
90//! the counter is never behind an identifier already on the wire; a write that lands before a
91//! crash burns an identifier, and a gap costs nothing. Keying the durable set by content instead
92//! would also work, at the price of the departure above.
93//!
94//! # The write cost is linear, not quadratic
95//!
96//! One append per message log-delivered, and one metadata rewrite per message sent — a single
97//! `u64`, whose cost does not grow with what preceded it. Rewriting the whole *set* on each
98//! arrival would cost `O(n²)` over a run, which is the failure mode `docs/bounded-space.md`
99//! calls unbounded *work*: small per item, and growing without limit. That is the distinction
100//! the split between metadata and entries exists to make.
101//!
102//! Deduplication is a lookup by identifier rather than a scan over pairs, for the same reason:
103//! a scan would make the *read* side quadratic while the write side was linear, and a cost note
104//! that mentions only one of them is not a cost note.
105//!
106//! What is left is the indication, which carries the whole set and therefore copies it on every
107//! arrival. That one is the book's interface rather than an implementation choice — see the note
108//! at the top on why the indication cannot be a message — and it goes away with the same change
109//! that bounds the record: a delivered *cursor* instead of a delivered *set*.
110//!
111//! The record still grows without limit, so this remains a transcription. What appending changes
112//! is the cost of adding to it, not its size.
113
114use core::time::Duration;
115use recon_core::{Child, NodeId, Position, ProtoCx, Protocol, TimerId};
116use serde::{Deserialize, Serialize};
117use std::collections::BTreeMap;
118
119use crate::perfect_link::MsgId;
120use crate::stubborn_link::{self as sl, SendId, StubbornLink};
121
122/// What goes on the wire: the payload and the identifier that names it.
123///
124/// The same shape [`crate::perfect_link`] uses, and for the same reason — deduplication is by
125/// identifier, so that identical content sent twice is two messages.
126pub type Wire<P> = crate::perfect_link::Wire<P>;
127
128/// Requests from the layer above.
129#[derive(Debug, Clone, PartialEq, Eq)]
130pub enum Cmd<P> {
131    Send { to: NodeId, msg: P },
132}
133
134/// Indications to the layer above.
135///
136/// One variant, carrying the durable log rather than a message. See the module note.
137#[derive(Debug, Clone, PartialEq, Eq)]
138pub enum Ind<P: Ord> {
139    /// The durable set of log-delivered messages may have changed. Here it is.
140    Delivered(Log<P>),
141}
142
143/// The set of messages log-delivered, in stable storage.
144///
145/// Ordered, so that reading it is deterministic and two processes holding the same set see it the
146/// same way — the ordered-maps rule applies to what is written down as much as to what is held.
147#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
148#[serde(bound(deserialize = "P: Ord + Deserialize<'de>"))]
149pub struct Log<P: Ord> {
150    /// Keyed by identifier, not a set of pairs: deduplication asks "have I seen this id", and
151    /// over a set of pairs that question is a scan. One arrival is `O(log n)` here and `O(n)`
152    /// there, which over a run is the difference between linear and quadratic *work* — the thing
153    /// `docs/bounded-space.md` counts alongside space.
154    entries: BTreeMap<MsgId, P>,
155}
156
157impl<P: Ord> Default for Log<P> {
158    fn default() -> Self {
159        Log { entries: BTreeMap::new() }
160    }
161}
162
163impl<P: Ord> Log<P> {
164    /// Every message log-delivered, with the identifier naming its sender.
165    pub fn entries(&self) -> impl Iterator<Item = (&MsgId, &P)> + '_ {
166        self.entries.iter()
167    }
168
169    /// How many messages have been log-delivered.
170    pub fn len(&self) -> usize {
171        self.entries.len()
172    }
173
174    pub fn is_empty(&self) -> bool {
175        self.entries.is_empty()
176    }
177
178    /// Whether `id` has been log-delivered.
179    pub fn contains(&self, id: MsgId) -> bool {
180        self.entries.contains_key(&id)
181    }
182}
183
184/// Perfect-link guarantees over log-delivery, so that they hold across a restart.
185#[derive(Debug)]
186pub struct LoggedLink<P: Clone + Ord> {
187    me: NodeId,
188    /// Durable, and restored on recovery: it keys a set that outlives this incarnation.
189    seq: u64,
190    /// The durable set. Volatile here, written down on every change, and retrieved on recovery.
191    delivered: Log<P>,
192    link: Child<StubbornLink<Wire<P>>>,
193}
194
195impl<P: Clone + Ord> LoggedLink<P> {
196    /// Log-deliver for `me`, retransmitting every `retransmit`.
197    pub fn new(me: NodeId, retransmit: Duration) -> Self {
198        LoggedLink {
199            me,
200            seq: 0,
201            delivered: Log::default(),
202            link: Child::new(StubbornLink::new(retransmit)),
203        }
204    }
205
206    /// What has been log-delivered. The same value the layer above is handed.
207    pub fn log(&self) -> &Log<P> {
208        &self.delivered
209    }
210}
211
212impl<P: Clone + Ord> LoggedLink<P> {
213    fn with_link(
214        &mut self,
215        cx: &mut ProtoCx<'_, Self>,
216        f: impl FnOnce(&mut StubbornLink<Wire<P>>, &mut ProtoCx<'_, StubbornLink<Wire<P>>>),
217    ) {
218        let mut inds = self.link.run(cx, core::convert::identity, f);
219        for sl::Ind::Deliver { from, msg } in inds.drain(..) {
220            self.on_arrival(from, msg, cx);
221        }
222        self.link.reclaim(inds);
223    }
224
225    /// `upon event ⟨ sl, Deliver | p, m ⟩`.
226    fn on_arrival(&mut self, _from: NodeId, wire: Wire<P>, cx: &mut ProtoCx<'_, Self>) {
227        if self.delivered.contains(wire.id) {
228            // Seen before, in this incarnation or an earlier one. Nothing changed, so nothing is
229            // written and nothing is announced.
230            return;
231        }
232        let entry = (wire.id, wire.payload);
233        self.delivered.entries.insert(entry.0, entry.1.clone());
234        // One append, not a rewrite — and durable before the layer above is told, so a crash in
235        // between leaves the message in the record rather than in a lost notification.
236        cx.storage().append(entry);
237        cx.indicate(Ind::Delivered(self.delivered.clone()));
238    }
239}
240
241impl<P: Clone + Ord> Protocol for LoggedLink<P> {
242    type Cmd = Cmd<P>;
243    type Ind = Ind<P>;
244    type Msg = Wire<P>;
245    type Scope = core::convert::Infallible;
246    type Note = crate::Note;
247    /// The send counter: one small value, rewritten before each send. See the module note.
248    type Meta = u64;
249    /// One log-delivered message; appended, so the write cost is linear rather than quadratic.
250    type Entry = (MsgId, P);
251
252    /// `⟨ lpl, Init ⟩ do delivered := ∅; store(delivered)`.
253    ///
254    /// The initial write is what makes the branch real: after it, storage holds something, so
255    /// every later restart takes the recovery path rather than starting afresh. What is written
256    /// is the send counter at zero — the empty set is the empty sequence of entries, which needs
257    /// no write of its own.
258    fn on_init(&mut self, cx: &mut ProtoCx<'_, Self>) {
259        cx.storage().set(self.seq);
260    }
261
262    fn on_cmd(&mut self, Cmd::Send { to, msg }: Cmd<P>, cx: &mut ProtoCx<'_, Self>) {
263        self.seq += 1;
264        // Durable before the identifier it names is on the wire: a restarted sender that resumed
265        // from zero would have its new messages discarded as duplicates by a log that survived.
266        cx.storage().set(self.seq);
267        let id = MsgId { src: self.me, seq: self.seq };
268        let wire = Wire { id, payload: msg };
269        // Stubbornly, and never stopped: retransmission is what reaches a process that was down.
270        self.with_link(cx, |link, ccx| {
271            link.on_cmd(sl::Cmd::Send { id: SendId(id.seq), to, msg: wire }, ccx)
272        });
273    }
274
275    fn on_msg(&mut self, from: NodeId, msg: Wire<P>, cx: &mut ProtoCx<'_, Self>) {
276        self.with_link(cx, |link, ccx| link.on_msg(from, msg, ccx));
277    }
278
279    fn on_timer(&mut self, id: TimerId, cx: &mut ProtoCx<'_, Self>) {
280        self.with_link(cx, |link, ccx| link.on_timer(id, ccx));
281    }
282
283    /// `upon event ⟨ lpl, Recovery ⟩ do retrieve(delivered); trigger ⟨ lpl, Deliver | delivered ⟩`.
284    ///
285    /// The record is read here rather than handed over, along with the send counter that keys it.
286    /// Nothing else is dispatched until this returns, which is what makes it safe to have held an
287    /// empty index a moment ago.
288    ///
289    /// The layer above is told again: the notification sent before the crash may have been lost
290    /// with the incarnation that sent it.
291    fn on_recovery(&mut self, cx: &mut ProtoCx<'_, Self>) {
292        self.seq = cx.storage().get().copied().unwrap_or(0);
293        let entries: Vec<(MsgId, P)> =
294            cx.storage().read_from(Position::START).into_iter().cloned().collect();
295        self.delivered = Log { entries: entries.into_iter().collect() };
296        cx.indicate(Ind::Delivered(self.delivered.clone()));
297    }
298}