Skip to main content

async_ach_spsc/
heapless.rs

1use ach_spsc::heapless as ach;
2use async_ach_notify::Notify;
3use futures_util::StreamExt;
4
5pub struct Spsc<T, const N: usize> {
6    buf: ach::Spsc<T, N>,
7    consumer: Notify<1>,
8    producer: Notify<1>,
9}
10impl<T, const N: usize> Spsc<T, N> {
11    pub const fn new() -> Self {
12        Self {
13            buf: ach::Spsc::new(),
14            consumer: Notify::new(),
15            producer: Notify::new(),
16        }
17    }
18}
19impl<T, const N: usize> Default for Spsc<T, N> {
20    fn default() -> Self {
21        Self::new()
22    }
23}
24impl<T: Unpin, const N: usize> Spsc<T, N> {
25    pub fn take_sender(&'_ self) -> Option<Sender<'_, T, N>> {
26        let sender = self.buf.take_sender()?;
27        Some(Sender {
28            parent: self,
29            sender,
30        })
31    }
32    pub fn take_recver(&'_ self) -> Option<Receiver<'_, T, N>> {
33        let recver = self.buf.take_recver()?;
34        Some(Receiver {
35            parent: self,
36            recver,
37        })
38    }
39}
40
41pub struct Sender<'a, T: Unpin, const N: usize> {
42    parent: &'a Spsc<T, N>,
43    sender: ach::Sender<'a, T, N>,
44}
45impl<'a, T: Unpin, const N: usize> Sender<'a, T, N> {
46    pub fn try_send(&mut self, val: T) -> Result<(), T> {
47        self.sender.try_send(val).map(|_| {
48            self.parent.producer.notify_one();
49        })
50    }
51    pub async fn send(&mut self, mut val: T) {
52        let mut wait_c = self.parent.consumer.listen();
53        while let Err(v) = self.try_send(val) {
54            val = v;
55            wait_c.next().await;
56        }
57    }
58}
59
60pub struct Receiver<'a, T, const N: usize> {
61    parent: &'a Spsc<T, N>,
62    recver: ach::Receiver<'a, T, N>,
63}
64impl<'a, T: Unpin, const N: usize> Receiver<'a, T, N> {
65    pub fn try_recv(&mut self) -> Option<T> {
66        self.recver.try_recv().inspect(|_x| {
67            self.parent.consumer.notify_one();
68        })
69    }
70    pub async fn recv(&mut self) -> T {
71        let mut wait_p = self.parent.producer.listen();
72        loop {
73            if let Some(v) = self.try_recv() {
74                break v;
75            } else {
76                wait_p.next().await;
77            }
78        }
79    }
80}