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}