Skip to main content

read_sample/
read_sample.rs

1//! Reads a signal with its provenance, using `ridl-rt` and nothing else.
2//!
3//! This program writes by hand what a code generator and a runtime would
4//! provide: one catalog with one interface, `Drivetrain`; one signal, `speed`;
5//! a `repr(C)` codec for its payload type; an in-memory runtime that
6//! implements `Clock`, `SignalReader` and `SignalWriter`; and the accessor a
7//! generated client would contain. `main` moves the signal through its
8//! provenance and freshness states and prints the sample read after each
9//! step.
10//!
11//! Run it with `cargo run -p ridl-rt --example read_sample`.
12//! `tests/read_sample.rs` runs the same steps and checks each sample.
13
14#![forbid(unsafe_code)]
15
16use std::cell::Cell;
17
18use ridl_rt::contract::{
19    CatalogHash, CatalogRef, EncodedSizes, Interaction, Interface, InterfaceNo, Kind, Member,
20    Ordinal, PayloadInfo, Signal, Timing, TimingMode,
21};
22use ridl_rt::encoding::ReprC;
23use ridl_rt::error::Contract;
24use ridl_rt::payload::{
25    EncodeError, Encoded, Malformed, Payload, Ref, Rule, VerifyError, Violation,
26};
27use ridl_rt::port::{
28    Attached, Clock, RawSample, ReadError, SignalReader, SignalWriter, WriteError,
29};
30use ridl_rt::sample::{
31    Cause, Detection, Duration, Envelope, Freshness, Provenance, Sample, Timestamp,
32};
33
34// ---- What a code generator writes -------------------------------------------
35
36/// Speed in km/h, declared in the range 0 to 300.
37#[derive(Clone, Copy, Debug, PartialEq, Eq)]
38pub struct Speed(pub u16);
39
40impl Payload<ReprC> for Speed {
41    const MAX_SIZE: usize = 2;
42    type View<'a> = &'a [u8];
43
44    fn encode<'o>(&self, out: &'o mut [u8]) -> Result<Encoded<'o, &'o [u8]>, EncodeError> {
45        let available = out.len();
46        let Some(front) = out.get_mut(..2) else {
47            return Err(EncodeError::Capacity {
48                needed: 2,
49                available,
50            });
51        };
52        front.copy_from_slice(&self.0.to_le_bytes());
53        let bytes: &'o [u8] = front;
54        Ok(Encoded { bytes, view: bytes })
55    }
56
57    fn verify(buf: &[u8]) -> Result<&[u8], VerifyError> {
58        let bytes: [u8; 2] = buf
59            .try_into()
60            .map_err(|_| VerifyError::Structure(Malformed::OutOfBounds))?;
61        if u16::from_le_bytes(bytes) > 300 {
62            return Err(VerifyError::Contract(Violation {
63                type_name: "Speed",
64                rule: Rule::Range,
65            }));
66        }
67        Ok(buf)
68    }
69
70    fn decode(r: Ref<'_, Self, ReprC>) -> Self {
71        let b = r.bytes();
72        Speed(u16::from_le_bytes([b[0], b[1]]))
73    }
74}
75
76/// The catalog of the package `vehicle`.
77pub const VEHICLE: CatalogRef = CatalogRef {
78    name: "vehicle",
79    hash: CatalogHash([7; 32]),
80};
81
82/// `interface Drivetrain { signal speed : Speed @[20ms..500ms] }`
83pub struct Drivetrain;
84
85impl Interface for Drivetrain {
86    const CATALOG: &'static CatalogRef = &VEHICLE;
87    const NUMBER: InterfaceNo = InterfaceNo(1);
88    const PROVISIONAL: bool = false;
89    const NAME: &'static str = "Drivetrain";
90    const MEMBERS: &'static [Member] = &[Member {
91        ordinal: Ordinal(1),
92        kind: Kind::Signal,
93        name: "speed",
94        timing: Some(Timing {
95            mode: TimingMode::Range,
96            min: Some(Duration(20_000)),
97            max: Some(Duration(500_000)),
98        }),
99        // A generated descriptor fills in every encoding that can carry the
100        // payload. This program writes only the `repr(C)` codec, so it fills
101        // in only that size.
102        payloads: &[PayloadInfo {
103            type_name: "Speed",
104            max_size: EncodedSizes {
105                proto3: None,
106                flatbuffers: None,
107                repr_c: Some(2),
108            },
109        }],
110    }];
111}
112
113/// The descriptor of `Drivetrain.speed`.
114pub struct SpeedSignal;
115
116impl Interaction for SpeedSignal {
117    type Iface = Drivetrain;
118    const MEMBER: &'static Member = &Drivetrain::MEMBERS[0];
119}
120
121impl Signal for SpeedSignal {
122    type Payload = Speed;
123    fn init() -> Speed {
124        Speed(30)
125    }
126}
127
128/// The accessor a generated client writes for `speed`: read the bytes, check
129/// them, decode them. The `Init`-or-`Invalid(Declared)`-with-no-bytes arm
130/// below handles a port whose init or invalidated sample carries no bytes;
131/// `walk` never produces that shape, because it seeds `Memory` with a real
132/// 2-byte init encoding and `invalidate` keeps those bytes rather than
133/// clearing them (`tests/read_sample.rs` exercises the arm directly, with a
134/// `Memory` seeded from an empty init buffer). Otherwise the check runs;
135/// when it fails, the accessor reports the detection as the provenance and
136/// substitutes the init value. A generated accessor would substitute its own
137/// last good value when it has one.
138pub fn read_speed(port: &dyn SignalReader) -> Result<Sample<Speed>, ReadError> {
139    let mut buf = [0u8; <Speed as Payload<ReprC>>::MAX_SIZE];
140    let raw = port.read(Drivetrain::NUMBER, SpeedSignal::MEMBER.ordinal, &mut buf)?;
141    let (value, provenance) = match raw.provenance {
142        Provenance::Init | Provenance::Invalid(Cause::Declared) if raw.len == 0 => {
143            (SpeedSignal::init(), raw.provenance)
144        }
145        _ => match Ref::<Speed, ReprC>::verify(&buf[..raw.len]) {
146            Ok(proof) => (proof.decode(), raw.provenance),
147            Err(VerifyError::Contract(violation)) => (
148                SpeedSignal::init(),
149                Provenance::Invalid(Cause::Detected(Detection::InvalidValue(violation))),
150            ),
151            Err(_) => (
152                SpeedSignal::init(),
153                Provenance::Invalid(Cause::Detected(Detection::Corrupt)),
154            ),
155        },
156    };
157    Ok(Sample {
158        value,
159        provenance,
160        freshness: raw.freshness,
161        envelope: raw.envelope,
162    })
163}
164
165// ---- What a runtime writes ----------------------------------------------------
166
167/// A runtime for one signal of at most 2 bytes: a manual clock, one slot, and
168/// one staged change.
169pub struct Memory {
170    now: Cell<i64>,
171    bytes: [u8; 2],
172    len: usize,
173    provenance: Provenance,
174    envelope: Envelope,
175    staged: Option<Staged>,
176}
177
178enum Staged {
179    Value([u8; 2], usize),
180    Invalid,
181    Touch,
182}
183
184impl Memory {
185    /// Creates the channel at `now`, holding the init value's bytes.
186    pub fn new(now: Timestamp, init: &[u8]) -> Memory {
187        let mut bytes = [0u8; 2];
188        bytes[..init.len()].copy_from_slice(init);
189        Memory {
190            now: Cell::new(now.0),
191            bytes,
192            len: init.len(),
193            provenance: Provenance::Init,
194            envelope: Envelope { stamp: now, seq: 0 },
195            staged: None,
196        }
197    }
198
199    /// Moves the clock forward.
200    pub fn advance(&self, by: Duration) {
201        self.now.set(self.now.get() + by.0);
202    }
203
204    fn member(iface: InterfaceNo, ord: Ordinal) -> Option<&'static Member> {
205        if iface != Drivetrain::NUMBER {
206            return None;
207        }
208        Drivetrain::MEMBERS
209            .iter()
210            .find(|m| m.ordinal == ord && m.kind == Kind::Signal)
211    }
212}
213
214impl Attached for Memory {
215    fn catalog(&self) -> &CatalogRef {
216        Drivetrain::CATALOG
217    }
218}
219
220impl Clock for Memory {
221    fn now(&self) -> Timestamp {
222        Timestamp(self.now.get())
223    }
224}
225
226impl SignalReader for Memory {
227    fn read(
228        &self,
229        iface: InterfaceNo,
230        ord: Ordinal,
231        out: &mut [u8],
232    ) -> Result<RawSample, ReadError> {
233        let member =
234            Memory::member(iface, ord).ok_or(ReadError::Contract(Contract::UnknownInteraction))?;
235        let Some(front) = out.get_mut(..self.len) else {
236            return Err(ReadError::Short { needed: self.len });
237        };
238        front.copy_from_slice(&self.bytes[..self.len]);
239        let freshness = match member.timing.and_then(|t| t.max) {
240            None => Freshness::Unbounded,
241            Some(max) => {
242                let age = self.now().0 - self.envelope.stamp.0;
243                if age > max.0 {
244                    Freshness::Stale {
245                        by: Duration(age - max.0),
246                    }
247                } else {
248                    Freshness::Fresh
249                }
250            }
251        };
252        Ok(RawSample {
253            provenance: self.provenance,
254            freshness,
255            envelope: self.envelope,
256            len: self.len,
257        })
258    }
259}
260
261impl SignalWriter for Memory {
262    fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
263        if Memory::member(iface, ord).is_none() {
264            return Err(WriteError::Contract(Contract::UnknownInteraction));
265        }
266        let mut value = [0u8; 2];
267        let Some(front) = value.get_mut(..bytes.len()) else {
268            return Err(WriteError::TooLarge { cap: 2 });
269        };
270        front.copy_from_slice(bytes);
271        self.staged = Some(Staged::Value(value, bytes.len()));
272        Ok(())
273    }
274
275    fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
276        if Memory::member(iface, ord).is_none() {
277            return Err(WriteError::Contract(Contract::UnknownInteraction));
278        }
279        self.staged = Some(Staged::Invalid);
280        Ok(())
281    }
282
283    fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
284        if Memory::member(iface, ord).is_none() {
285            return Err(WriteError::Contract(Contract::UnknownInteraction));
286        }
287        if self.staged.is_none() {
288            self.staged = Some(Staged::Touch);
289        }
290        Ok(())
291    }
292
293    fn commit(&mut self) {
294        let Some(staged) = self.staged.take() else {
295            return;
296        };
297        match staged {
298            Staged::Value(bytes, len) => {
299                self.bytes = bytes;
300                self.len = len;
301                self.provenance = Provenance::Live;
302            }
303            Staged::Invalid => self.provenance = Provenance::Invalid(Cause::Declared),
304            Staged::Touch => {}
305        }
306        self.envelope = Envelope {
307            stamp: self.now(),
308            seq: self.envelope.seq + 1,
309        };
310    }
311}
312
313// ---- The program --------------------------------------------------------------
314
315/// Publishes `speed` the way a generated publisher does: encode, stage, commit.
316fn publish(runtime: &mut Memory, speed: Speed) {
317    let mut buf = [0u8; <Speed as Payload<ReprC>>::MAX_SIZE];
318    let proof = Ref::<Speed, ReprC>::encode(&speed, &mut buf).expect("the buffer is MAX_SIZE");
319    runtime
320        .set(
321            Drivetrain::NUMBER,
322            SpeedSignal::MEMBER.ordinal,
323            proof.bytes(),
324        )
325        .expect("the runtime owns speed");
326    runtime.commit();
327}
328
329/// Moves `speed` through six states and returns the sample read after each:
330/// the init value; a live value; the same value once its staleness bound has
331/// passed; the invalid state the provider declares; a value the accessor's
332/// check rejects; and bytes the accessor cannot decode.
333pub fn walk() -> [Sample<Speed>; 6] {
334    let mut init = [0u8; <Speed as Payload<ReprC>>::MAX_SIZE];
335    let proof = Ref::<Speed, ReprC>::encode(&SpeedSignal::init(), &mut init)
336        .expect("the buffer is MAX_SIZE");
337    let mut runtime = Memory::new(Timestamp(1_000_000), proof.bytes());
338    let read = |runtime: &Memory| read_speed(runtime).expect("speed is in the catalog");
339
340    let at_init = read(&runtime);
341
342    runtime.advance(Duration(10_000));
343    publish(&mut runtime, Speed(88));
344    let live = read(&runtime);
345
346    runtime.advance(Duration(600_000));
347    let stale = read(&runtime);
348
349    runtime
350        .invalidate(Drivetrain::NUMBER, SpeedSignal::MEMBER.ordinal)
351        .expect("the runtime owns speed");
352    runtime.commit();
353    let declared = read(&runtime);
354
355    runtime
356        .set(
357            Drivetrain::NUMBER,
358            SpeedSignal::MEMBER.ordinal,
359            &400u16.to_le_bytes(),
360        )
361        .expect("the runtime owns speed");
362    runtime.commit();
363    let detected = read(&runtime);
364
365    // One byte where the encoding needs two.
366    runtime
367        .set(Drivetrain::NUMBER, SpeedSignal::MEMBER.ordinal, &[1])
368        .expect("the runtime owns speed");
369    runtime.commit();
370    let corrupt = read(&runtime);
371
372    [at_init, live, stale, declared, detected, corrupt]
373}
374
375fn main() {
376    let steps = [
377        "init",
378        "live",
379        "stale",
380        "declared invalid",
381        "detected invalid",
382        "corrupt",
383    ];
384    for (step, sample) in steps.iter().zip(walk()) {
385        println!("{step}: {sample:?} usable={}", sample.usable());
386    }
387}