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. When the check fails, the accessor reports the detection
130/// as the provenance and substitutes the init value. A generated accessor
131/// would substitute its own last good value when it has one.
132pub fn read_speed(port: &dyn SignalReader) -> Result<Sample<Speed>, ReadError> {
133    let mut buf = [0u8; <Speed as Payload<ReprC>>::MAX_SIZE];
134    let raw = port.read(Drivetrain::NUMBER, SpeedSignal::MEMBER.ordinal, &mut buf)?;
135    let (value, provenance) = match Ref::<Speed, ReprC>::verify(&buf[..raw.len]) {
136        Ok(proof) => (proof.decode(), raw.provenance),
137        Err(VerifyError::Contract(violation)) => (
138            SpeedSignal::init(),
139            Provenance::Invalid(Cause::Detected(Detection::InvalidValue(violation))),
140        ),
141        Err(_) => (
142            SpeedSignal::init(),
143            Provenance::Invalid(Cause::Detected(Detection::Corrupt)),
144        ),
145    };
146    Ok(Sample {
147        value,
148        provenance,
149        freshness: raw.freshness,
150        envelope: raw.envelope,
151    })
152}
153
154// ---- What a runtime writes ----------------------------------------------------
155
156/// A runtime for one signal of at most 2 bytes: a manual clock, one slot, and
157/// one staged change.
158pub struct Memory {
159    now: Cell<i64>,
160    bytes: [u8; 2],
161    len: usize,
162    provenance: Provenance,
163    envelope: Envelope,
164    staged: Option<Staged>,
165}
166
167enum Staged {
168    Value([u8; 2], usize),
169    Invalid,
170    Touch,
171}
172
173impl Memory {
174    /// Creates the channel at `now`, holding the init value's bytes.
175    pub fn new(now: Timestamp, init: &[u8]) -> Memory {
176        let mut bytes = [0u8; 2];
177        bytes[..init.len()].copy_from_slice(init);
178        Memory {
179            now: Cell::new(now.0),
180            bytes,
181            len: init.len(),
182            provenance: Provenance::Init,
183            envelope: Envelope { stamp: now, seq: 0 },
184            staged: None,
185        }
186    }
187
188    /// Moves the clock forward.
189    pub fn advance(&self, by: Duration) {
190        self.now.set(self.now.get() + by.0);
191    }
192
193    fn member(iface: InterfaceNo, ord: Ordinal) -> Option<&'static Member> {
194        if iface != Drivetrain::NUMBER {
195            return None;
196        }
197        Drivetrain::MEMBERS
198            .iter()
199            .find(|m| m.ordinal == ord && m.kind == Kind::Signal)
200    }
201}
202
203impl Attached for Memory {
204    fn catalog(&self) -> &CatalogRef {
205        Drivetrain::CATALOG
206    }
207}
208
209impl Clock for Memory {
210    fn now(&self) -> Timestamp {
211        Timestamp(self.now.get())
212    }
213}
214
215impl SignalReader for Memory {
216    fn read(
217        &self,
218        iface: InterfaceNo,
219        ord: Ordinal,
220        out: &mut [u8],
221    ) -> Result<RawSample, ReadError> {
222        let member =
223            Memory::member(iface, ord).ok_or(ReadError::Contract(Contract::UnknownInteraction))?;
224        let Some(front) = out.get_mut(..self.len) else {
225            return Err(ReadError::Short { needed: self.len });
226        };
227        front.copy_from_slice(&self.bytes[..self.len]);
228        let freshness = match member.timing.and_then(|t| t.max) {
229            None => Freshness::Unbounded,
230            Some(max) => {
231                let age = self.now().0 - self.envelope.stamp.0;
232                if age > max.0 {
233                    Freshness::Stale {
234                        by: Duration(age - max.0),
235                    }
236                } else {
237                    Freshness::Fresh
238                }
239            }
240        };
241        Ok(RawSample {
242            provenance: self.provenance,
243            freshness,
244            envelope: self.envelope,
245            len: self.len,
246        })
247    }
248}
249
250impl SignalWriter for Memory {
251    fn set(&mut self, iface: InterfaceNo, ord: Ordinal, bytes: &[u8]) -> Result<(), WriteError> {
252        if Memory::member(iface, ord).is_none() {
253            return Err(WriteError::Contract(Contract::UnknownInteraction));
254        }
255        let mut value = [0u8; 2];
256        let Some(front) = value.get_mut(..bytes.len()) else {
257            return Err(WriteError::TooLarge { cap: 2 });
258        };
259        front.copy_from_slice(bytes);
260        self.staged = Some(Staged::Value(value, bytes.len()));
261        Ok(())
262    }
263
264    fn invalidate(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
265        if Memory::member(iface, ord).is_none() {
266            return Err(WriteError::Contract(Contract::UnknownInteraction));
267        }
268        self.staged = Some(Staged::Invalid);
269        Ok(())
270    }
271
272    fn touch(&mut self, iface: InterfaceNo, ord: Ordinal) -> Result<(), WriteError> {
273        if Memory::member(iface, ord).is_none() {
274            return Err(WriteError::Contract(Contract::UnknownInteraction));
275        }
276        if self.staged.is_none() {
277            self.staged = Some(Staged::Touch);
278        }
279        Ok(())
280    }
281
282    fn commit(&mut self) {
283        let Some(staged) = self.staged.take() else {
284            return;
285        };
286        match staged {
287            Staged::Value(bytes, len) => {
288                self.bytes = bytes;
289                self.len = len;
290                self.provenance = Provenance::Live;
291            }
292            Staged::Invalid => self.provenance = Provenance::Invalid(Cause::Declared),
293            Staged::Touch => {}
294        }
295        self.envelope = Envelope {
296            stamp: self.now(),
297            seq: self.envelope.seq + 1,
298        };
299    }
300}
301
302// ---- The program --------------------------------------------------------------
303
304/// Publishes `speed` the way a generated publisher does: encode, stage, commit.
305fn publish(runtime: &mut Memory, speed: Speed) {
306    let mut buf = [0u8; <Speed as Payload<ReprC>>::MAX_SIZE];
307    let proof = Ref::<Speed, ReprC>::encode(&speed, &mut buf).expect("the buffer is MAX_SIZE");
308    runtime
309        .set(
310            Drivetrain::NUMBER,
311            SpeedSignal::MEMBER.ordinal,
312            proof.bytes(),
313        )
314        .expect("the runtime owns speed");
315    runtime.commit();
316}
317
318/// Moves `speed` through six states and returns the sample read after each:
319/// the init value; a live value; the same value once its staleness bound has
320/// passed; the invalid state the provider declares; a value the accessor's
321/// check rejects; and bytes the accessor cannot decode.
322pub fn walk() -> [Sample<Speed>; 6] {
323    let mut init = [0u8; <Speed as Payload<ReprC>>::MAX_SIZE];
324    let proof = Ref::<Speed, ReprC>::encode(&SpeedSignal::init(), &mut init)
325        .expect("the buffer is MAX_SIZE");
326    let mut runtime = Memory::new(Timestamp(1_000_000), proof.bytes());
327    let read = |runtime: &Memory| read_speed(runtime).expect("speed is in the catalog");
328
329    let at_init = read(&runtime);
330
331    runtime.advance(Duration(10_000));
332    publish(&mut runtime, Speed(88));
333    let live = read(&runtime);
334
335    runtime.advance(Duration(600_000));
336    let stale = read(&runtime);
337
338    runtime
339        .invalidate(Drivetrain::NUMBER, SpeedSignal::MEMBER.ordinal)
340        .expect("the runtime owns speed");
341    runtime.commit();
342    let declared = read(&runtime);
343
344    runtime
345        .set(
346            Drivetrain::NUMBER,
347            SpeedSignal::MEMBER.ordinal,
348            &400u16.to_le_bytes(),
349        )
350        .expect("the runtime owns speed");
351    runtime.commit();
352    let detected = read(&runtime);
353
354    // One byte where the encoding needs two.
355    runtime
356        .set(Drivetrain::NUMBER, SpeedSignal::MEMBER.ordinal, &[1])
357        .expect("the runtime owns speed");
358    runtime.commit();
359    let corrupt = read(&runtime);
360
361    [at_init, live, stale, declared, detected, corrupt]
362}
363
364fn main() {
365    let steps = [
366        "init",
367        "live",
368        "stale",
369        "declared invalid",
370        "detected invalid",
371        "corrupt",
372    ];
373    for (step, sample) in steps.iter().zip(walk()) {
374        println!("{step}: {sample:?} usable={}", sample.usable());
375    }
376}