Skip to main content

ridl_rt/
sample.rs

1//! Time, the envelope, and the values a read returns (ridl §3.1, §4.5, §9),
2//! with the two computations a consuming runtime makes on an envelope: a
3//! value's freshness ([`Freshness::of`]) and event loss ([`EventSeqTracker`]).
4
5use crate::contract::{InterfaceNo, Ordinal, Timing};
6use crate::payload::Violation;
7
8/// A point in time, in microseconds since the PTP epoch, on the TAI time scale
9/// (ridl §3.1).
10#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
11pub struct Timestamp(pub i64);
12
13/// A length of time, in microseconds.
14///
15/// This is not `core::time::Duration`: generated code writes
16/// `ridl_rt::sample::Duration` in full.
17#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
18pub struct Duration(pub i64);
19
20/// The sender's timestamp and sequence number (ridl §3.1). The runtime stamps
21/// it from its clock, and no relay changes it.
22///
23/// Before a signal's first publication, its envelope has `seq` 0 and the time
24/// at which the channel was created. On a call, `seq` is unique for each
25/// caller, not for each channel.
26#[derive(Clone, Copy, Debug, PartialEq, Eq)]
27pub struct Envelope {
28    /// When the sender published, raised or called the instance (ridl §3.1).
29    pub stamp: Timestamp,
30    /// The sender's sequence number. On an event channel, a gap is a loss
31    /// (ridl §3.1).
32    pub seq: u64,
33}
34
35/// Where a signal's value comes from (ridl §4.5).
36#[derive(Clone, Copy, Debug, PartialEq, Eq)]
37pub enum Provenance {
38    /// No publication yet: the value is the init value.
39    Init,
40    /// The value is the latest publication.
41    Live,
42    /// The channel is in the invalid state: the value is the last good value,
43    /// or the init value when there is none.
44    Invalid(Cause),
45}
46
47/// Why a channel is invalid.
48#[derive(Clone, Copy, Debug, PartialEq, Eq)]
49pub enum Cause {
50    /// The provider declared the invalid state (ridl §4.5).
51    Declared,
52    /// The consumer's binding detected an invalid payload.
53    Detected(Detection),
54}
55
56/// What a consumer's binding detected in a payload.
57#[derive(Clone, Copy, Debug, PartialEq, Eq)]
58pub enum Detection {
59    /// The payload breaks a typl constraint: `INVALID_VALUE`, ridl §10.2.
60    InvalidValue(Violation),
61    /// The payload is not a well-formed encoding: a serialization failure,
62    /// ridl §10.3.
63    Corrupt,
64}
65
66/// How old a value is, measured against its staleness bound (ridl §9).
67#[derive(Clone, Copy, Debug, PartialEq, Eq)]
68pub enum Freshness {
69    /// Within the bound.
70    Fresh,
71    /// Older than the bound.
72    Stale {
73        /// How far past the bound the value is.
74        by: Duration,
75    },
76    /// The value has no staleness bound: its member's timing has no `max`, as
77    /// under `@[1s..]`. A signal with no `@` annotation is not unbounded,
78    /// because it receives the default range (ridl §9.1).
79    Unbounded,
80}
81
82impl Freshness {
83    /// The freshness of a value whose envelope is `envelope`, read at `now`,
84    /// against the `max` of its member's `timing` (ridl §4, §9; frame
85    /// specification §8).
86    ///
87    /// With `stamp` the envelope's timestamp: `Fresh` while
88    /// `now − stamp ≤ max`, `Stale { by: now − stamp − max }` past it, and
89    /// `Unbounded` when `timing` is `None` or carries no `max`. `min` plays no
90    /// part. Under a strict period `@Xms`, `max` holds the period, so the
91    /// period is the bound. A `stamp` later than `now` gives a negative age,
92    /// which the formula makes `Fresh`. The age subtraction `now − stamp`
93    /// saturates at the ends of the `i64` range instead of overflowing.
94    ///
95    /// Pass the member's `timing` as the descriptor holds it:
96    /// `Freshness::of(&envelope, now, member.timing)`.
97    pub fn of(envelope: &Envelope, now: Timestamp, timing: Option<Timing>) -> Freshness {
98        let Some(max) = timing.and_then(|timing| timing.max) else {
99            return Freshness::Unbounded;
100        };
101        let age = now.0.saturating_sub(envelope.stamp.0);
102        if age <= max.0 {
103            Freshness::Fresh
104        } else {
105            Freshness::Stale {
106                by: Duration(age.saturating_sub(max.0)),
107            }
108        }
109    }
110}
111
112/// A signal value with its provenance, its freshness and its envelope.
113#[derive(Clone, Copy, Debug, PartialEq, Eq)]
114pub struct Sample<T> {
115    /// The value. Never absent: the init value under `Init`, and the last good
116    /// value or the init value under `Invalid`.
117    pub value: T,
118    /// Where the value comes from.
119    pub provenance: Provenance,
120    /// How old the value is.
121    pub freshness: Freshness,
122    /// The sender's timestamp and sequence number.
123    pub envelope: Envelope,
124}
125
126impl<T> Sample<T> {
127    /// `true` when the provenance is `Live` and the freshness is not `Stale`.
128    pub fn usable(&self) -> bool {
129        matches!(self.provenance, Provenance::Live)
130            && !matches!(self.freshness, Freshness::Stale { .. })
131    }
132}
133
134/// An event occurrence as a consumer receives it.
135#[derive(Clone, Copy, Debug, PartialEq, Eq)]
136pub struct Occurrence<T> {
137    /// The payload, or what the binding detected when the payload failed its
138    /// check.
139    pub payload: Result<T, Detection>,
140    /// The sender's timestamp and sequence number.
141    pub envelope: Envelope,
142}
143
144/// What [`EventSeqTracker::observe`] found when it compared an occurrence's
145/// `seq` with the last one its channel accepted (ridl §3.1; frame
146/// specification §7).
147#[derive(Clone, Copy, Debug, PartialEq, Eq)]
148pub enum Continuity {
149    /// The first occurrence the tracker has seen on this channel. No loss is
150    /// reported, because a consumer receives only the occurrences raised after
151    /// its subscription was answered (frame specification §6.2), so the first
152    /// one may carry any `seq`.
153    First,
154    /// The `seq` is one more than the last: nothing was lost.
155    Next,
156    /// The `seq` is more than one past the last: `count` occurrences between
157    /// the two were lost (ridl §3.1: "sequence gaps make loss detectable").
158    Lost {
159        /// How many sequence numbers the gap skipped: `seq − last − 1`.
160        count: u64,
161    },
162    /// The `seq` is not greater than the last: a duplicate or a reordered
163    /// occurrence. The tracker keeps `last`. What the caller does with the
164    /// occurrence is the caller's decision.
165    NotNewer {
166        /// The last `seq` the channel accepted, unchanged.
167        last: u64,
168    },
169}
170
171/// The error [`EventSeqTracker::observe`] returns when the occurrence is on a
172/// channel the tracker does not hold and every one of its `N` slots is taken.
173#[derive(Clone, Copy, Debug, PartialEq, Eq)]
174pub struct TrackerFull;
175
176/// The last `seq` accepted on each event channel of one session, and the loss
177/// each new occurrence reveals (ridl §3.1; frame specification §5.2 and §7).
178///
179/// The sequence counter of an event is per channel, per provider instance, and
180/// starts over with each session (frame specification §3, §7). A channel is
181/// one `(InterfaceNo, Ordinal)` in the session's catalog, so the tracker keys
182/// on that pair and never on the interface alone: two channels of one
183/// interface whose occurrences interleave are not a loss. One tracker serves
184/// one session, and so one provider instance; a consumer holding several
185/// sessions holds one tracker for each, and a new session starts with a new
186/// tracker.
187///
188/// The storage is the caller's: `N` slots, one per channel, held inline with
189/// no allocation. A channel that finds no free slot is refused with
190/// [`TrackerFull`]; [`forget`](Self::forget) frees a slot, for example when
191/// the consumer unsubscribes (frame specification §6.2).
192///
193/// Feed the tracker every occurrence the channel accepts, in the order
194/// received. A loss is a gap between two accepted occurrences of one channel
195/// (frame specification §5.2, §7, which attribute it to ridl §3.1).
196#[derive(Clone, Copy, Debug, PartialEq, Eq)]
197pub struct EventSeqTracker<const N: usize> {
198    slots: [Option<Channel>; N],
199}
200
201/// One tracked channel and the last `seq` it accepted.
202#[derive(Clone, Copy, Debug, PartialEq, Eq)]
203struct Channel {
204    interface: InterfaceNo,
205    ordinal: Ordinal,
206    last: u64,
207}
208
209impl<const N: usize> EventSeqTracker<N> {
210    /// A tracker holding no channel.
211    pub const fn new() -> Self {
212        EventSeqTracker { slots: [None; N] }
213    }
214
215    /// Records `seq` as the latest occurrence on the channel
216    /// `(interface, ordinal)` and reports what it reveals.
217    ///
218    /// The first occurrence on a channel takes a free slot and reports
219    /// [`Continuity::First`]. A later one reports [`Continuity::Next`] or
220    /// [`Continuity::Lost`] and becomes the channel's last `seq`, or reports
221    /// [`Continuity::NotNewer`] and leaves the last `seq` as it was.
222    ///
223    /// # Errors
224    ///
225    /// [`TrackerFull`] when the channel is not tracked and no slot is free.
226    /// Every tracked channel is unchanged.
227    pub fn observe(
228        &mut self,
229        interface: InterfaceNo,
230        ordinal: Ordinal,
231        seq: u64,
232    ) -> Result<Continuity, TrackerFull> {
233        let mut free = None;
234        for (index, slot) in self.slots.iter_mut().enumerate() {
235            match slot {
236                Some(channel) if channel.interface == interface && channel.ordinal == ordinal => {
237                    let last = channel.last;
238                    if seq <= last {
239                        return Ok(Continuity::NotNewer { last });
240                    }
241                    channel.last = seq;
242                    let count = seq - last - 1;
243                    return Ok(if count == 0 {
244                        Continuity::Next
245                    } else {
246                        Continuity::Lost { count }
247                    });
248                }
249                None if free.is_none() => free = Some(index),
250                _ => {}
251            }
252        }
253        let index = free.ok_or(TrackerFull)?;
254        self.slots[index] = Some(Channel {
255            interface,
256            ordinal,
257            last: seq,
258        });
259        Ok(Continuity::First)
260    }
261
262    /// Stops tracking the channel `(interface, ordinal)` and frees its slot.
263    /// The next occurrence on that channel reports [`Continuity::First`]. A
264    /// channel the tracker does not hold is left alone.
265    pub fn forget(&mut self, interface: InterfaceNo, ordinal: Ordinal) {
266        for slot in &mut self.slots {
267            if matches!(slot, Some(channel) if channel.interface == interface && channel.ordinal == ordinal)
268            {
269                *slot = None;
270            }
271        }
272    }
273}
274
275impl<const N: usize> Default for EventSeqTracker<N> {
276    fn default() -> Self {
277        Self::new()
278    }
279}
280
281#[cfg(test)]
282mod tests {
283    use super::{Cause, Detection, Duration, Envelope, Freshness, Provenance, Sample, Timestamp};
284    use crate::payload::{Rule, Violation};
285
286    fn sample(provenance: Provenance, freshness: Freshness) -> Sample<u8> {
287        Sample {
288            value: 0,
289            provenance,
290            freshness,
291            envelope: Envelope {
292                stamp: Timestamp(0),
293                seq: 0,
294            },
295        }
296    }
297
298    #[test]
299    fn only_a_live_value_that_is_not_stale_is_usable() {
300        let violation = Violation {
301            type_name: "Speed",
302            rule: Rule::Range,
303        };
304        let provenances = [
305            Provenance::Init,
306            Provenance::Live,
307            Provenance::Invalid(Cause::Declared),
308            Provenance::Invalid(Cause::Detected(Detection::InvalidValue(violation))),
309            Provenance::Invalid(Cause::Detected(Detection::Corrupt)),
310        ];
311        let freshnesses = [
312            Freshness::Fresh,
313            Freshness::Stale { by: Duration(1) },
314            Freshness::Unbounded,
315        ];
316        for provenance in provenances {
317            for freshness in freshnesses {
318                let expected =
319                    provenance == Provenance::Live && !matches!(freshness, Freshness::Stale { .. });
320                assert_eq!(
321                    sample(provenance, freshness).usable(),
322                    expected,
323                    "{provenance:?} with {freshness:?}"
324                );
325            }
326        }
327    }
328}