Skip to main content

ridl_loopback/
lib.rs

1//! `ridl-loopback` — the in-process reference runtime.
2//!
3//! This crate implements the eleven 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 the 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.
30//! [`Loopback`] is the aggregate: it implements all eleven port traits by delegating to the six
31//! role handles it holds, and it is what a generated `Client`, `Publisher` or
32//! `dispatch` is normally built over.
33//!
34//! ```
35//! use ridl_loopback::Loopback;
36//! use ridl_rt::contract::{CatalogHash, CatalogRef, InterfaceNo, Ordinal};
37//! use ridl_rt::port::{SignalReader, SignalWriter};
38//!
39//! let catalog = CatalogRef { name: "face.demo", hash: CatalogHash([0u8; 32]) };
40//! let mut rt = Loopback::new(catalog);
41//!
42//! rt.set(InterfaceNo(1), Ordinal(1), &[42]).expect("staged");
43//! rt.commit();
44//!
45//! let mut out = [0u8; 8];
46//! let sample = rt.read(InterfaceNo(1), Ordinal(1), &mut out).expect("read");
47//! assert_eq!(&out[..sample.len], &[42]);
48//! ```
49//!
50//! [`Loopback::split`] hands out the six role handles for a program that wants
51//! them apart — one thread reading while another publishes, or two callers on
52//! one provider. [`ReaderHandle`] is `Send + Sync`; the other five are `Send`
53//! and driven by one thread each. Nothing here declares either: both follow
54//! from the fields, and the assertions at the bottom of this file pin them.
55//!
56//! # What it reports, and what it cannot
57//!
58//! The loopback holds no catalog descriptor — the descriptor and the catalog
59//! hash arrive with story E16.2 (driftsys/ridl#378) — so it has no member
60//! table, and there is no ordinal it can call unknown, no member it can call
61//! unowned, and no timing annotation it can measure a value's freshness or a
62//! call's remaining time against. What it therefore never returns:
63//! `WriteError::NotOwner`, `RaiseError::NotOwner`, `ServeError::NotOwner`, any
64//! port error's `Contract` variant except
65//! [`FixedReader::read_fixed`](ridl_rt::port::FixedReader::read_fixed)'s, and
66//! `Freshness::Fresh` or `Freshness::Stale`. Nothing detaches, because every
67//! handle holds the store alive, so `Detached` never appears either; and
68//! nothing is bounded, so `Busy` and `TooLarge` do not appear outside
69//! [`Loopback::fail_next_settle`].
70//!
71//! [`Attached::catalog`](ridl_rt::port::Attached::catalog) returns the
72//! `CatalogRef` the runtime was built with, unexamined. ADR-0021 decision 3
73//! places a check of it against the interface's own `CATALOG` in a generated
74//! face's constructor, once, when the face is built; the constructor the Rust
75//! backend emits today performs no such check (driftsys/ridl#448). Either way
76//! it is the face's check and not the runtime's: the loopback carries the
77//! value and compares nothing.
78//!
79//! The crate's as-built design record, with the reasoning behind each of these
80//! choices, is `docs/design/ridl-loopback.md` in this repository.
81
82use std::sync::{Arc, Mutex};
83
84use ridl_rt::contract::{CatalogRef, InterfaceNo, Ordinal};
85use ridl_rt::error::CallError;
86use ridl_rt::port::{
87    Attached, Caller, Changed, Claim, ClaimId, Clock, CoherentSignals, Correlation, EventSink,
88    EventSource, FixedReader, Handler, RaiseError, RawOccurrence, RawSample, ReadError,
89    ScannableSignals, SendError, ServeError, SettleError, SignalReader, SignalWriter,
90    SubscribeError, Watermark, WriteError,
91};
92use ridl_rt::sample::{Duration, Timestamp};
93
94mod handle;
95mod store;
96
97pub use handle::{
98    CallerHandle, HandlerHandle, ReaderHandle, SinkHandle, SourceHandle, WriterHandle,
99};
100
101use handle::{Shared, lock};
102use store::Store;
103
104/// The six role handles of one runtime, as [`Loopback::split`] hands them out.
105pub struct Handles {
106    /// `Attached`, `Clock`, `SignalReader`, `FixedReader`, `ScannableSignals`
107    /// and `CoherentSignals`.
108    pub reader: ReaderHandle,
109    /// `SignalWriter`.
110    pub writer: WriterHandle,
111    /// `EventSource`.
112    pub source: SourceHandle,
113    /// `EventSink`.
114    pub sink: SinkHandle,
115    /// `Caller`.
116    pub caller: CallerHandle,
117    /// `Handler`.
118    pub handler: HandlerHandle,
119}
120
121/// The aggregate handle: one value implementing all eleven port traits by
122/// delegating to the six role handles it holds.
123///
124/// A face is built over one value implementing at least the port traits its
125/// interface needs, and a generated `Client` is commonly bound over
126/// `SignalReader + EventSource + Caller` at once, which no single role handle
127/// satisfies. That is what an aggregate is for (ADR-0021 decision 12). Pass it
128/// by value, or as `&mut` under the forwarding impls of ADR-0021 decision 11.
129pub struct Loopback {
130    shared: Shared,
131    catalog: CatalogRef,
132    handles: Handles,
133}
134
135impl Loopback {
136    /// A runtime attached to `catalog`, with an empty store and its clock at
137    /// [`Timestamp`] 0.
138    ///
139    /// The catalog is carried, not checked: see the crate documentation.
140    #[must_use]
141    pub fn new(catalog: CatalogRef) -> Self {
142        let shared: Shared = Arc::new(Mutex::new(Store::new()));
143        let handles = Handles {
144            reader: ReaderHandle::new(Arc::clone(&shared), catalog),
145            writer: WriterHandle::new(Arc::clone(&shared), catalog),
146            source: SourceHandle::new(Arc::clone(&shared), catalog),
147            sink: SinkHandle::new(Arc::clone(&shared), catalog),
148            caller: CallerHandle::new(Arc::clone(&shared), catalog),
149            handler: HandlerHandle::new(Arc::clone(&shared), catalog),
150        };
151        Loopback {
152            shared,
153            catalog,
154            handles,
155        }
156    }
157
158    /// Hands out the six role handles this aggregate holds. Every one of them
159    /// keeps the same store, so a value published through `writer` is read
160    /// through `reader`.
161    #[must_use]
162    pub fn split(self) -> Handles {
163        self.handles
164    }
165
166    /// An additional reader handle on the same store.
167    #[must_use]
168    pub fn reader(&self) -> ReaderHandle {
169        ReaderHandle::new(Arc::clone(&self.shared), self.catalog)
170    }
171
172    /// An additional writer handle on the same store, with its own staging
173    /// area and its own per channel sequence counters.
174    #[must_use]
175    pub fn writer(&self) -> WriterHandle {
176        WriterHandle::new(Arc::clone(&self.shared), self.catalog)
177    }
178
179    /// An additional event source on the same store, with its own
180    /// subscription set and its own queue.
181    #[must_use]
182    pub fn source(&self) -> SourceHandle {
183        SourceHandle::new(Arc::clone(&self.shared), self.catalog)
184    }
185
186    /// An additional event sink on the same store, with its own sequence
187    /// counter.
188    #[must_use]
189    pub fn sink(&self) -> SinkHandle {
190        SinkHandle::new(Arc::clone(&self.shared), self.catalog)
191    }
192
193    /// An additional caller on the same store, with its own sequence counter.
194    /// Two callers on one provider therefore never collide on a sequence
195    /// number (driftsys/ridl#308).
196    #[must_use]
197    pub fn caller(&self) -> CallerHandle {
198        CallerHandle::new(Arc::clone(&self.shared), self.catalog)
199    }
200
201    /// An additional handler on the same store.
202    ///
203    /// Every handler draws from the one queue of waiting calls, filtered by
204    /// what it has served: a handler that has served nothing is presented
205    /// every waiting call, and one that has served members is presented only
206    /// those. A claim belongs to the handler it was presented to.
207    #[must_use]
208    pub fn handler(&self) -> HandlerHandle {
209        HandlerHandle::new(Arc::clone(&self.shared), self.catalog)
210    }
211
212    /// Advances the clock by `by`.
213    ///
214    /// The clock is a counter this method moves and nothing else moves. It
215    /// never reads wall-clock time, so a round trip over this runtime produces
216    /// the same timestamps on every run and on every machine. It saturates
217    /// rather than overflowing.
218    ///
219    /// # Panics
220    ///
221    /// When `by` is negative. A clock that ran backwards would stamp an
222    /// envelope before one already stamped.
223    pub fn advance(&mut self, by: Duration) {
224        lock(&self.shared).advance(by);
225    }
226
227    /// Supplies the value of a `fixed` (ridl §8), which is provisioned into a
228    /// runtime rather than published by application code.
229    ///
230    /// Until a `fixed` is provisioned, reading it answers
231    /// `ReadError::Contract(Contract::UnknownInteraction)`: the loopback holds
232    /// no value and, having no member table, cannot say anything narrower.
233    pub fn provision_fixed(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) {
234        lock(&self.shared).provision_fixed(iface, ord, bytes);
235    }
236
237    /// Makes the next `Handler::settle` on this runtime fail with
238    /// `SettleError::TooLarge` and record no outcome.
239    ///
240    /// This is the one fault this runtime injects, and the one place a
241    /// `SettleError` other than `UnknownClaim` comes from. It exists because
242    /// the generated `dispatch` counts a claim only once the handler has
243    /// accepted its settlement, and nothing else in an in-process runtime can
244    /// make that path fail.
245    pub fn fail_next_settle(&mut self) {
246        lock(&self.shared).fail_next_settle();
247    }
248}
249
250// ---------------------------------------------------------------------------
251// The aggregate's eleven port implementations, each one a delegation.
252// ---------------------------------------------------------------------------
253
254impl Attached for Loopback {
255    fn catalog(&self) -> &CatalogRef {
256        self.handles.reader.catalog()
257    }
258}
259
260impl Clock for Loopback {
261    fn now(&self) -> Timestamp {
262        self.handles.reader.now()
263    }
264}
265
266impl SignalReader for Loopback {
267    fn read(
268        &self,
269        iface: InterfaceNo,
270        ord: Ordinal,
271        out: &mut [u8],
272    ) -> Result<RawSample, ReadError> {
273        self.handles.reader.read(iface, ord, out)
274    }
275}
276
277impl FixedReader for Loopback {
278    fn read_fixed(
279        &self,
280        iface: InterfaceNo,
281        ord: Ordinal,
282        out: &mut [u8],
283    ) -> Result<usize, ReadError> {
284        self.handles.reader.read_fixed(iface, ord, out)
285    }
286}
287
288impl ScannableSignals for Loopback {
289    fn generation(&self, iface: InterfaceNo) -> u64 {
290        self.handles.reader.generation(iface)
291    }
292
293    fn scan(&self, marks: &mut [Watermark], out: &mut [Changed]) -> usize {
294        self.handles.reader.scan(marks, out)
295    }
296}
297
298impl CoherentSignals for Loopback {
299    fn read_coherent(
300        &self,
301        iface: InterfaceNo,
302        ords: &[Ordinal],
303        out: &mut [u8],
304        samples: &mut [RawSample],
305    ) -> Result<usize, ReadError> {
306        self.handles.reader.read_coherent(iface, ords, out, samples)
307    }
308}
309
310impl SignalWriter for Loopback {
311    fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
312        self.handles.writer.set(iface, ord, bytes)
313    }
314
315    fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
316        self.handles.writer.invalidate(iface, ord)
317    }
318
319    fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
320        self.handles.writer.touch(iface, ord)
321    }
322
323    fn commit(&mut self) {
324        self.handles.writer.commit();
325    }
326}
327
328impl EventSource for Loopback {
329    fn subscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), SubscribeError> {
330        self.handles.source.subscribe(iface, ords)
331    }
332
333    fn unsubscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) {
334        self.handles.source.unsubscribe(iface, ords);
335    }
336
337    fn next(&mut self, out: &mut [u8]) -> Result<Option<RawOccurrence>, ReadError> {
338        self.handles.source.next(out)
339    }
340}
341
342impl EventSink for Loopback {
343    fn raise(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), RaiseError> {
344        self.handles.sink.raise(iface, ord, bytes)
345    }
346}
347
348impl Caller for Loopback {
349    fn command(
350        &mut self,
351        iface: InterfaceNo,
352        ord: Ordinal,
353        args: &[u8],
354    ) -> Result<Correlation, SendError> {
355        self.handles.caller.command(iface, ord, args)
356    }
357
358    fn query(
359        &mut self,
360        iface: InterfaceNo,
361        ord: Ordinal,
362        args: &[u8],
363    ) -> Result<Correlation, SendError> {
364        self.handles.caller.query(iface, ord, args)
365    }
366
367    fn ack(&mut self, c: Correlation) -> Option<Result<(), CallError>> {
368        self.handles.caller.ack(c)
369    }
370
371    fn reply(
372        &mut self,
373        c: Correlation,
374        out: &mut [u8],
375    ) -> Result<Option<Result<usize, CallError>>, ReadError> {
376        self.handles.caller.reply(c, out)
377    }
378
379    fn forget(&mut self, c: Correlation) {
380        self.handles.caller.forget(c);
381    }
382}
383
384impl Handler for Loopback {
385    fn serve(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), ServeError> {
386        self.handles.handler.serve(iface, ords)
387    }
388
389    fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
390        self.handles.handler.next_claim(out)
391    }
392
393    fn settle(
394        &mut self,
395        claim: ClaimId,
396        outcome: Result<&[u8], CallError>,
397    ) -> Result<(), SettleError> {
398        self.handles.handler.settle(claim, outcome)
399    }
400}
401
402// ---------------------------------------------------------------------------
403// The threading model, asserted rather than documented (ADR-0021 decision 12).
404// ---------------------------------------------------------------------------
405
406const _: () = {
407    const fn assert_sync<T: Sync>() {}
408    const fn assert_send<T: Send>() {}
409
410    // Several threads may read one store at once.
411    assert_sync::<ReaderHandle>();
412    assert_send::<ReaderHandle>();
413
414    // One thread drives each of the rest.
415    assert_send::<WriterHandle>();
416    assert_send::<SourceHandle>();
417    assert_send::<SinkHandle>();
418    assert_send::<CallerHandle>();
419    assert_send::<HandlerHandle>();
420
421    // The aggregate is as `Send` as the handles it holds.
422    assert_send::<Loopback>();
423};