Skip to main content

async_ringbuf/traits/
consumer.rs

1use 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    /// Whether the corresponding producer was closed.
16    fn is_closed(&self) -> bool {
17        !self.write_is_held()
18    }
19
20    /// Pop item from the ring buffer waiting asynchronously if the buffer is empty.
21    ///
22    /// Future returns:
23    /// + `Some(item)` - an item is taken.
24    /// + `None` - the buffer is empty and the corresponding producer was dropped.
25    ///
26    /// # Cancel safety
27    ///
28    /// If future is cancelled then no item removed from the ring buffer.
29    fn pop(&mut self) -> PopFuture<'_, Self> {
30        PopFuture { owner: self, done: false }
31    }
32
33    /// Wait for the buffer to contain at least `count` items or to close.
34    ///
35    /// In debug mode panics if `count` is greater than buffer capacity.
36    ///
37    /// The method takes `&mut self` because only single [`WaitOccupiedFuture`] is allowed at a time.
38    ///
39    /// # Cancel safety
40    ///
41    /// The future can be safely cancelled.
42    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    /// Fill slice with items from the ring buffer waiting asynchronously until slice filled or corresponding producer closed.
52    ///
53    /// Future returns:
54    /// + `Ok` - the whole slice is filled with the items from the buffer.
55    /// + `Err(count)` - the buffer is empty and the corresponding producer was dropped, number of items copied to slice is returned.
56    ///
57    /// # Cancel safety
58    ///
59    /// If future is cancelled then slice can be partially filled.
60    /// The number of items already copied can be examined by [`PopSliceFuture::count`].
61    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    /// Fill `vec` with items from the ring buffer waiting asynchronously until corresponding producer closed.
73    ///
74    /// # Cancel safety
75    ///
76    /// If future is cancelled then `vec` contains items taken from RB before cancellation.
77    #[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    /// Poll for the next item in the ring buffer.
86    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    /// Poll reading bytes from byte buffer.
108    #[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/// # Cancel safety
130///
131/// If future is cancelled then no item removed from the ring buffer.
132#[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/// # Cancel safety
168///
169/// If future is cancelled then slice can be partially filled.
170#[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    /// Number of items already copied from the ring buufer to the slice provided.
222    pub fn count(&self) -> usize {
223        self.count
224    }
225}
226
227/// # Cancel safety
228///
229/// If future is cancelled then `vec` contains items taken from RB before cancellation.
230#[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/// # Cancel safety
279///
280/// The future can be safely cancelled.
281#[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}