Skip to main content

QueueStreamedFn

Trait QueueStreamedFn 

Source
pub trait QueueStreamedFn<T>
where T: Deserialize + Sized,
{ // Required methods fn fetch(&self, buffer: &mut T, time: TickType) -> Result<()>; fn fetch_from_isr(&self, buffer: &mut T) -> Result<()>; fn post(&self, item: &T, time: TickType) -> Result<()>; fn post_from_isr(&self, item: &T) -> Result<()>; fn delete(&mut self); }
Expand description

All OSAL trait definitions for advanced usage. Type-safe queue for structured message passing.

This trait provides a queue that works with specific types, offering compile-time type safety for queue operations.

§Type Safety

Unlike raw Queue, QueueStreamed ensures that only messages of type T can be sent and received, preventing type confusion at compile time.

§Serialization

Messages are automatically serialized when sent and deserialized when received. The type T must implement the Deserialize trait.

§Type Parameters

  • T - The message type (must implement Deserialize)

§Examples

use osal_rs::os::*;
use osal_rs::utils::{Error, Result};

// Wire format: 4 bytes of `id`, 2 of `temperature`, 1 of `humidity`,
// little-endian. Keeping the message in its encoded form is what lets
// `to_bytes` hand out a borrowed slice.
const SENSOR_DATA_LEN: usize = 7;

#[derive(Clone, Copy, Debug, PartialEq)]
struct SensorData([u8; SENSOR_DATA_LEN]);

impl SensorData {
    fn new(id: u32, temperature: i16, humidity: u8) -> Self {
        let mut raw = [0u8; SENSOR_DATA_LEN];
        raw[..4].copy_from_slice(&id.to_le_bytes());
        raw[4..6].copy_from_slice(&temperature.to_le_bytes());
        raw[6] = humidity;
        Self(raw)
    }

    fn id(&self) -> u32 {
        u32::from_le_bytes(self.0[..4].try_into().unwrap())
    }
}

impl BytesHasLen for SensorData {
    fn len(&self) -> usize { SENSOR_DATA_LEN }
}

impl Serialize for SensorData {
    fn to_bytes(&self) -> &[u8] { &self.0 }
}

impl Deserialize for SensorData {
    fn from_bytes(bytes: &[u8]) -> Result<Self> {
        if bytes.len() < SENSOR_DATA_LEN {
            return Err(Error::OutOfIndex);
        }
        let mut raw = [0u8; SENSOR_DATA_LEN];
        raw.copy_from_slice(&bytes[..SENSOR_DATA_LEN]);
        Ok(Self(raw))
    }
}

let queue = QueueStreamed::<SensorData>::new(10, SENSOR_DATA_LEN as _).unwrap();

// Producer
let data = SensorData::new(1, 235, 65);
queue.post(&data, 100).unwrap();

// Consumer
let mut received = SensorData::new(0, 0, 0);
queue.fetch(&mut received, 100).unwrap();
assert_eq!(received.id(), 1);
assert_eq!(received, data);

Required Methods§

Source

fn fetch(&self, buffer: &mut T, time: TickType) -> Result<()>

Fetches a typed message from the queue (blocking).

Removes and deserializes the oldest message from the queue. Blocks the calling task if the queue is empty.

§Parameters
  • buffer - Mutable reference to receive the deserialized message
  • time - Maximum ticks to wait for a message:
    • 0: Return immediately if empty
    • n: Wait up to n ticks
    • TickType::MAX: Wait forever
§Returns
  • Ok(()) - Message received and deserialized successfully
  • Err(Error::Timeout) - Queue was empty for entire timeout period
  • Err(Error) - Deserialization error or other error
§Examples
use osal_rs::os::*;
let queue = QueueStreamed::<Message>::new(4, 4).unwrap();
queue.post(&Message::new(42), 100).unwrap();

let mut msg = Message::default();

match queue.fetch(&mut msg, 1000) {
    Ok(()) => assert_eq!(msg.id(), 42),
    Err(_) => panic!("no message available"),
}
Source

fn fetch_from_isr(&self, buffer: &mut T) -> Result<()>

Fetches a typed message from ISR context (non-blocking).

ISR-safe version of fetch(). Returns immediately without blocking. Must only be called from interrupt context.

§Parameters
  • buffer - Mutable reference to receive the deserialized message
§Returns
  • Ok(()) - Message received and deserialized successfully
  • Err(Error) - Queue is empty or deserialization failed
§Examples
use osal_rs::os::*;
let queue = QueueStreamed::<Message>::new(4, 4).unwrap();
queue.post(&Message::new(7), 100).unwrap();

// In interrupt handler
let mut msg = Message::default();
if queue.fetch_from_isr(&mut msg).is_ok() {
    // Process message
    assert_eq!(msg.id(), 7);
}

// Queue empty: reported immediately instead of blocking the "ISR".
assert!(queue.fetch_from_isr(&mut msg).is_err());
Source

fn post(&self, item: &T, time: TickType) -> Result<()>

Posts a typed message to the queue (blocking).

Serializes and adds a new message to the end of the queue. Blocks the calling task if the queue is full.

§Parameters
  • item - Reference to the message to serialize and send
  • time - Maximum ticks to wait if queue is full:
    • 0: Return immediately if full
    • n: Wait up to n ticks for space
    • TickType::MAX: Wait forever
§Returns
  • Ok(()) - Message serialized and sent successfully
  • Err(Error::Timeout) - Queue was full for entire timeout period
  • Err(Error) - Serialization error or other error
§Examples
use osal_rs::os::*;
// Room for a single message.
let queue = QueueStreamed::<Message>::new(1, 4).unwrap();
let msg = Message::new(42);

match queue.post(&msg, 1000) {
    Ok(()) => (), // sent successfully
    Err(_) => panic!("failed to send"),
}

// The only slot is taken and nobody is fetching: this one times out.
assert!(queue.post(&msg, 10).is_err());
Source

fn post_from_isr(&self, item: &T) -> Result<()>

Posts a typed message from ISR context (non-blocking).

ISR-safe version of post(). Returns immediately without blocking. Must only be called from interrupt context.

§Parameters
  • item - Reference to the message to serialize and send
§Returns
  • Ok(()) - Message serialized and sent successfully
  • Err(Error) - Queue is full or serialization failed
§Examples
use osal_rs::os::*;
let queue = QueueStreamed::<Message>::new(1, 4).unwrap();

// In interrupt handler
let msg = Message::new(1);
if queue.post_from_isr(&msg).is_err() {
    // Queue full, message dropped
}

// The single slot is now taken, so the next one really is dropped.
assert!(queue.post_from_isr(&msg).is_err());
Source

fn delete(&mut self)

Deletes the queue and frees its resources.

§Safety

Ensure no tasks are blocked on this queue before deletion. Calling this while tasks are waiting may cause undefined behavior.

§Examples
use osal_rs::os::*;
let mut queue = QueueStreamed::<Message>::new(10, core::mem::size_of::<Message>() as _).unwrap();

// Use queue...
queue.post(&Message::new(1), 100).unwrap();

queue.delete();

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§

Source§

impl<T> QueueStreamed<T> for QueueStreamed<T>
where T: StructSerde,

Available on crate feature serde only.