1use std::cell::RefCell;
5use std::collections::VecDeque;
6use std::future::Future;
7use std::pin::Pin;
8use std::rc::Rc;
9use std::task::{Context, Poll, Waker};
10
11struct QInner<T> {
12 buf: RefCell<VecDeque<T>>,
13 cap: Option<usize>, get_waiters: RefCell<Vec<Waker>>,
15 put_waiters: RefCell<Vec<Waker>>,
16}
17
18impl<T> QInner<T> {
19 fn has_space(&self) -> bool {
20 match self.cap {
21 None => true,
22 Some(c) => self.buf.borrow().len() < c,
23 }
24 }
25 fn wake_getters(&self) {
26 for w in self.get_waiters.borrow_mut().drain(..) {
27 w.wake();
28 }
29 }
30 fn wake_putters(&self) {
31 for w in self.put_waiters.borrow_mut().drain(..) {
32 w.wake();
33 }
34 }
35}
36
37pub struct Queue<T> {
39 inner: Rc<QInner<T>>,
40}
41
42impl<T> Clone for Queue<T> {
43 fn clone(&self) -> Self {
44 Queue { inner: self.inner.clone() }
45 }
46}
47
48impl<T> Queue<T> {
49 pub fn new(cap: Option<usize>) -> Queue<T> {
51 Queue {
52 inner: Rc::new(QInner {
53 buf: RefCell::new(VecDeque::new()),
54 cap,
55 get_waiters: RefCell::new(Vec::new()),
56 put_waiters: RefCell::new(Vec::new()),
57 }),
58 }
59 }
60
61 pub fn unbounded() -> Queue<T> {
62 Self::new(None)
63 }
64
65 pub fn len(&self) -> usize {
66 self.inner.buf.borrow().len()
67 }
68 pub fn is_empty(&self) -> bool {
69 self.inner.buf.borrow().is_empty()
70 }
71
72 pub fn has_space(&self) -> bool {
75 self.inner.has_space()
76 }
77
78 pub fn try_put(&self, item: T) -> Result<(), T> {
79 if self.inner.has_space() {
80 self.inner.buf.borrow_mut().push_back(item);
81 self.inner.wake_getters();
82 Ok(())
83 } else {
84 Err(item)
85 }
86 }
87
88 pub fn try_get(&self) -> Option<T> {
89 let item = self.inner.buf.borrow_mut().pop_front();
90 if item.is_some() {
91 self.inner.wake_putters();
92 }
93 item
94 }
95
96 pub fn put(&self, item: T) -> Put<T> {
97 Put { inner: self.inner.clone(), item: Some(item) }
98 }
99
100 pub fn get(&self) -> Get<T> {
101 Get { inner: self.inner.clone() }
102 }
103
104 pub fn wait_for_space(&self) -> Space<T> {
117 Space { inner: self.inner.clone() }
118 }
119}
120
121pub struct Space<T> {
123 inner: Rc<QInner<T>>,
124}
125
126impl<T> Future for Space<T> {
127 type Output = ();
128 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
129 if self.inner.has_space() {
130 Poll::Ready(())
131 } else {
132 self.inner.put_waiters.borrow_mut().push(cx.waker().clone());
133 Poll::Pending
134 }
135 }
136}
137
138impl<T: Clone> Queue<T> {
139 pub fn try_peek(&self) -> Option<T> {
145 self.inner.buf.borrow().front().cloned()
146 }
147
148 pub fn peek(&self) -> Peek<T> {
150 Peek { inner: self.inner.clone() }
151 }
152}
153
154pub struct Peek<T> {
155 inner: Rc<QInner<T>>,
156}
157
158impl<T: Clone> Future for Peek<T> {
159 type Output = T;
160 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<T> {
161 let item = self.inner.buf.borrow().front().cloned();
162 match item {
163 Some(v) => Poll::Ready(v),
165 None => {
166 self.inner.get_waiters.borrow_mut().push(cx.waker().clone());
167 Poll::Pending
168 }
169 }
170 }
171}
172
173pub struct Put<T> {
174 inner: Rc<QInner<T>>,
175 item: Option<T>,
176}
177
178impl<T> Future for Put<T> {
179 type Output = ();
180 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
181 let this = unsafe { self.get_unchecked_mut() };
185 if this.inner.has_space() {
186 let item = this.item.take().expect("Put polled after completion");
187 this.inner.buf.borrow_mut().push_back(item);
188 this.inner.wake_getters();
189 Poll::Ready(())
190 } else {
191 this.inner.put_waiters.borrow_mut().push(cx.waker().clone());
192 Poll::Pending
193 }
194 }
195}
196
197pub struct Get<T> {
198 inner: Rc<QInner<T>>,
199}
200
201impl<T> Future for Get<T> {
202 type Output = T;
203 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<T> {
204 let item = self.inner.buf.borrow_mut().pop_front();
205 match item {
206 Some(v) => {
207 self.inner.wake_putters();
208 Poll::Ready(v)
209 }
210 None => {
211 self.inner.get_waiters.borrow_mut().push(cx.waker().clone());
212 Poll::Pending
213 }
214 }
215 }
216}
217
218#[cfg(test)]
223mod tests {
224 use super::*;
225 use crate::testing::{assert_pending, block_on};
226
227 #[test]
228 fn put_get_round_trips_in_order() {
229 block_on(async {
230 let q: Queue<u8> = Queue::unbounded();
231 for n in 0..5 {
232 q.put(n).await;
233 }
234 for n in 0..5 {
235 assert_eq!(q.get().await, n, "a queue is FIFO");
236 }
237 });
238 }
239
240 #[test]
243 fn try_put_on_a_full_queue_returns_the_item() {
244 block_on(async {
245 let q: Queue<String> = Queue::new(Some(1));
246 assert!(q.try_put(String::from("first")).is_ok());
247 match q.try_put(String::from("second")) {
248 Err(back) => assert_eq!(back, "second", "the item comes home"),
249 Ok(()) => panic!("a depth-1 queue accepted a second item"),
250 }
251 });
252 }
253
254 #[test]
255 fn try_get_on_an_empty_queue_is_none() {
256 block_on(async {
257 let q: Queue<u8> = Queue::unbounded();
258 assert!(q.try_get().is_none());
259 q.put(1).await;
260 assert_eq!(q.try_get(), Some(1));
261 assert!(q.try_get().is_none(), "and it was consumed");
262 });
263 }
264
265 #[test]
266 fn a_blocked_get_wakes_on_a_put() {
267 block_on(async {
268 let q: Queue<u8> = Queue::unbounded();
269 let producer = q.clone();
270 crate::executor::spawn(async move {
271 producer.put(42).await;
272 });
273 assert_eq!(q.get().await, 42, "the get was waiting before the put");
274 });
275 }
276
277 #[test]
278 fn a_blocked_put_wakes_when_space_appears() {
279 block_on(async {
280 let q: Queue<u8> = Queue::new(Some(1));
281 q.put(1).await;
282 let consumer = q.clone();
283 crate::executor::spawn(async move {
284 let _ = consumer.get().await;
285 });
286 q.put(2).await; assert_eq!(q.len(), 1);
288 });
289 }
290
291 #[test]
292 fn a_get_on_an_empty_queue_waits() {
293 let q: Queue<u8> = Queue::unbounded();
294 assert_pending(async move { q.get().await });
295 }
296
297 #[test]
300 fn peek_does_not_consume() {
301 block_on(async {
302 let q: Queue<u8> = Queue::unbounded();
303 q.put(9).await;
304 assert_eq!(q.peek().await, 9);
305 assert_eq!(q.len(), 1, "peek left it there");
306 assert_eq!(q.get().await, 9, "and get takes the same item");
307 });
308 }
309
310 #[test]
311 fn try_peek_is_none_when_empty() {
312 block_on(async {
313 let q: Queue<u8> = Queue::unbounded();
314 assert!(q.try_peek().is_none());
315 q.put(3).await;
316 assert_eq!(q.try_peek(), Some(3));
317 assert_eq!(q.len(), 1);
318 });
319 }
320
321 #[test]
322 fn has_space_agrees_with_len() {
323 block_on(async {
324 let q: Queue<u8> = Queue::new(Some(2));
325 assert!(q.has_space());
326 q.put(1).await;
327 assert!(q.has_space());
328 q.put(2).await;
329 assert!(!q.has_space(), "full at its declared depth");
330 assert_eq!(q.len(), 2);
331 });
332 }
333
334 #[test]
335 fn wait_for_space_returns_when_drained() {
336 block_on(async {
337 let q: Queue<u8> = Queue::new(Some(1));
338 q.put(1).await;
339 let consumer = q.clone();
340 crate::executor::spawn(async move {
341 let _ = consumer.get().await;
342 });
343 q.wait_for_space().await;
344 assert!(q.has_space());
345 });
346 }
347
348 #[test]
349 fn unbounded_never_blocks_a_put() {
350 block_on(async {
351 let q: Queue<u32> = Queue::unbounded();
352 for n in 0..1000 {
353 assert!(q.try_put(n).is_ok(), "an unbounded queue always has room");
354 }
355 assert_eq!(q.len(), 1000);
356 });
357 }
358
359 #[test]
360 fn is_empty_tracks_contents() {
361 block_on(async {
362 let q: Queue<u8> = Queue::unbounded();
363 assert!(q.is_empty());
364 q.put(1).await;
365 assert!(!q.is_empty());
366 let _ = q.get().await;
367 assert!(q.is_empty());
368 });
369 }
370
371 #[test]
374 fn a_dropped_get_consumes_nothing() {
375 block_on(async {
376 let q: Queue<u8> = Queue::unbounded();
377 {
378 let pending = q.get();
379 drop(pending);
380 }
381 q.put(5).await;
382 assert_eq!(q.get().await, 5, "the item survived the dropped get");
383 });
384 }
385}