Line data Source code
1 : use std::{collections::VecDeque, sync::Arc};
2 :
3 : use parking_lot::{Mutex, MutexGuard};
4 :
5 : use crate::executor::{self, PollSome, Waker};
6 :
7 : /// FIFO channel with blocking send and receive. Can be cloned and shared between threads.
8 : /// Blocking functions should be used only from threads that are managed by the executor.
9 : pub struct Chan<T> {
10 : shared: Arc<State<T>>,
11 : }
12 :
13 : impl<T> Clone for Chan<T> {
14 4841278 : fn clone(&self) -> Self {
15 4841278 : Chan {
16 4841278 : shared: self.shared.clone(),
17 4841278 : }
18 4841278 : }
19 : }
20 :
21 : impl<T> Default for Chan<T> {
22 0 : fn default() -> Self {
23 0 : Self::new()
24 0 : }
25 : }
26 :
27 : impl<T> Chan<T> {
28 632330 : pub fn new() -> Chan<T> {
29 632330 : Chan {
30 632330 : shared: Arc::new(State {
31 632330 : queue: Mutex::new(VecDeque::new()),
32 632330 : waker: Waker::new(),
33 632330 : }),
34 632330 : }
35 632330 : }
36 :
37 : /// Get a message from the front of the queue, block if the queue is empty.
38 : /// If not called from the executor thread, it can block forever.
39 4270 : pub fn recv(&self) -> T {
40 4270 : self.shared.recv()
41 4270 : }
42 :
43 : /// Panic if the queue is empty.
44 518903 : pub fn must_recv(&self) -> T {
45 518903 : self.shared
46 518903 : .try_recv()
47 518903 : .expect("message should've been ready")
48 518903 : }
49 :
50 : /// Get a message from the front of the queue, return None if the queue is empty.
51 : /// Never blocks.
52 450263 : pub fn try_recv(&self) -> Option<T> {
53 450263 : self.shared.try_recv()
54 450263 : }
55 :
56 : /// Send a message to the back of the queue.
57 955851 : pub fn send(&self, t: T) {
58 955851 : self.shared.send(t);
59 955851 : }
60 : }
61 :
62 : struct State<T> {
63 : queue: Mutex<VecDeque<T>>,
64 : waker: Waker,
65 : }
66 :
67 : impl<T> State<T> {
68 955851 : fn send(&self, t: T) {
69 955851 : self.queue.lock().push_back(t);
70 955851 : self.waker.wake_all();
71 955851 : }
72 :
73 969166 : fn try_recv(&self) -> Option<T> {
74 969166 : let mut q = self.queue.lock();
75 969166 : q.pop_front()
76 969166 : }
77 :
78 4270 : fn recv(&self) -> T {
79 4270 : // interrupt the receiver to prevent consuming everything at once
80 4270 : executor::yield_me(0);
81 4270 :
82 4270 : let mut queue = self.queue.lock();
83 4270 : if let Some(t) = queue.pop_front() {
84 0 : return t;
85 4270 : }
86 : loop {
87 10686 : self.waker.wake_me_later();
88 10686 : if let Some(t) = queue.pop_front() {
89 3705 : return t;
90 6416 : }
91 6416 : MutexGuard::unlocked(&mut queue, || {
92 6416 : executor::yield_me(-1);
93 6416 : });
94 : }
95 3705 : }
96 : }
97 :
98 : impl<T> PollSome for Chan<T> {
99 : /// Schedules a wakeup for the current thread.
100 6057061 : fn wake_me(&self) {
101 6057061 : self.shared.waker.wake_me_later();
102 6057061 : }
103 :
104 : /// Checks if chan has any pending messages.
105 4885712 : fn has_some(&self) -> bool {
106 4885712 : !self.shared.queue.lock().is_empty()
107 4885712 : }
108 : }
|