Skip to main content

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}