async_ringbuf/traits/
consumer.rs1use core::{
2 future::Future,
3 pin::Pin,
4 task::{Context, Poll, Waker},
5};
6use futures_util::future::FusedFuture;
7use ringbuf::traits::Consumer;
8#[cfg(feature = "std")]
9use std::io;
10
11pub trait AsyncConsumer: Consumer {
12 fn register_waker(&self, waker: &Waker);
13
14 fn close(&mut self);
15 fn is_closed(&self) -> bool {
17 !self.write_is_held()
18 }
19
20 fn pop(&mut self) -> PopFuture<'_, Self> {
30 PopFuture { owner: self, done: false }
31 }
32
33 fn wait_occupied(&mut self, count: usize) -> WaitOccupiedFuture<'_, Self> {
43 debug_assert!(count <= self.capacity().get());
44 WaitOccupiedFuture {
45 owner: self,
46 count,
47 done: false,
48 }
49 }
50
51 fn pop_exact<'a: 'b, 'b>(&'a mut self, slice: &'b mut [Self::Item]) -> PopSliceFuture<'a, 'b, Self>
62 where
63 Self::Item: Copy,
64 {
65 PopSliceFuture {
66 owner: self,
67 slice: Some(slice),
68 count: 0,
69 }
70 }
71
72 #[cfg(feature = "alloc")]
78 fn pop_until_end<'a: 'b, 'b>(&'a mut self, vec: &'b mut alloc::vec::Vec<Self::Item>) -> PopVecFuture<'a, 'b, Self> {
79 PopVecFuture {
80 owner: self,
81 vec: Some(vec),
82 }
83 }
84
85 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>>
87 where
88 Self: Unpin,
89 {
90 let mut waker_registered = false;
91 loop {
92 let closed = self.is_closed();
93 if let Some(item) = self.try_pop() {
94 break Poll::Ready(Some(item));
95 }
96 if closed {
97 break Poll::Ready(None);
98 }
99 if waker_registered {
100 break Poll::Pending;
101 }
102 self.register_waker(cx.waker());
103 waker_registered = true;
104 }
105 }
106
107 #[cfg(feature = "std")]
109 fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut [u8]) -> Poll<io::Result<usize>>
110 where
111 Self: AsyncConsumer<Item = u8> + Unpin,
112 {
113 let mut waker_registered = false;
114 loop {
115 let closed = self.is_closed();
116 let len = self.pop_slice(buf);
117 if len != 0 || closed {
118 break Poll::Ready(Ok(len));
119 }
120 if waker_registered {
121 break Poll::Pending;
122 }
123 self.register_waker(cx.waker());
124 waker_registered = true;
125 }
126 }
127}
128
129#[must_use = "futures do nothing unless you `.await` or poll them"]
133pub struct PopFuture<'a, A: AsyncConsumer + ?Sized> {
134 owner: &'a mut A,
135 done: bool,
136}
137impl<A: AsyncConsumer> Unpin for PopFuture<'_, A> {}
138impl<A: AsyncConsumer> FusedFuture for PopFuture<'_, A> {
139 fn is_terminated(&self) -> bool {
140 self.done || self.owner.is_closed()
141 }
142}
143impl<A: AsyncConsumer> Future for PopFuture<'_, A> {
144 type Output = Option<A::Item>;
145
146 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
147 let mut waker_registered = false;
148 loop {
149 assert!(!self.done);
150 let closed = self.owner.is_closed();
151 if let Some(item) = self.owner.try_pop() {
152 self.done = true;
153 break Poll::Ready(Some(item));
154 }
155 if closed {
156 break Poll::Ready(None);
157 }
158 if waker_registered {
159 break Poll::Pending;
160 }
161 self.owner.register_waker(cx.waker());
162 waker_registered = true;
163 }
164 }
165}
166
167#[must_use = "futures do nothing unless you `.await` or poll them"]
171pub struct PopSliceFuture<'a, 'b, A: AsyncConsumer + ?Sized>
172where
173 A::Item: Copy,
174{
175 owner: &'a mut A,
176 slice: Option<&'b mut [A::Item]>,
177 count: usize,
178}
179impl<A: AsyncConsumer> Unpin for PopSliceFuture<'_, '_, A> where A::Item: Copy {}
180impl<A: AsyncConsumer> FusedFuture for PopSliceFuture<'_, '_, A>
181where
182 A::Item: Copy,
183{
184 fn is_terminated(&self) -> bool {
185 self.slice.is_none()
186 }
187}
188impl<A: AsyncConsumer> Future for PopSliceFuture<'_, '_, A>
189where
190 A::Item: Copy,
191{
192 type Output = Result<(), usize>;
193
194 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
195 let mut waker_registered = false;
196 loop {
197 let closed = self.owner.is_closed();
198 let mut slice = self.slice.take().unwrap();
199 let len = self.owner.pop_slice(slice);
200 slice = &mut slice[len..];
201 self.count += len;
202 if slice.is_empty() {
203 break Poll::Ready(Ok(()));
204 }
205 if closed {
206 break Poll::Ready(Err(self.count));
207 }
208 self.slice.replace(slice);
209 if waker_registered {
210 break Poll::Pending;
211 }
212 self.owner.register_waker(cx.waker());
213 waker_registered = true;
214 }
215 }
216}
217impl<A: AsyncConsumer> PopSliceFuture<'_, '_, A>
218where
219 A::Item: Copy,
220{
221 pub fn count(&self) -> usize {
223 self.count
224 }
225}
226
227#[cfg(feature = "alloc")]
231#[must_use = "futures do nothing unless you `.await` or poll them"]
232pub struct PopVecFuture<'a, 'b, A: AsyncConsumer + ?Sized> {
233 owner: &'a mut A,
234 vec: Option<&'b mut alloc::vec::Vec<A::Item>>,
235}
236#[cfg(feature = "alloc")]
237impl<A: AsyncConsumer> Unpin for PopVecFuture<'_, '_, A> {}
238#[cfg(feature = "alloc")]
239impl<A: AsyncConsumer> FusedFuture for PopVecFuture<'_, '_, A> {
240 fn is_terminated(&self) -> bool {
241 self.vec.is_none()
242 }
243}
244#[cfg(feature = "alloc")]
245impl<A: AsyncConsumer> Future for PopVecFuture<'_, '_, A> {
246 type Output = ();
247
248 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
249 let mut waker_registered = false;
250 loop {
251 let closed = self.owner.is_closed();
252 let vec = self.vec.take().unwrap();
253
254 loop {
255 if vec.len() == vec.capacity() {
256 vec.reserve(vec.capacity().max(16));
257 }
258 let n = self.owner.pop_slice_uninit(vec.spare_capacity_mut());
259 if n == 0 {
260 break;
261 }
262 unsafe { vec.set_len(vec.len() + n) };
263 }
264
265 if closed {
266 break Poll::Ready(());
267 }
268 self.vec.replace(vec);
269 if waker_registered {
270 break Poll::Pending;
271 }
272 self.owner.register_waker(cx.waker());
273 waker_registered = true;
274 }
275 }
276}
277
278#[must_use = "futures do nothing unless you `.await` or poll them"]
282pub struct WaitOccupiedFuture<'a, A: AsyncConsumer + ?Sized> {
283 owner: &'a A,
284 count: usize,
285 done: bool,
286}
287impl<A: AsyncConsumer> Unpin for WaitOccupiedFuture<'_, A> {}
288impl<A: AsyncConsumer> FusedFuture for WaitOccupiedFuture<'_, A> {
289 fn is_terminated(&self) -> bool {
290 self.done
291 }
292}
293impl<A: AsyncConsumer> Future for WaitOccupiedFuture<'_, A> {
294 type Output = ();
295
296 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
297 let mut waker_registered = false;
298 loop {
299 assert!(!self.done);
300 let closed = self.owner.is_closed();
301 if self.count <= self.owner.occupied_len() || closed {
302 self.done = true;
303 break Poll::Ready(());
304 }
305 if waker_registered {
306 break Poll::Pending;
307 }
308 self.owner.register_waker(cx.waker());
309 waker_registered = true;
310 }
311 }
312}