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}