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};