Skip to main content

ridl_loopback/
lib.rs

1//! `ridl-loopback` — the in-process reference runtime.
2//!
3//! This crate implements the twelve port traits of
4//! [`ridl_rt::port`](https://docs.rs/ridl-rt) over one in-memory store. It is
5//! the first runtime in this workspace (ADR-0020 decision 6, which names the
6//! crate and fixes `ridl-rt` as its only dependency), and it is what the
7//! generated interaction face runs against in a test, an example or a
8//! single-process application.
9//!
10//! What it is not:
11//!
12//! - **Not a transport.** It carries no frame, opens no socket and has no wire
13//!   format. A value published here is read here, in the same process.
14//! - **Not the engine.** The store below is a map and a queue, not the seqlock
15//!   store, the sans-IO session, the platform traits and the scheduler, which
16//!   are outside this repository.
17//! - **Not a checker.** Payload bytes are opaque to it: it carries them and
18//!   never verifies one. `Payload::verify` is the generated face's, on both
19//!   sides.
20//!
21//! # The handles
22//!
23//! A runtime presents one handle type per port role rather than one type
24//! implementing them all, and may also offer an aggregate handle covering the
25//! port set one interface's face needs (ADR-0021 decision 12). This crate
26//! offers both. Its six handles group eleven roles the way that decision
27//! derives the threading split: a handle each for the five roles with a
28//! `&mut self` method, and one handle for the six whose methods all take
29//! `&self`, which are exactly the roles several threads may hold at once. The
30//! twelfth port trait, `Wakeable`, is implemented on every handle, because
31//! each handle wakes its own waiters (see "Waking" below).
32//! [`Loopback`] is the aggregate: it implements all twelve port traits by
33//! delegating to the six role handles it holds, and it is what a generated
34//! `Client`, `Publisher` or `dispatch` is normally built over.
35//!
36//! ```
37//! use ridl_loopback::Loopback;
38//! use ridl_rt::contract::{CatalogHash, CatalogRef, InterfaceNo, Ordinal};
39//! use ridl_rt::port::{SignalReader, SignalWriter};
40//!
41//! let catalog = CatalogRef { name: "face.demo", hash: CatalogHash([0u8; 32]) };
42//! let mut rt = Loopback::new(catalog);
43//!
44//! rt.set(InterfaceNo(1), Ordinal(1), &[42]).expect("staged");
45//! rt.commit();
46//!
47//! let mut out = [0u8; 8];
48//! let sample = rt.read(InterfaceNo(1), Ordinal(1), &mut out).expect("read");
49//! assert_eq!(&out[..sample.len], &[42]);
50//! ```
51//!
52//! [`Loopback::split`] hands out the six role handles for a program that wants
53//! them apart — one thread reading while another publishes, or two callers on
54//! one provider. [`ReaderHandle`] is `Send + Sync`; the other five are `Send`
55//! and driven by one thread each. Nothing here declares either: both follow
56//! from the fields, and the assertions at the bottom of this file pin them.
57//!
58//! # What it reports, and what it cannot
59//!
60//! The loopback holds no catalog descriptor — the descriptor and the catalog
61//! hash arrive with story E16.2 (driftsys/ridl#378) — so it has no member
62//! table, and there is no ordinal it can call unknown, no member it can call
63//! unowned, and no timing annotation it can measure a value's freshness or a
64//! call's remaining time against. What it therefore never returns:
65//! `WriteError::NotOwner`, `RaiseError::NotOwner`, `ServeError::NotOwner`, any
66//! port error's `Contract` variant except
67//! [`FixedReader::read_fixed`](ridl_rt::port::FixedReader::read_fixed)'s, and
68//! `Freshness::Fresh` or `Freshness::Stale`. Nothing detaches, because every
69//! handle holds the store alive, so `Detached` never appears either. The one
70//! bound is the call table's [`Loopback::SLOTS`]: a send with every slot
71//! taken answers `SendError::Busy`. Nothing else is bounded, so `TooLarge`
72//! appears only from [`Loopback::fail_next_settle`], and `Transport::Busy`
73//! reaches a caller only when a provider settles a call with it.
74//!
75//! # Waking
76//!
77//! Every handle implements [`Wakeable`], and a handle stores a waker only
78//! under a kind of key one of its roles observes: a caller handle under
79//! `Interest::Outcome`, kept with each call, and one waker under
80//! `Interest::Slot`, a source handle one waker under `Interest::Event`, and a
81//! handler handle one waker under `Interest::Claim`.
82//! For `Event` and `Claim` the rule is one waker per kind, and a change to
83//! any key of the kind wakes the stored waker, whatever interface it was
84//! registered under (ADR-0021 decision 13); an `Outcome` waker is per call,
85//! and no change but that call's settlement, its `forget`, or the drop of the
86//! caller handle that sent it wakes it (a displacement by another task does,
87//! as for every kind). A settlement wakes the
88//! call's waiter, a raise wakes each source it queues the occurrence for, and
89//! a send wakes each handler that serves the member, and a reclaimed slot of
90//! the call table wakes every caller's `Slot` waiter. A registration whose key
91//! already holds — the outcome is known, an occurrence or a call is waiting, a
92//! slot is free — is woken at once, and so is one under a kind the handle does
93//! not observe. A registration of the waker already stored refreshes it
94//! without waking it. No waker is woken while the store is locked.
95//!
96//! [`Attached::catalog`](ridl_rt::port::Attached::catalog) returns the
97//! `CatalogRef` the runtime was built with, unexamined. ADR-0021 decision 3
98//! places a check of it against the interface's own `CATALOG` in a generated
99//! face's constructor, once, when the face is built; the constructor the Rust
100//! backend emits today performs no such check (driftsys/ridl#448). Either way
101//! it is the face's check and not the runtime's: the loopback carries the
102//! value and compares nothing.
103//!
104//! The crate's as-built design record, with the reasoning behind each of these
105//! choices, is `docs/design/ridl-loopback.md` in this repository.
106
107use std::sync::{Arc, Mutex};
108use std::task::Waker;
109
110use ridl_rt::contract::{CatalogRef, InterfaceNo, Ordinal};
111use ridl_rt::error::CallError;
112use ridl_rt::port::{
113    Attached, Caller, Changed, Claim, ClaimId, Clock, CoherentSignals, Correlation, EventSink,
114    EventSource, FixedReader, Handler, Interest, RaiseError, RawOccurrence, RawSample, ReadError,
115    ScannableSignals, SendError, ServeError, SettleError, SignalReader, SignalWriter,
116    SubscribeError, Wakeable, Watermark, WriteError,
117};
118use ridl_rt::sample::{Duration, Timestamp};
119
120mod handle;
121mod store;
122
123pub use handle::{
124    CallerHandle, HandlerHandle, ReaderHandle, SinkHandle, SourceHandle, WriterHandle,
125};
126
127use handle::{Shared, lock};
128use store::Store;
129
130/// The six role handles of one runtime, as [`Loopback::split`] hands them out.
131pub struct Handles {
132    /// `Attached`, `Clock`, `SignalReader`, `FixedReader`, `ScannableSignals`
133    /// and `CoherentSignals`.
134    pub reader: ReaderHandle,
135    /// `SignalWriter`.
136    pub writer: WriterHandle,
137    /// `EventSource`.
138    pub source: SourceHandle,
139    /// `EventSink`.
140    pub sink: SinkHandle,
141    /// `Caller`.
142    pub caller: CallerHandle,
143    /// `Handler`.
144    pub handler: HandlerHandle,
145}
146
147/// The aggregate handle: one value implementing all twelve port traits by
148/// delegating to the six role handles it holds.
149///
150/// A face is built over one value implementing at least the port traits its
151/// interface needs, and a generated `Client` is commonly bound over
152/// `SignalReader + EventSource + Caller` at once, which no single role handle
153/// satisfies. That is what an aggregate is for (ADR-0021 decision 12). Pass it
154/// by value, or as `&mut` under the forwarding impls of ADR-0021 decision 11.
155pub struct Loopback {
156    shared: Shared,
157    catalog: CatalogRef,
158    handles: Handles,
159}
160
161impl Loopback {
162    /// The number of calls the runtime holds at once. A call holds its slot
163    /// from its send until one of these reclaims it:
164    ///
165    /// - [`Caller::forget`], or the drop of the caller handle that sent it,
166    ///   while it is settled or still waiting for a handler and offered to
167    ///   none;
168    /// - its settlement, when it was forgotten while a handler held it —
169    ///   taken, or offered through `ReadError::ShortClaim` and still waiting
170    ///   (driftsys/ridl#569);
171    /// - the drop of the handler that held it unsettled, when it was
172    ///   forgotten.
173    ///
174    /// So a settled call keeps its slot until it is forgotten, and a claimed
175    /// call keeps its slot after its `forget` until it is settled or its
176    /// handler is dropped. With every slot taken, [`Caller::command`] and
177    /// [`Caller::query`] answer [`SendError::Busy`], on every caller handle,
178    /// because the table is the runtime's.
179    ///
180    /// Sixteen is a small bound, chosen so that a test reaches it in a few
181    /// sends and a program that never forgets a call finds out at once rather
182    /// than after its memory grows. The loopback has no catalog descriptor to
183    /// size a byte budget from (story E16.2), so the slot count is its only
184    /// bound (note F-9 of the async face design).
185    ///
186    /// The generated async client's future (story E11.21) forgets its call
187    /// when it leaves the waiting phase: in the poll that takes the outcome,
188    /// at the call's deadline, and on drop while the call is still waiting.
189    /// So a program that calls through the generated face holds one slot per
190    /// call in flight. A program that sends through the `Caller` port
191    /// directly must forget each correlation itself once it has read the
192    /// outcome, or drop the caller handle, which forgets that handle's calls;
193    /// otherwise it gets `SendError::Busy` from its seventeenth call on.
194    pub const SLOTS: usize = 16;
195
196    /// A runtime attached to `catalog`, with an empty store and its clock at
197    /// [`Timestamp`] 0.
198    ///
199    /// The catalog is carried, not checked: see the crate documentation.
200    #[must_use]
201    pub fn new(catalog: CatalogRef) -> Self {
202        let shared: Shared = Arc::new(Mutex::new(Store::new()));
203        let handles = Handles {
204            reader: ReaderHandle::new(Arc::clone(&shared), catalog),
205            writer: WriterHandle::new(Arc::clone(&shared), catalog),
206            source: SourceHandle::new(Arc::clone(&shared), catalog),
207            sink: SinkHandle::new(Arc::clone(&shared), catalog),
208            caller: CallerHandle::new(Arc::clone(&shared), catalog),
209            handler: HandlerHandle::new(Arc::clone(&shared), catalog),
210        };
211        Loopback {
212            shared,
213            catalog,
214            handles,
215        }
216    }
217
218    /// Hands out the six role handles this aggregate holds. Every one of them
219    /// keeps the same store, so a value published through `writer` is read
220    /// through `reader`.
221    #[must_use]
222    pub fn split(self) -> Handles {
223        self.handles
224    }
225
226    /// An additional reader handle on the same store.
227    #[must_use]
228    pub fn reader(&self) -> ReaderHandle {
229        ReaderHandle::new(Arc::clone(&self.shared), self.catalog)
230    }
231
232    /// An additional writer handle on the same store, with its own staging
233    /// area and its own per channel sequence counters.
234    #[must_use]
235    pub fn writer(&self) -> WriterHandle {
236        WriterHandle::new(Arc::clone(&self.shared), self.catalog)
237    }
238
239    /// An additional event source on the same store, with its own
240    /// subscription set and its own queue.
241    #[must_use]
242    pub fn source(&self) -> SourceHandle {
243        SourceHandle::new(Arc::clone(&self.shared), self.catalog)
244    }
245
246    /// An additional event sink on the same store, with its own sequence
247    /// counter.
248    #[must_use]
249    pub fn sink(&self) -> SinkHandle {
250        SinkHandle::new(Arc::clone(&self.shared), self.catalog)
251    }
252
253    /// An additional caller on the same store, with its own sequence counter.
254    /// Two callers on one provider therefore never collide on a sequence
255    /// number (driftsys/ridl#308).
256    #[must_use]
257    pub fn caller(&self) -> CallerHandle {
258        CallerHandle::new(Arc::clone(&self.shared), self.catalog)
259    }
260
261    /// An additional handler on the same store.
262    ///
263    /// Every handler draws from the one queue of waiting calls, filtered by
264    /// what it has served: a handler that has served nothing is presented
265    /// every waiting call, and one that has served members is presented only
266    /// those. A claim belongs to the handler it was presented to.
267    #[must_use]
268    pub fn handler(&self) -> HandlerHandle {
269        HandlerHandle::new(Arc::clone(&self.shared), self.catalog)
270    }
271
272    /// Advances the clock by `by`.
273    ///
274    /// The clock is a counter this method moves and nothing else moves. It
275    /// never reads wall-clock time, so a round trip over this runtime produces
276    /// the same timestamps on every run and on every machine. It saturates
277    /// rather than overflowing.
278    ///
279    /// # Panics
280    ///
281    /// When `by` is negative. A clock that ran backwards would stamp an
282    /// envelope before one already stamped.
283    pub fn advance(&mut self, by: Duration) {
284        lock(&self.shared).advance(by);
285    }
286
287    /// Supplies the value of a `fixed` (ridl §8), which is provisioned into a
288    /// runtime rather than published by application code.
289    ///
290    /// Until a `fixed` is provisioned, reading it answers
291    /// `ReadError::Contract(Contract::UnknownInteraction)`: the loopback holds
292    /// no value and, having no member table, cannot say anything narrower.
293    pub fn provision_fixed(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) {
294        lock(&self.shared).provision_fixed(iface, ord, bytes);
295    }
296
297    /// Makes the next `Handler::settle` on this runtime fail with
298    /// `SettleError::TooLarge` and record no outcome.
299    ///
300    /// This is the one fault this runtime injects, and the one place a
301    /// `SettleError` other than `UnknownClaim` comes from. It exists because
302    /// the generated `dispatch` counts a claim only once the handler has
303    /// accepted its settlement, and nothing else in an in-process runtime can
304    /// make that path fail.
305    pub fn fail_next_settle(&mut self) {
306        lock(&self.shared).fail_next_settle();
307    }
308}
309
310// ---------------------------------------------------------------------------
311// The aggregate's twelve port implementations, each one a delegation.
312// ---------------------------------------------------------------------------
313
314impl Attached for Loopback {
315    fn catalog(&self) -> &CatalogRef {
316        self.handles.reader.catalog()
317    }
318}
319
320impl Clock for Loopback {
321    fn now(&self) -> Timestamp {
322        self.handles.reader.now()
323    }
324}
325
326impl SignalReader for Loopback {
327    fn read(
328        &self,
329        iface: InterfaceNo,
330        ord: Ordinal,
331        out: &mut [u8],
332    ) -> Result<RawSample, ReadError> {
333        self.handles.reader.read(iface, ord, out)
334    }
335}
336
337impl FixedReader for Loopback {
338    fn read_fixed(
339        &self,
340        iface: InterfaceNo,
341        ord: Ordinal,
342        out: &mut [u8],
343    ) -> Result<usize, ReadError> {
344        self.handles.reader.read_fixed(iface, ord, out)
345    }
346}
347
348impl ScannableSignals for Loopback {
349    fn generation(&self, iface: InterfaceNo) -> u64 {
350        self.handles.reader.generation(iface)
351    }
352
353    fn scan(&self, marks: &mut [Watermark], out: &mut [Changed]) -> usize {
354        self.handles.reader.scan(marks, out)
355    }
356}
357
358impl CoherentSignals for Loopback {
359    fn read_coherent(
360        &self,
361        iface: InterfaceNo,
362        ords: &[Ordinal],
363        out: &mut [u8],
364        samples: &mut [RawSample],
365    ) -> Result<usize, ReadError> {
366        self.handles.reader.read_coherent(iface, ords, out, samples)
367    }
368}
369
370impl SignalWriter for Loopback {
371    fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
372        self.handles.writer.set(iface, ord, bytes)
373    }
374
375    fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
376        self.handles.writer.invalidate(iface, ord)
377    }
378
379    fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
380        self.handles.writer.touch(iface, ord)
381    }
382
383    fn commit(&mut self) {
384        self.handles.writer.commit();
385    }
386}
387
388impl EventSource for Loopback {
389    fn subscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), SubscribeError> {
390        self.handles.source.subscribe(iface, ords)
391    }
392
393    fn unsubscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) {
394        self.handles.source.unsubscribe(iface, ords);
395    }
396
397    fn next(&mut self, out: &mut [u8]) -> Result<Option<RawOccurrence>, ReadError> {
398        self.handles.source.next(out)
399    }
400}
401
402impl EventSink for Loopback {
403    fn raise(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), RaiseError> {
404        self.handles.sink.raise(iface, ord, bytes)
405    }
406}
407
408impl Caller for Loopback {
409    fn command(
410        &mut self,
411        iface: InterfaceNo,
412        ord: Ordinal,
413        args: &[u8],
414    ) -> Result<Correlation, SendError> {
415        self.handles.caller.command(iface, ord, args)
416    }
417
418    fn query(
419        &mut self,
420        iface: InterfaceNo,
421        ord: Ordinal,
422        args: &[u8],
423    ) -> Result<Correlation, SendError> {
424        self.handles.caller.query(iface, ord, args)
425    }
426
427    fn ack(&mut self, c: Correlation) -> Option<Result<(), CallError>> {
428        self.handles.caller.ack(c)
429    }
430
431    fn reply(
432        &mut self,
433        c: Correlation,
434        out: &mut [u8],
435    ) -> Result<Option<Result<usize, CallError>>, ReadError> {
436        self.handles.caller.reply(c, out)
437    }
438
439    fn forget(&mut self, c: Correlation) {
440        self.handles.caller.forget(c);
441    }
442}
443
444impl Handler for Loopback {
445    fn serve(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), ServeError> {
446        self.handles.handler.serve(iface, ords)
447    }
448
449    fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
450        self.handles.handler.next_claim(out)
451    }
452
453    fn settle(
454        &mut self,
455        claim: ClaimId,
456        outcome: Result<&[u8], CallError>,
457    ) -> Result<(), SettleError> {
458        self.handles.handler.settle(claim, outcome)
459    }
460}
461
462/// Each key goes to the handle that observes it: `Outcome` and `Slot` to the
463/// caller, `Event` to the source, and `Claim` to the handler.
464impl Wakeable for Loopback {
465    fn wake_on(&self, what: Interest, waker: &Waker) {
466        match what {
467            Interest::Outcome(_) | Interest::Slot => self.handles.caller.wake_on(what, waker),
468            Interest::Event(_) => self.handles.source.wake_on(what, waker),
469            Interest::Claim(_) => self.handles.handler.wake_on(what, waker),
470        }
471    }
472}
473
474// ---------------------------------------------------------------------------
475// The threading model, asserted rather than documented (ADR-0021 decision 12).
476// ---------------------------------------------------------------------------
477
478const _: () = {
479    const fn assert_sync<T: Sync>() {}
480    const fn assert_send<T: Send>() {}
481
482    // Several threads may read one store at once.
483    assert_sync::<ReaderHandle>();
484    assert_send::<ReaderHandle>();
485
486    // One thread drives each of the rest.
487    assert_send::<WriterHandle>();
488    assert_send::<SourceHandle>();
489    assert_send::<SinkHandle>();
490    assert_send::<CallerHandle>();
491    assert_send::<HandlerHandle>();
492
493    // The aggregate is as `Send` as the handles it holds.
494    assert_send::<Loopback>();
495};