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