async_ach_spsc/
heapless.rs1use 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}