osal_rs/async_primitives/
queue.rs1use core::future::Future;
29use core::pin::Pin;
30use core::task::{Context, Poll};
31use core::time::Duration;
32
33
34use alloc::sync::Arc;
35
36use crate::os::types::{UBaseType, TickType};
37use crate::os::Queue;
38use crate::traits::QueueFn;
39use crate::utils::{Error, Result};
40
41use super::waker_slot::WakerSlot;
42
43pub struct AsyncQueue {
48 inner: Arc<Queue>,
49 rx_waker: Arc<WakerSlot>,
50 tx_waker: Arc<WakerSlot>,
51}
52
53impl AsyncQueue {
54 pub fn new(size: u32, message_size: u32) -> Result<Self> {
63 let inner = Queue::new(size as UBaseType, message_size as UBaseType)?;
64 Ok(Self {
65 inner: Arc::new(inner),
66 rx_waker: Arc::new(WakerSlot::new()),
67 tx_waker: Arc::new(WakerSlot::new()),
68 })
69 }
70
71 pub fn post(&self, item: &[u8], timeout_ms: u64) -> Result<()> {
73 let ticks = TickType::try_from(
74 Duration::from_millis(timeout_ms)
75 .as_millis()
76 .try_into()
77 .unwrap_or(TickType::MAX),
78 )
79 .unwrap_or(TickType::MAX);
80 let result = self.inner.post(item, ticks);
81 if result.is_ok() {
82 self.rx_waker.wake();
83 }
84 result
85 }
86
87 pub fn fetch(&self, buf: &mut [u8], timeout_ms: u64) -> Result<()> {
89 let ticks = TickType::try_from(
90 Duration::from_millis(timeout_ms)
91 .as_millis()
92 .try_into()
93 .unwrap_or(TickType::MAX),
94 )
95 .unwrap_or(TickType::MAX);
96 let result = self.inner.fetch(buf, ticks);
97 if result.is_ok() {
98 self.tx_waker.wake();
99 }
100 result
101 }
102
103 pub fn fetch_async<'a>(&'a self, buf: &'a mut [u8]) -> FetchFuture<'a> {
105 FetchFuture {
106 queue: self,
107 buf,
108 }
109 }
110
111 pub fn post_async<'a>(&'a self, item: &'a [u8]) -> PostFuture<'a> {
113 PostFuture {
114 queue: self,
115 item,
116 }
117 }
118}
119
120pub struct FetchFuture<'a> {
122 queue: &'a AsyncQueue,
123 buf: &'a mut [u8],
124}
125
126impl<'a> Future for FetchFuture<'a> {
127 type Output = Result<()>;
128
129 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
130 let this = self.get_mut();
131
132 match this.queue.inner.fetch(this.buf, 0) {
134 Ok(()) => {
135 this.queue.tx_waker.wake();
136 return Poll::Ready(Ok(()));
137 }
138 Err(Error::Timeout) => {}
139 Err(e) => return Poll::Ready(Err(e)),
140 }
141
142 this.queue.rx_waker.store(cx.waker());
144
145 match this.queue.inner.fetch(this.buf, 0) {
146 Ok(()) => {
147 this.queue.tx_waker.wake();
148 Poll::Ready(Ok(()))
149 }
150 Err(Error::Timeout) => Poll::Pending,
151 Err(e) => Poll::Ready(Err(e)),
152 }
153 }
154}
155
156pub struct PostFuture<'a> {
158 queue: &'a AsyncQueue,
159 item: &'a [u8],
160}
161
162impl<'a> Future for PostFuture<'a> {
163 type Output = Result<()>;
164
165 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
166 let this = self.get_mut();
167
168 match this.queue.inner.post(this.item, 0) {
169 Ok(()) => {
170 this.queue.rx_waker.wake();
171 return Poll::Ready(Ok(()));
172 }
173 Err(Error::Timeout) => {}
174 Err(e) => return Poll::Ready(Err(e)),
175 }
176
177 this.queue.tx_waker.store(cx.waker());
178
179 match this.queue.inner.post(this.item, 0) {
180 Ok(()) => {
181 this.queue.rx_waker.wake();
182 Poll::Ready(Ok(()))
183 }
184 Err(Error::Timeout) => Poll::Pending,
185 Err(e) => Poll::Ready(Err(e)),
186 }
187 }
188}