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