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 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. The twelfth port
9//! trait, `Wakeable`, is on every handle, because each handle wakes its own
10//! waiters. Every handle here
11//! holds the same `Arc<Mutex<Store>>` and its own copy of the
12//! [`CatalogRef`](ridl_rt::contract::CatalogRef) the runtime was built with,
13//! so `Attached::catalog` can return a reference without reaching through the
14//! lock.
15//!
16//! Every handle implements `Attached`, which every port trait but `Clock` and
17//! `Wakeable` has as a supertrait, and `Wakeable`. The table lists what each
18//! handle adds to them, and the split follows the receiver of those methods,
19//! as ADR-0021 decision 12 derives it:
20//!
21//! | Handle            | Port roles beside `Attached` and `Wakeable`                                   | Kinds of key it stores                            | Threading     |
22//! | ----------------- | ----------------------------------------------------------------------------- | ------------------------------------------------- | ------------- |
23//! | [`ReaderHandle`]  | `Clock`, `SignalReader`, `FixedReader`, `ScannableSignals`, `CoherentSignals` | none                                              | `Send + Sync` |
24//! | [`WriterHandle`]  | `SignalWriter`                                                                | none                                              | `Send`        |
25//! | [`SourceHandle`]  | `EventSource`                                                                 | `Event`, one waker                                | `Send`        |
26//! | [`SinkHandle`]    | `EventSink`                                                                   | none                                              | `Send`        |
27//! | [`CallerHandle`]  | `Clock`, `Caller`                                                             | `Outcome`, kept with each call; `Slot`, one waker | `Send`        |
28//! | [`HandlerHandle`] | `Handler`                                                                     | `Claim`, one waker                                | `Send`        |
29//!
30//! A handle stores a waker only under a kind of key one of its roles
31//! observes. For `Event` and `Claim` that is one waker per kind, and a
32//! change to any key of that kind wakes it (ADR-0021 decision 13). An
33//! `Outcome` waker is per call: it is kept with its call, and no change but
34//! that call's settlement, its `forget`, or the drop of the caller handle
35//! that sent it wakes it; a displacement by another task does, as for every
36//! kind. A caller stores one `Slot` waker, woken at once while a slot of the
37//! call table is free, and otherwise by the next reclaim of a slot — a
38//! `forget`, a settlement of a forgotten call, a caller handle's drop, or a
39//! handler handle's drop while it holds a claim whose call was forgotten —
40//! which wakes every other caller's `Slot` waker. A registration
41//! under any other kind is woken at once, because nothing that handle could
42//! read changes under it, and a stored waker would never be woken.
43//!
44//! Every method on the reader handle takes `&self`, so several threads may
45//! read one store at once; every other handle carries a trait with a
46//! `&mut self` method and is driven by one thread at a time. Neither property
47//! is declared: both are derived by the compiler from the fields, and
48//! `crates/ridl-loopback/src/lib.rs` asserts them at compile time.
49
50use std::collections::BTreeMap;
51use std::sync::{Arc, Mutex, MutexGuard};
52use std::task::Waker;
53
54use ridl_rt::contract::{CatalogRef, InterfaceNo, Ordinal};
55use ridl_rt::error::CallError;
56use ridl_rt::port::{
57    Attached, Caller, Changed, Claim, ClaimId, Clock, CoherentSignals, Correlation, EventSink,
58    EventSource, FixedReader, Handler, Interest, RaiseError, RawOccurrence, RawSample, ReadError,
59    ScannableSignals, SendError, ServeError, SettleError, SignalReader, SignalWriter,
60    SubscribeError, Wakeable, Watermark, WriteError,
61};
62use ridl_rt::sample::Timestamp;
63
64use crate::store::{CallKind, Key, Staged, Store};
65
66/// The shared store, as every handle holds it.
67pub(crate) type Shared = Arc<Mutex<Store>>;
68
69/// Takes the one lock.
70///
71/// A panic while the lock is held poisons it. This runtime recovers the guard
72/// rather than propagating the poison, because every critical section here is
73/// a read or a write of the maps that leaves them well formed, and a poisoned
74/// lock would otherwise turn one panic in one test into a panic in every later
75/// port call over the same runtime.
76pub(crate) fn lock(shared: &Shared) -> MutexGuard<'_, Store> {
77    shared
78        .lock()
79        .unwrap_or_else(|poisoned| poisoned.into_inner())
80}
81
82/// Takes the one lock, runs `f` with a list of wakers to wake, releases the
83/// lock, and then wakes each waker `f` put on the list.
84///
85/// A waker runs code the runtime does not control — an executor's scheduling,
86/// or a test's own — and that code may call a port method on this runtime. So
87/// no waker is woken while the lock is held.
88pub(crate) fn locked<R>(shared: &Shared, f: impl FnOnce(&mut Store, &mut Vec<Waker>) -> R) -> R {
89    let mut wake = Vec::new();
90    let result = {
91        let mut store = lock(shared);
92        f(&mut store, &mut wake)
93    };
94    for waker in wake {
95        waker.wake();
96    }
97    result
98}
99
100/// `Wakeable` on a handle none of whose roles observes any key: every
101/// registration is woken at once.
102fn wake_at_once(waker: &Waker) {
103    waker.wake_by_ref();
104}
105
106// ---------------------------------------------------------------------------
107// The reader handle
108// ---------------------------------------------------------------------------
109
110/// The `&self` port roles: `Attached`, `Clock`, `SignalReader`,
111/// `FixedReader`, and the two signal extensions.
112///
113/// It is `Send + Sync`, so several threads may hold one and read at once, and
114/// a face that needs only `SignalReader` can be built over this handle
115/// directly rather than over the aggregate.
116pub struct ReaderHandle {
117    shared: Shared,
118    catalog: CatalogRef,
119}
120
121impl ReaderHandle {
122    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
123        ReaderHandle { shared, catalog }
124    }
125}
126
127impl Attached for ReaderHandle {
128    fn catalog(&self) -> &CatalogRef {
129        &self.catalog
130    }
131}
132
133impl Clock for ReaderHandle {
134    fn now(&self) -> Timestamp {
135        lock(&self.shared).now()
136    }
137}
138
139impl Wakeable for ReaderHandle {
140    fn wake_on(&self, _: Interest, waker: &Waker) {
141        wake_at_once(waker);
142    }
143}
144
145impl SignalReader for ReaderHandle {
146    fn read(
147        &self,
148        iface: InterfaceNo,
149        ord: Ordinal,
150        out: &mut [u8],
151    ) -> Result<RawSample, ReadError> {
152        lock(&self.shared).read(iface, ord, out)
153    }
154}
155
156impl FixedReader for ReaderHandle {
157    fn read_fixed(
158        &self,
159        iface: InterfaceNo,
160        ord: Ordinal,
161        out: &mut [u8],
162    ) -> Result<usize, ReadError> {
163        lock(&self.shared).read_fixed(iface, ord, out)
164    }
165}
166
167impl ScannableSignals for ReaderHandle {
168    fn generation(&self, iface: InterfaceNo) -> u64 {
169        lock(&self.shared).generation(iface)
170    }
171
172    fn scan(&self, marks: &mut [Watermark], out: &mut [Changed]) -> usize {
173        lock(&self.shared).scan(marks, out)
174    }
175}
176
177impl CoherentSignals for ReaderHandle {
178    fn read_coherent(
179        &self,
180        iface: InterfaceNo,
181        ords: &[Ordinal],
182        out: &mut [u8],
183        samples: &mut [RawSample],
184    ) -> Result<usize, ReadError> {
185        lock(&self.shared).read_coherent(iface, ords, out, samples)
186    }
187}
188
189// ---------------------------------------------------------------------------
190// The writer handle
191// ---------------------------------------------------------------------------
192
193/// The `SignalWriter` port role.
194///
195/// The staged changes and the per channel sequence counters live on the
196/// handle, not in the store: staging is one provider's private state until its
197/// `commit`, and a sequence number is assigned by the sender (ridl §3.1).
198pub struct WriterHandle {
199    shared: Shared,
200    catalog: CatalogRef,
201    staged: BTreeMap<Key, Staged>,
202    seqs: BTreeMap<Key, u64>,
203}
204
205impl WriterHandle {
206    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
207        WriterHandle {
208            shared,
209            catalog,
210            staged: BTreeMap::new(),
211            seqs: BTreeMap::new(),
212        }
213    }
214}
215
216impl Attached for WriterHandle {
217    fn catalog(&self) -> &CatalogRef {
218        &self.catalog
219    }
220}
221
222impl SignalWriter for WriterHandle {
223    fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
224        self.staged
225            .insert((iface, ord), Staged::Set(bytes.to_vec()));
226        Ok(())
227    }
228
229    fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
230        self.staged.insert((iface, ord), Staged::Invalidate);
231        Ok(())
232    }
233
234    fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
235        // A touch re-affirms the current value. It stages one only when
236        // nothing else is staged for the channel: a `set` or an `invalidate`
237        // already staged is itself a publication, and a re-affirmation adds
238        // nothing to it. Replacing one here would discard a value this writer
239        // staged, which is not what `touch` means. A later `set` or
240        // `invalidate` does replace a staged touch, because each is a newer
241        // decision about the same channel.
242        self.staged.entry((iface, ord)).or_insert(Staged::Touch);
243        Ok(())
244    }
245
246    fn commit(&mut self) {
247        lock(&self.shared).commit(&mut self.staged, &mut self.seqs);
248    }
249}
250
251impl Wakeable for WriterHandle {
252    fn wake_on(&self, _: Interest, waker: &Waker) {
253        wake_at_once(waker);
254    }
255}
256
257// ---------------------------------------------------------------------------
258// The event handles
259// ---------------------------------------------------------------------------
260
261/// The `EventSource` port role: one subscription set and one queue, both in
262/// the store under this handle's identity.
263///
264/// Each source handle receives its own copy of every occurrence raised while
265/// it is subscribed, so two consumers of one event never consume each other's
266/// occurrences.
267pub struct SourceHandle {
268    shared: Shared,
269    catalog: CatalogRef,
270    id: usize,
271}
272
273impl SourceHandle {
274    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
275        let id = lock(&shared).open_source();
276        SourceHandle {
277            shared,
278            catalog,
279            id,
280        }
281    }
282}
283
284/// The source's waker is dropped after the lock is released: it may be the
285/// last reference to a task that owns another handle of this runtime, and
286/// that handle's own `Drop` takes the lock.
287impl Drop for SourceHandle {
288    fn drop(&mut self) {
289        let waiters = lock(&self.shared).close_source(self.id);
290        drop(waiters);
291    }
292}
293
294impl Attached for SourceHandle {
295    fn catalog(&self) -> &CatalogRef {
296        &self.catalog
297    }
298}
299
300impl EventSource for SourceHandle {
301    fn subscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), SubscribeError> {
302        lock(&self.shared).subscribe(self.id, iface, ords);
303        Ok(())
304    }
305
306    fn unsubscribe(&mut self, iface: InterfaceNo, ords: &[Ordinal]) {
307        lock(&self.shared).unsubscribe(self.id, iface, ords);
308    }
309
310    fn next(&mut self, out: &mut [u8]) -> Result<Option<RawOccurrence>, ReadError> {
311        lock(&self.shared).next_event(self.id, out)
312    }
313}
314
315/// Stores one `Event` waker, whatever interface it was registered under, woken
316/// by any raise that queues an occurrence for this source.
317impl Wakeable for SourceHandle {
318    fn wake_on(&self, what: Interest, waker: &Waker) {
319        match what {
320            Interest::Event(_) => locked(&self.shared, |store, wake| {
321                store.wait_event(self.id, what, waker, wake);
322            }),
323            Interest::Outcome(_) | Interest::Slot | Interest::Claim(_) => wake_at_once(waker),
324        }
325    }
326}
327
328/// The `EventSink` port role.
329///
330/// The sequence counters are the handle's own, one per event channel, for the
331/// same reason the writer handle's are: ridl §3.1 scopes the number to the
332/// channel and has the sender assign it. One counter for the whole handle
333/// would make a consumer subscribed to some of this sink's events see a gap in
334/// `seq` where nothing was lost, and `EventSource::next` states that a gap is
335/// a loss.
336pub struct SinkHandle {
337    shared: Shared,
338    catalog: CatalogRef,
339    seqs: BTreeMap<Key, u64>,
340}
341
342impl SinkHandle {
343    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
344        SinkHandle {
345            shared,
346            catalog,
347            seqs: BTreeMap::new(),
348        }
349    }
350
351    fn take_seq(&mut self, key: Key) -> u64 {
352        let counter = self.seqs.entry(key).or_insert(0);
353        *counter += 1;
354        *counter
355    }
356}
357
358impl Attached for SinkHandle {
359    fn catalog(&self) -> &CatalogRef {
360        &self.catalog
361    }
362}
363
364impl EventSink for SinkHandle {
365    fn raise(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), RaiseError> {
366        let seq = self.take_seq((iface, ord));
367        locked(&self.shared, |store, wake| {
368            store.raise(iface, ord, bytes, seq, wake);
369        });
370        Ok(())
371    }
372}
373
374impl Wakeable for SinkHandle {
375    fn wake_on(&self, _: Interest, waker: &Waker) {
376        wake_at_once(waker);
377    }
378}
379
380// ---------------------------------------------------------------------------
381// The call handles
382// ---------------------------------------------------------------------------
383
384/// The `Caller` port role.
385///
386/// One counter for the whole handle, not one per channel: on a call the scope
387/// is the caller instance (ADR-0021 decision 5, and driftsys/ridl#308's own
388/// report), so every call this caller sends draws from one sequence. That is
389/// what keeps two callers on one provider from colliding.
390///
391/// The call table is the runtime's, shared by every caller: with
392/// [`Loopback::SLOTS`](crate::Loopback::SLOTS) calls sent and not forgotten,
393/// a send answers [`SendError::Busy`] and draws no sequence number, because
394/// nothing was sent.
395pub struct CallerHandle {
396    shared: Shared,
397    catalog: CatalogRef,
398    id: usize,
399    next_seq: u64,
400}
401
402impl CallerHandle {
403    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
404        let id = lock(&shared).open_caller();
405        CallerHandle {
406            shared,
407            catalog,
408            id,
409            next_seq: 0,
410        }
411    }
412
413    /// Sends one call. The sequence number is drawn only when the table
414    /// admits the call.
415    fn send(
416        &mut self,
417        kind: CallKind,
418        iface: InterfaceNo,
419        ord: Ordinal,
420        args: &[u8],
421    ) -> Result<Correlation, SendError> {
422        let seq = self.next_seq + 1;
423        let c = locked(&self.shared, |store, wake| {
424            store.send(self.id, kind, (iface, ord), args, seq, wake)
425        })?;
426        self.next_seq = seq;
427        Ok(c)
428    }
429}
430
431/// A dropped caller's `Slot` waker leaves the store, and every call it sent
432/// and did not forget is forgotten, so its slot comes back. The waker is
433/// dropped after the lock is released, as the source's is.
434impl Drop for CallerHandle {
435    fn drop(&mut self) {
436        let waiters = locked(&self.shared, |store, wake| {
437            store.close_caller(self.id, wake)
438        });
439        drop(waiters);
440    }
441}
442
443impl Attached for CallerHandle {
444    fn catalog(&self) -> &CatalogRef {
445        &self.catalog
446    }
447}
448
449/// The caller reads the runtime's one clock, so a generated async `Client`
450/// can compute a call's deadline over this handle alone.
451impl Clock for CallerHandle {
452    fn now(&self) -> Timestamp {
453        lock(&self.shared).now()
454    }
455}
456
457/// Stores one `Outcome` waker per call, kept with the call and woken by its
458/// settlement, its `forget`, the drop of this handle, or a displacement by
459/// another task. Stores one `Slot` waker, woken at once while a slot is free
460/// and otherwise by the next reclaim.
461impl Wakeable for CallerHandle {
462    fn wake_on(&self, what: Interest, waker: &Waker) {
463        match what {
464            Interest::Outcome(c) => locked(&self.shared, |store, wake| {
465                store.wait_outcome(c, waker, wake);
466            }),
467            Interest::Slot => locked(&self.shared, |store, wake| {
468                store.wait_slot(self.id, waker, wake);
469            }),
470            Interest::Event(_) | Interest::Claim(_) => wake_at_once(waker),
471        }
472    }
473}
474
475impl Caller for CallerHandle {
476    fn command(
477        &mut self,
478        iface: InterfaceNo,
479        ord: Ordinal,
480        args: &[u8],
481    ) -> Result<Correlation, SendError> {
482        self.send(CallKind::Command, iface, ord, args)
483    }
484
485    fn query(
486        &mut self,
487        iface: InterfaceNo,
488        ord: Ordinal,
489        args: &[u8],
490    ) -> Result<Correlation, SendError> {
491        self.send(CallKind::Query, iface, ord, args)
492    }
493
494    fn ack(&mut self, c: Correlation) -> Option<Result<(), CallError>> {
495        lock(&self.shared).ack(c)
496    }
497
498    fn reply(
499        &mut self,
500        c: Correlation,
501        out: &mut [u8],
502    ) -> Result<Option<Result<usize, CallError>>, ReadError> {
503        lock(&self.shared).reply(c, out)
504    }
505
506    fn forget(&mut self, c: Correlation) {
507        locked(&self.shared, |store, wake| store.forget(c, wake));
508    }
509}
510
511/// The `Handler` port role.
512///
513/// A handler that has served nothing is presented every call waiting in the
514/// store; once it has served anything, it is presented only the members it
515/// served, and another handler's calls stay waiting for that handler. A claim
516/// belongs to the handler it was presented to, so another handler's `settle`
517/// of it answers [`SettleError::UnknownClaim`].
518///
519/// The empty set meaning no filter is a deliberate deviation from
520/// [`Handler::serve`], which says delivery starts at the members listed. The
521/// generated `dispatch` never calls `serve`
522/// (`crates/ridl-backend-rust/src/face/dispatch.rs`), so a handler that always filtered
523/// would be presented nothing at all by it. `serve` with an empty slice
524/// records nothing and so leaves the handler unfiltered, the same as never
525/// having called it. [`served`](HandlerHandle::served) reads the set back.
526///
527/// The served set is kept twice: here, for `served` to return a slice, and in
528/// the store, where `next_claim` filters by it and a caller's send reads it to
529/// know which handler to wake. `serve` is the one method that changes it, and
530/// it changes both.
531pub struct HandlerHandle {
532    shared: Shared,
533    catalog: CatalogRef,
534    id: usize,
535    served: Vec<Key>,
536}
537
538impl HandlerHandle {
539    pub(crate) fn new(shared: Shared, catalog: CatalogRef) -> Self {
540        let id = lock(&shared).open_handler();
541        HandlerHandle {
542            shared,
543            catalog,
544            id,
545            served: Vec::new(),
546        }
547    }
548
549    /// The members `serve` was called with, in the order they were served,
550    /// with no duplicate.
551    #[must_use]
552    pub fn served(&self) -> &[(InterfaceNo, Ordinal)] {
553        &self.served
554    }
555}
556
557/// A claim this handler has taken and not settled returns to the waiting
558/// calls, and every handler that serves its member is woken. A claim it was
559/// only offered, through `ReadError::ShortClaim`, never left the waiting calls,
560/// so it is neither re-inserted nor woken for (driftsys/ridl#569). A claim whose
561/// call the caller forgot is withdrawn instead: its slot is reclaimed, and no
562/// handler is presented it again or woken for it. The handler's waker is
563/// dropped after the lock is released, as the source's is.
564impl Drop for HandlerHandle {
565    fn drop(&mut self) {
566        let waiters = locked(&self.shared, |store, wake| {
567            store.close_handler(self.id, wake)
568        });
569        drop(waiters);
570    }
571}
572
573impl Attached for HandlerHandle {
574    fn catalog(&self) -> &CatalogRef {
575        &self.catalog
576    }
577}
578
579impl Handler for HandlerHandle {
580    fn serve(&mut self, iface: InterfaceNo, ords: &[Ordinal]) -> Result<(), ServeError> {
581        for ord in ords {
582            if !self.served.contains(&(iface, *ord)) {
583                self.served.push((iface, *ord));
584            }
585        }
586        locked(&self.shared, |store, wake| {
587            store.serve(self.id, iface, ords, wake);
588        });
589        Ok(())
590    }
591
592    fn next_claim(&mut self, out: &mut [u8]) -> Result<Option<Claim>, ReadError> {
593        lock(&self.shared).next_claim(self.id, out)
594    }
595
596    fn settle(
597        &mut self,
598        claim: ClaimId,
599        outcome: Result<&[u8], CallError>,
600    ) -> Result<(), SettleError> {
601        locked(&self.shared, |store, wake| {
602            store.settle(self.id, claim, outcome, wake)
603        })
604    }
605}
606
607/// Stores one `Claim` waker, whatever interface it was registered under, woken
608/// by a send of a member this handler serves, by a `serve` that admits a call
609/// already waiting, or by another handler's drop that returns a claim this
610/// handler serves.
611impl Wakeable for HandlerHandle {
612    fn wake_on(&self, what: Interest, waker: &Waker) {
613        match what {
614            Interest::Claim(_) => locked(&self.shared, |store, wake| {
615                store.wait_claim(self.id, what, waker, wake);
616            }),
617            Interest::Outcome(_) | Interest::Slot | Interest::Event(_) => wake_at_once(waker),
618        }
619    }
620}