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