recon_core/cx.rs
1//! Where a protocol's effects go, and how it is told the time.
2
3use crate::store::{
4 FullSlotStore, KeyedSlot, KeyedSlotStore, NoStore, SeqSlot, Slot, SlotStore, Store,
5};
6use crate::{Effect, NodeId, Time, TimerId};
7use core::convert::Infallible;
8use core::marker::PhantomData;
9use core::time::Duration;
10use rand::RngCore;
11
12/// Receives the effects a protocol emits.
13///
14/// The core deliberately has no opinion about how effects are stored. A driver that wants
15/// amortised allocation passes a `Vec` and reuses it across events; a `no_std` or
16/// latency-sensitive driver passes a fixed-capacity sink; a test passes one that merely counts.
17/// Protocol code is identical in every case, because a protocol only ever calls
18/// [`Cx::send`] and friends.
19pub trait EffectSink<M, I> {
20 fn emit(&mut self, effect: Effect<M, I>);
21}
22
23impl<M, I> EffectSink<M, I> for Vec<Effect<M, I>> {
24 fn emit(&mut self, effect: Effect<M, I>) {
25 self.push(effect);
26 }
27}
28
29/// Receives the decisions a protocol narrates.
30///
31/// Separate from [`EffectSink`], and deliberately. An effect is *deferred*: it describes something
32/// the driver will do on the protocol's behalf. A note describes something that has **already
33/// happened**, at a point inside the handler — so it is recorded at the moment of the call, in the
34/// handler's own text, for the same reason [`crate::Store`] is not an effect either.
35///
36/// Nothing a protocol can observe reveals whether anything is listening, so no behaviour can
37/// depend on it: a run is identical whether or not it was read.
38pub trait NoteSink<N> {
39 fn note(&mut self, note: N);
40}
41
42impl<N> NoteSink<N> for Vec<N> {
43 fn note(&mut self, note: N) {
44 self.push(note);
45 }
46}
47
48/// Nobody is listening.
49///
50/// The ordinary case: a run pays for narration only when something asked to read it. A driver
51/// keeps one of these and hands it to every context it builds, exactly as it would a real sink —
52/// so that a protocol's code is identical either way, which is what makes narrating unable to
53/// change the run.
54pub struct NoNotes;
55
56impl<N> NoteSink<N> for NoNotes {
57 fn note(&mut self, _note: N) {}
58}
59
60/// Translates a child's effects into a parent's terms as they are emitted.
61///
62/// This is what makes composition free of intermediate buffers: the child pushes, the mapper
63/// re-wraps, and the parent's sink receives — in one step, with nothing collected in between.
64///
65/// The mappers are `Fn` rather than `fn`, so a parent can stamp with its own state. Normally they
66/// are still enum variant constructors and a wrong one is a type error; what a bare pointer
67/// forbade is the case where the *parent* owns the child's identity. `total_order_broadcast` is the
68/// first: it holds one consensus instance per round, and the round is its concept rather than the
69/// consensus's, so unlike `epoch_consensus` — whose instance stamps its own messages because the
70/// epoch is its identity — the stamp cannot live in the child. `Effect::map` already took `FnOnce`;
71/// the sinks were the last place pinning a pointer.
72struct MapSink<'p, PM, PI, CM, CI, FM, FI> {
73 parent: &'p mut dyn EffectSink<PM, PI>,
74 msg: FM,
75 ind: FI,
76 child: PhantomData<(CM, CI)>,
77}
78
79impl<PM, PI, CM, CI, FM, FI> EffectSink<CM, CI> for MapSink<'_, PM, PI, CM, CI, FM, FI>
80where
81 FM: Fn(CM) -> PM,
82 FI: Fn(CI) -> PI,
83{
84 fn emit(&mut self, effect: Effect<CM, CI>) {
85 self.parent.emit(effect.map(&self.msg, &self.ind));
86 }
87}
88
89/// Translates a child's outgoing effects, but hands its indications back to the parent.
90///
91/// A parent almost never wants to forward a child's indications untouched: the child reporting
92/// "a message arrived" is an *input* to the parent's logic, not an output of it. The parent
93/// cannot react during the call — it is already borrowed by the child — so indications are
94/// collected and processed once the child returns.
95struct ConsumeSink<'p, 'c, PM, PI, CM, CI, FM> {
96 parent: &'p mut dyn EffectSink<PM, PI>,
97 collected: &'c mut Vec<CI>,
98 msg: FM,
99 child: PhantomData<CM>,
100}
101
102impl<PM, PI, CM, CI, FM> EffectSink<CM, CI> for ConsumeSink<'_, '_, PM, PI, CM, CI, FM>
103where
104 FM: Fn(CM) -> PM,
105{
106 fn emit(&mut self, effect: Effect<CM, CI>) {
107 match effect {
108 Effect::Send { to, msg } => self.parent.emit(Effect::Send { to, msg: (self.msg)(msg) }),
109 // A timer belongs to whoever registered it, so it passes straight through: there is
110 // nothing in it for a parent to re-wrap.
111 Effect::SetTimer { after, id } => self.parent.emit(Effect::SetTimer { after, id }),
112 Effect::Indicate(ind) => self.collected.push(ind),
113 }
114 }
115}
116
117/// A protocol's window onto the world.
118///
119/// Supplies the current time and a seeded randomness source, and receives every effect the
120/// protocol emits. Nothing else reaches a protocol: given the same state, event, `now`, and RNG
121/// stream, it behaves identically every time.
122pub struct Cx<'a, M, I, N, Me, En> {
123 sink: &'a mut dyn EffectSink<M, I>,
124 now: Time,
125 rng: &'a mut dyn RngCore,
126 store: &'a mut dyn Store<Me, En>,
127 /// Where narrated decisions go. Shared down the whole composition like `next_timer`, so a
128 /// child's note reaches the run without the parent handling it. [`NoNotes`] when nobody is
129 /// listening, which is the ordinary case.
130 notes: &'a mut dyn NoteSink<N>,
131 /// Where registered timers get their identities. Owned by the driver and shared down the whole
132 /// composition, so an identity is unique to a run rather than to a layer — two layers each
133 /// starting from zero would each accept the other's expiry.
134 next_timer: &'a mut u64,
135}
136
137impl<'a, M, I, N, Me, En> Cx<'a, M, I, N, Me, En> {
138 /// Build a context over any sink, any store, and any audience for what it narrates.
139 ///
140 /// Pass [`NoNotes`] when nothing is listening, which is the ordinary case. It is a parameter
141 /// rather than an option so that the protocol's own code is the same either way — which is
142 /// what makes narrating unable to change the run.
143 pub fn new(
144 sink: &'a mut dyn EffectSink<M, I>,
145 now: Time,
146 rng: &'a mut dyn RngCore,
147 store: &'a mut dyn Store<Me, En>,
148 next_timer: &'a mut u64,
149 notes: &'a mut dyn NoteSink<N>,
150 ) -> Self {
151 Cx { sink, now, rng, store, next_timer, notes }
152 }
153
154 /// Transmit `msg` to `to`.
155 pub fn send(&mut self, to: NodeId, msg: M) {
156 self.sink.emit(Effect::Send { to, msg });
157 }
158
159 /// Raise an indication to the layer above.
160 pub fn indicate(&mut self, ind: I) {
161 self.sink.emit(Effect::Indicate(ind));
162 }
163
164 /// Register a timer for `after`, and take a handle naming it.
165 ///
166 /// The handle is what a later expiry is compared against, so a protocol can tell the timer it
167 /// is waiting on from one it has superseded. It says nothing about which protocol registered
168 /// it or where that protocol sits in a composition.
169 pub fn set_timer(&mut self, after: Duration) -> TimerId {
170 let id = TimerId(*self.next_timer);
171 *self.next_timer += 1;
172 self.sink.emit(Effect::SetTimer { after, id });
173 id
174 }
175
176 /// Record a decision this protocol has taken.
177 ///
178 /// Recorded at the point of the call, not deferred like an effect: a note says what has already
179 /// happened, and its place in the run is where the handler put it.
180 ///
181 /// **Worth narrating only where the record of effects cannot say it.** A note beside
182 /// `indicate` restating the same thing adds nothing a reader could not already see and can
183 /// drift from it. What no trace can hold is a decision that produced *no* effect — a message
184 /// refused, a timestamp already passed, an announcement not made — and that is the case whose
185 /// absence has cost this project the most.
186 ///
187 /// Uncallable for a protocol whose `Note` is uninhabited: there is no value to pass. The
188 /// absence is checked rather than trusted, exactly as it is for a scope event.
189 ///
190 /// ```compile_fail
191 /// # use recon_core::{Cx, NodeId, ProtoCx, Protocol, TimerId};
192 /// struct Quiet;
193 /// impl Protocol for Quiet {
194 /// type Cmd = ();
195 /// type Ind = ();
196 /// type Msg = ();
197 /// type Scope = core::convert::Infallible;
198 /// // Narrates nothing, and so cannot.
199 /// type Note = core::convert::Infallible;
200 /// type Meta = core::convert::Infallible;
201 /// type Entry = core::convert::Infallible;
202 /// fn on_cmd(&mut self, (): (), cx: &mut ProtoCx<'_, Self>) {
203 /// cx.note(());
204 /// }
205 /// fn on_msg(&mut self, _: NodeId, (): (), _: &mut ProtoCx<'_, Self>) {}
206 /// fn on_timer(&mut self, _: TimerId, _: &mut ProtoCx<'_, Self>) {}
207 /// }
208 /// ```
209 pub fn note(&mut self, note: N) {
210 self.notes.note(note);
211 }
212
213 /// This protocol's durable state. A write does not return until it would survive a crash, so
214 /// a protocol may record a promise and then make it. See [`crate::store`].
215 pub fn storage(&mut self) -> &mut dyn Store<Me, En> {
216 &mut *self.store
217 }
218
219 /// The current time — virtual under simulation, real under a live driver.
220 pub fn now(&self) -> Time {
221 self.now
222 }
223
224 /// The seeded randomness source. Reproducible for a given seed.
225 pub fn rng(&mut self) -> &mut dyn RngCore {
226 self.rng
227 }
228
229 /// Run `f` against a child's context, translating everything it emits into this protocol's
230 /// terms on the way out.
231 ///
232 /// The mappers are normally enum variant constructors, so a wrong one is a type error.
233 ///
234 /// The child is handed a store it cannot write to, so only a child keeping nothing durably can
235 /// be composed. See [`crate::store::NoStore`].
236 /// Nothing is buffered: a child's effect is re-wrapped as it is emitted and passed straight
237 /// on to this context's sink.
238 pub fn with_child<CM, CI>(
239 &mut self,
240 msg: impl Fn(CM) -> M,
241 ind: impl Fn(CI) -> I,
242 f: impl FnOnce(&mut Cx<'_, CM, CI, N, Infallible, Infallible>),
243 ) {
244 let mut mapped = MapSink { parent: &mut *self.sink, msg, ind, child: PhantomData };
245 let mut none = NoStore;
246 let mut child = Cx {
247 sink: &mut mapped,
248 now: self.now,
249 rng: &mut *self.rng,
250 store: &mut none,
251 next_timer: &mut *self.next_timer,
252 notes: &mut *self.notes,
253 };
254 f(&mut child);
255 }
256
257 /// Run `f` against a child's context, forwarding what it sends and schedules but collecting
258 /// its indications into `collected` for this protocol to handle itself.
259 ///
260 /// This is the usual shape. A child's indication is the parent's input — the stubborn link
261 /// reporting a delivery is what the perfect link deduplicates — so it must be consumed, not
262 /// passed through. `collected` belongs to the caller and is reused across events.
263 pub fn with_child_consuming<CM, CI>(
264 &mut self,
265 msg: impl Fn(CM) -> M,
266 collected: &mut Vec<CI>,
267 f: impl FnOnce(&mut Cx<'_, CM, CI, N, Infallible, Infallible>),
268 ) {
269 let mut sink = ConsumeSink { parent: &mut *self.sink, collected, msg, child: PhantomData };
270 let mut none = NoStore;
271 let mut child = Cx {
272 sink: &mut sink,
273 now: self.now,
274 rng: &mut *self.rng,
275 store: &mut none,
276 next_timer: &mut *self.next_timer,
277 notes: &mut *self.notes,
278 };
279 f(&mut child);
280 }
281
282 /// [`Cx::with_child_consuming`], for a child that keeps durable state of its own.
283 ///
284 /// The child is handed a view of `slot` — the part of this protocol's record that belongs to
285 /// it. Its `get` projects, and its `set` reads this record back, replaces the child's part, and
286 /// writes the whole thing down again. **That is one write, not two**, which is what stops a
287 /// crash landing between a parent's record and its child's.
288 ///
289 /// Prefer [`Cx::with_child_consuming`] wherever the child keeps nothing: it hands a
290 /// [`NoStore`], and a child that cannot write is one fewer thing to reason about. This exists
291 /// because Algorithm 5.10 needs it — a protocol that keeps `(ets, ℓ, decision)` of its own and
292 /// composes two children that each keep a record too, and whose recovery reads its children's
293 /// records by name.
294 ///
295 /// The child's `Entry` is uninhabited: a child composed this way cannot append. That is the
296 /// right default and most durable children want nothing else; one that *does* append is
297 /// composed with [`Cx::with_durable_child`] and a [`SeqSlot`] as well.
298 pub fn with_durable_child_consuming<CM, CI, CMe>(
299 &mut self,
300 msg: impl Fn(CM) -> M,
301 collected: &mut Vec<CI>,
302 slot: Slot<Me, CMe>,
303 f: impl FnOnce(&mut Cx<'_, CM, CI, N, CMe, Infallible>),
304 ) {
305 let mut sink = ConsumeSink { parent: &mut *self.sink, collected, msg, child: PhantomData };
306 let mut store = SlotStore { parent: &mut *self.store, slot };
307 let mut child = Cx {
308 sink: &mut sink,
309 now: self.now,
310 rng: &mut *self.rng,
311 store: &mut store,
312 next_timer: &mut *self.next_timer,
313 notes: &mut *self.notes,
314 };
315 f(&mut child);
316 }
317
318 /// [`Cx::with_durable_child_consuming`], for one member of a *family* of durable children.
319 ///
320 /// A [`Slot`] names a fixed place, which is right when a parent has one child of a kind. A
321 /// parent holding one instance per round needs a place per round, and `key` supplies it — as
322 /// data rather than as something the slot captured, so the slot is still one fixed function.
323 ///
324 /// The parent owns the keyspace: two members handed the same key share a record.
325 pub fn with_keyed_durable_child_consuming<CM, CI, CMe, K>(
326 &mut self,
327 msg: impl Fn(CM) -> M,
328 collected: &mut Vec<CI>,
329 slot: KeyedSlot<Me, CMe, K>,
330 key: K,
331 f: impl FnOnce(&mut Cx<'_, CM, CI, N, CMe, Infallible>),
332 ) {
333 let mut sink = ConsumeSink { parent: &mut *self.sink, collected, msg, child: PhantomData };
334 let mut store = KeyedSlotStore { parent: &mut *self.store, slot, key };
335 let mut child = Cx {
336 sink: &mut sink,
337 now: self.now,
338 rng: &mut *self.rng,
339 store: &mut store,
340 next_timer: &mut *self.next_timer,
341 notes: &mut *self.notes,
342 };
343 f(&mut child);
344 }
345
346 /// [`Cx::with_durable_child_consuming`], for a child that keeps metadata **and appends**.
347 ///
348 /// The child's record lives in `slot` of this protocol's, exactly as before, and its entries go
349 /// into *this* protocol's sequence through `entries` — **one sequence, not two**. Two sequences
350 /// would have no order between them, and a recovery replaying both would be inventing one.
351 ///
352 /// The positions the child reads back are this protocol's, and therefore sparse: its third entry
353 /// may sit at position seven. [`SeqSlot`] says why that is all a cursor needs.
354 ///
355 /// [`Cx::with_durable_child_consuming`] remains the one to reach for. A child that cannot append
356 /// is one fewer thing to reason about, and its signature says it cannot; this exists because the
357 /// fail-recovery total-order broadcast keeps a durable record of its own and composes a child
358 /// that appends, which nothing did before.
359 pub fn with_durable_child<CM, CI, CMe, CEn>(
360 &mut self,
361 msg: impl Fn(CM) -> M,
362 collected: &mut Vec<CI>,
363 slot: Slot<Me, CMe>,
364 entries: SeqSlot<En, CEn>,
365 f: impl FnOnce(&mut Cx<'_, CM, CI, N, CMe, CEn>),
366 ) {
367 let mut sink = ConsumeSink { parent: &mut *self.sink, collected, msg, child: PhantomData };
368 let mut store = FullSlotStore { parent: &mut *self.store, slot, entries };
369 let mut child = Cx {
370 sink: &mut sink,
371 now: self.now,
372 rng: &mut *self.rng,
373 store: &mut store,
374 next_timer: &mut *self.next_timer,
375 notes: &mut *self.notes,
376 };
377 f(&mut child);
378 }
379}