Skip to main content

ridl_loopback/
handle.rs

1//! The six handle types and their port implementations.
2//!
3//! A **port role** is one port trait, and a runtime presents one handle type
4//! per role rather than one type implementing them all (ADR-0021 decision 12).
5//! The six here group the eleven roles the way that decision derives the
6//! threading split: the five roles with a `&mut self` method take a handle
7//! each, and the six whose methods all take `&self` share one, because they
8//! are exactly the roles several threads may hold at once. Every handle here
9//! holds the same `Arc<Mutex<Store>>` and its own copy of the
10//! [`CatalogRef`](ridl_rt::contract::CatalogRef) the runtime was built with,
11//! so `Attached::catalog` can return a reference without reaching through the
12//! lock.
13//!
14//! Every handle implements `Attached`, which every port trait but `Clock` has
15//! as a supertrait. The table lists what each handle adds to it, and the split
16//! follows the receiver of those methods, as ADR-0021 decision 12 derives it:
17//!
18//! | Handle            | Port roles beside `Attached`                                                  | Threading     |
19//! | ----------------- | ----------------------------------------------------------------------------- | ------------- |
20//! | [`ReaderHandle`]  | `Clock`, `SignalReader`, `FixedReader`, `ScannableSignals`, `CoherentSignals` | `Send + Sync` |
21//! | [`WriterHandle`]  | `SignalWriter`                                                                | `Send`        |
22//! | [`SourceHandle`]  | `EventSource`                                                                 | `Send`        |
23//! | [`SinkHandle`]    | `EventSink`                                                                   | `Send`        |
24//! | [`CallerHandle`]  | `Caller`                                                                      | `Send`        |
25//! | [`HandlerHandle`] | `Handler`                                                                     | `Send`        |
26//!
27//! Every method on the reader handle takes `&self`, so several threads may
28//! read one store at once; every other handle carries a trait with a
29//! `&mut self` method and is driven by one thread at a time. Neither property
30//! is declared: both are derived by the compiler from the fields, and
31//! `crates/ridl-loopback/src/lib.rs` asserts them at compile time.
32
33use std::collections::BTreeMap;
34use std::sync::{Arc, Mutex, MutexGuard};
35
36use ridl_rt::contract::{CatalogRef, InterfaceNo, Ordinal};
37use ridl_rt::error::CallError;
38use ridl_rt::port::{
39    Attached, Caller, Changed, Claim, ClaimId, Clock, CoherentSignals, Correlation, EventSink,
40    EventSource, FixedReader, Handler, RaiseError, RawOccurrence, RawSample, ReadError,
41    ScannableSignals, SendError, ServeError, SettleError, SignalReader, SignalWriter,
42    SubscribeError, Watermark, WriteError,
43};
44use ridl_rt::sample::Timestamp;
45
46use crate::store::{CallKind, Key, Staged, Store};
47
48/// The shared store, as every handle holds it.
49pub(crate) type Shared = Arc<Mutex<Store>>;
50
51/// Takes the one lock.
52///
53/// A panic while the lock is held poisons it. This runtime recovers the guard
54/// rather than propagating the poison, because every critical section here is
55/// a read or a write of the maps that leaves them well formed, and a poisoned
56/// lock would otherwise turn one panic in one test into a panic in every later
57/// port call over the same runtime.
58pub(crate) fn lock(shared: &Shared) -> MutexGuard<'_, Store> {
59    shared
60        .lock()
61        .unwrap_or_else(|poisoned| poisoned.into_inner())
62}
63
64// ---------------------------------------------------------------------------
65// The reader handle
66// ---------------------------------------------------------------------------
67
68/// The `&self` port roles: `Attached`, `Clock`, `SignalReader`,
69/// `FixedReader`, and the two signal extensions.
70///
71/// It is `Send + Sync`, so several threads may hold one and read at once, and
72/// a face that needs only `SignalReader` can be built over this handle
73/// directly rather than over the aggregate.
74pub struct ReaderHandle {
75    shared: Shared,
76    catalog: CatalogRef,
77}
78
79impl ReaderHandle {
80    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
81        ReaderHandle { shared, catalog }
82    }
83}
84
85impl Attached for ReaderHandle {
86    fn catalog(&self) -> &CatalogRef {
87        &self.catalog
88    }
89}
90
91impl Clock for ReaderHandle {
92    fn now(&self) -> Timestamp {
93        lock(&self.shared).now()
94    }
95}
96
97impl SignalReader for ReaderHandle {
98    fn read(
99        &self,
100        iface: InterfaceNo,
101        ord: Ordinal,
102        out: &mut [u8],
103    ) -> Result<RawSample, ReadError> {
104        lock(&self.shared).read(iface, ord, out)
105    }
106}
107
108impl FixedReader for ReaderHandle {
109    fn read_fixed(
110        &self,
111        iface: InterfaceNo,
112        ord: Ordinal,
113        out: &mut [u8],
114    ) -> Result<usize, ReadError> {
115        lock(&self.shared).read_fixed(iface, ord, out)
116    }
117}
118
119impl ScannableSignals for ReaderHandle {
120    fn generation(&self, iface: InterfaceNo) -> u64 {
121        lock(&self.shared).generation(iface)
122    }
123
124    fn scan(&self, marks: &mut [Watermark], out: &mut [Changed]) -> usize {
125        lock(&self.shared).scan(marks, out)
126    }
127}
128
129impl CoherentSignals for ReaderHandle {
130    fn read_coherent(
131        &self,
132        iface: InterfaceNo,
133        ords: &[Ordinal],
134        out: &mut [u8],
135        samples: &mut [RawSample],
136    ) -> Result<usize, ReadError> {
137        lock(&self.shared).read_coherent(iface, ords, out, samples)
138    }
139}
140
141// ---------------------------------------------------------------------------
142// The writer handle
143// ---------------------------------------------------------------------------
144
145/// The `SignalWriter` port role.
146///
147/// The staged changes and the per channel sequence counters live on the
148/// handle, not in the store: staging is one provider's private state until its
149/// `commit`, and a sequence number is assigned by the sender (ridl §3.1).
150pub struct WriterHandle {
151    shared: Shared,
152    catalog: CatalogRef,
153    staged: BTreeMap<Key, Staged>,
154    seqs: BTreeMap<Key, u64>,
155}
156
157impl WriterHandle {
158    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
159        WriterHandle {
160            shared,
161            catalog,
162            staged: BTreeMap::new(),
163            seqs: BTreeMap::new(),
164        }
165    }
166}
167
168impl Attached for WriterHandle {
169    fn catalog(&self) -> &CatalogRef {
170        &self.catalog
171    }
172}
173
174impl SignalWriter for WriterHandle {
175    fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
176        self.staged
177            .insert((iface, ord), Staged::Set(bytes.to_vec()));
178        Ok(())
179    }
180
181    fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
182        self.staged.insert((iface, ord), Staged::Invalidate);
183        Ok(())
184    }
185
186    fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
187        // A touch re-affirms the current value. It stages one only when
188        // nothing else is staged for the channel: a `set` or an `invalidate`
189        // already staged is itself a publication, and a re-affirmation adds
190        // nothing to it. Replacing one here would discard a value this writer
191        // staged, which is not what `touch` means. A later `set` or
192        // `invalidate` does replace a staged touch, because each is a newer
193        // decision about the same channel.
194        self.staged.entry((iface, ord)).or_insert(Staged::Touch);
195        Ok(())
196    }
197
198    fn commit(&mut self) {
199        lock(&self.shared).commit(&mut self.staged, &mut self.seqs);
200    }
201}
202
203// ---------------------------------------------------------------------------
204// The event handles
205// ---------------------------------------------------------------------------
206
207/// The `EventSource` port role: one subscription set and one queue, both in
208/// the store under this handle's identity.
209///
210/// Each source handle receives its own copy of every occurrence raised while
211/// it is subscribed, so two consumers of one event never consume each other's
212/// occurrences.
213pub struct SourceHandle {
214    shared: Shared,
215    catalog: CatalogRef,
216    id: usize,
217}
218
219impl SourceHandle {
220    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
221        let id = lock(&shared).open_source();
222        SourceHandle {
223            shared,
224            catalog,
225            id,
226        }
227    }
228}
229
230impl Drop for SourceHandle {
231    fn drop(&mut self) {
232        lock(&self.shared).close_source(self.id);
233    }
234}
235
236impl Attached for SourceHandle {
237    fn catalog(&self) -> &CatalogRef {
238        &self.catalog
239    }
240}
241
242impl EventSource for SourceHandle {
243    fn subscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), SubscribeError> {
244        lock(&self.shared).subscribe(self.id, iface, ords);
245        Ok(())
246    }
247
248    fn unsubscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) {
249        lock(&self.shared).unsubscribe(self.id, iface, ords);
250    }
251
252    fn next(&mut self, out: &mut [u8]) -> Result<Option<RawOccurrence>, ReadError> {
253        lock(&self.shared).next_event(self.id, out)
254    }
255}
256
257/// The `EventSink` port role.
258///
259/// The sequence counters are the handle's own, one per event channel, for the
260/// same reason the writer handle's are: ridl §3.1 scopes the number to the
261/// channel and has the sender assign it. One counter for the whole handle
262/// would make a consumer subscribed to some of this sink's events see a gap in
263/// `seq` where nothing was lost, and `EventSource::next` states that a gap is
264/// a loss.
265pub struct SinkHandle {
266    shared: Shared,
267    catalog: CatalogRef,
268    seqs: BTreeMap<Key, u64>,
269}
270
271impl SinkHandle {
272    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
273        SinkHandle {
274            shared,
275            catalog,
276            seqs: BTreeMap::new(),
277        }
278    }
279
280    fn take_seq(&mut self, key: Key) -> u64 {
281        let counter = self.seqs.entry(key).or_insert(0);
282        *counter += 1;
283        *counter
284    }
285}
286
287impl Attached for SinkHandle {
288    fn catalog(&self) -> &CatalogRef {
289        &self.catalog
290    }
291}
292
293impl EventSink for SinkHandle {
294    fn raise(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), RaiseError> {
295        let seq = self.take_seq((iface, ord));
296        lock(&self.shared).raise(iface, ord, bytes, seq);
297        Ok(())
298    }
299}
300
301// ---------------------------------------------------------------------------
302// The call handles
303// ---------------------------------------------------------------------------
304
305/// The `Caller` port role.
306///
307/// One counter for the whole handle, not one per channel: on a call the scope
308/// is the caller instance (ADR-0021 decision 5, and driftsys/ridl#308's own
309/// report), so every call this caller sends draws from one sequence. That is
310/// what keeps two callers on one provider from colliding.
311pub struct CallerHandle {
312    shared: Shared,
313    catalog: CatalogRef,
314    next_seq: u64,
315}
316
317impl CallerHandle {
318    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
319        CallerHandle {
320            shared,
321            catalog,
322            next_seq: 0,
323        }
324    }
325
326    fn take_seq(&mut self) -> u64 {
327        self.next_seq += 1;
328        self.next_seq
329    }
330}
331
332impl Attached for CallerHandle {
333    fn catalog(&self) -> &CatalogRef {
334        &self.catalog
335    }
336}
337
338impl Caller for CallerHandle {
339    fn command(
340        &mut self,
341        iface: InterfaceNo,
342        ord: Ordinal,
343        args: &[u8],
344    ) -> Result<Correlation, SendError> {
345        let seq = self.take_seq();
346        Ok(lock(&self.shared).send(CallKind::Command, iface, ord, args, seq))
347    }
348
349    fn query(
350        &mut self,
351        iface: InterfaceNo,
352        ord: Ordinal,
353        args: &[u8],
354    ) -> Result<Correlation, SendError> {
355        let seq = self.take_seq();
356        Ok(lock(&self.shared).send(CallKind::Query, iface, ord, args, seq))
357    }
358
359    fn ack(&mut self, c: Correlation) -> Option<Result<(), CallError>> {
360        lock(&self.shared).ack(c)
361    }
362
363    fn reply(
364        &mut self,
365        c: Correlation,
366        out: &mut [u8],
367    ) -> Result<Option<Result<usize, CallError>>, ReadError> {
368        lock(&self.shared).reply(c, out)
369    }
370
371    fn forget(&mut self, c: Correlation) {
372        lock(&self.shared).forget(c);
373    }
374}
375
376/// The `Handler` port role.
377///
378/// A handler that has served nothing is presented every call waiting in the
379/// store; once it has served anything, it is presented only the members it
380/// served, and another handler's calls stay waiting for that handler. A claim
381/// belongs to the handler it was presented to, so another handler's `settle`
382/// of it answers [`SettleError::UnknownClaim`].
383///
384/// The empty set meaning no filter is a deliberate deviation from
385/// [`Handler::serve`], which says delivery starts at the members listed. The
386/// generated `dispatch` never calls `serve`
387/// (`crates/ridl-backend-rust/src/face.rs`), so a handler that always filtered
388/// would be presented nothing at all by it. `serve` with an empty slice
389/// records nothing and so leaves the handler unfiltered, the same as never
390/// having called it. [`served`](HandlerHandle::served) reads the set back.
391pub struct HandlerHandle {
392    shared: Shared,
393    catalog: CatalogRef,
394    id: usize,
395    served: Vec<Key>,
396}
397
398impl HandlerHandle {
399    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
400        let id = lock(&shared).open_handler();
401        HandlerHandle {
402            shared,
403            catalog,
404            id,
405            served: Vec::new(),
406        }
407    }
408
409    /// The members `serve` was called with, in the order they were served,
410    /// with no duplicate.
411    #[must_use]
412    pub fn served(&self) -> &[(InterfaceNo, Ordinal)] {
413        &self.served
414    }
415}
416
417impl Attached for HandlerHandle {
418    fn catalog(&self) -> &CatalogRef {
419        &self.catalog
420    }
421}
422
423impl Handler for HandlerHandle {
424    fn serve(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), ServeError> {
425        for ord in ords {
426            if !self.served.contains(&(iface, *ord)) {
427                self.served.push((iface, *ord));
428            }
429        }
430        Ok(())
431    }
432
433    fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
434        let served = if self.served.is_empty() {
435            None
436        } else {
437            Some(self.served.as_slice())
438        };
439        lock(&self.shared).next_claim(self.id, served, out)
440    }
441
442    fn settle(
443        &mut self,
444        claim: ClaimId,
445        outcome: Result<&[u8], CallError>,
446    ) -> Result<(), SettleError> {
447        lock(&self.shared).settle(self.id, claim, outcome)
448    }
449}