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