tracexec_core/primitives/
local_chan.rs1use std::{
9 cell::RefCell,
10 collections::VecDeque,
11 rc::Rc,
12};
13
14pub fn unbounded<T>() -> (LocalUnboundedSender<T>, LocalUnboundedReceiver<T>) {
15 let inner = Rc::new(RefCell::new(VecDeque::new()));
16 (
17 LocalUnboundedSender {
18 inner: inner.clone(),
19 },
20 LocalUnboundedReceiver { inner },
21 )
22}
23
24#[derive(Debug, Clone)]
25pub struct LocalUnboundedSender<T> {
26 inner: Rc<RefCell<VecDeque<T>>>,
27}
28
29#[derive(Debug, Clone)]
30pub struct LocalUnboundedReceiver<T> {
31 inner: Rc<RefCell<VecDeque<T>>>,
32}
33
34impl<T> LocalUnboundedSender<T> {
35 pub fn send(&self, v: T) {
36 self.inner.borrow_mut().push_back(v);
37 }
38}
39impl<T> LocalUnboundedReceiver<T> {
40 pub fn receive(&self) -> Option<T> {
41 self.inner.borrow_mut().pop_front()
42 }
43}
44
45#[cfg(test)]
46mod tests {
47 use test_that::prelude::*;
48
49 use super::*;
50
51 #[test]
52 fn test_basic_send_receive() {
53 let (sender, receiver) = unbounded();
54
55 sender.send(42);
56 sender.send(100);
57
58 assert_that!(receiver.receive(), some(eq(42)));
59 assert_that!(receiver.receive(), some(eq(100)));
60 assert_that!(receiver.receive(), none());
62 }
63
64 #[test]
65 fn test_receive_empty_channel() {
66 let (_sender, receiver): (LocalUnboundedSender<i32>, _) = unbounded();
67 assert_that!(receiver.receive(), none());
68 }
69
70 #[test]
71 fn test_multiple_senders_receive() {
72 let (sender1, receiver1) = unbounded();
73 let sender2 = sender1.clone();
74 let receiver2 = receiver1.clone();
75
76 sender1.send(1);
77 sender2.send(2);
78
79 assert_that!(receiver1.receive(), some(eq(1)));
81 assert_that!(receiver2.receive(), some(eq(2)));
82
83 assert_that!(receiver1.receive(), none());
85 assert_that!(receiver2.receive(), none());
86 }
87
88 #[test]
89 fn test_fifo_order_with_multiple_elements() {
90 let (sender, receiver) = unbounded();
91
92 for i in 0..10 {
93 sender.send(i);
94 }
95
96 for i in 0..10 {
97 assert_that!(receiver.receive(), some(eq(i)));
98 }
99
100 assert_that!(receiver.receive(), none());
101 }
102
103 #[test]
104 fn test_clone_sender_receiver_shares_queue() {
105 let (sender, receiver) = unbounded();
106 let sender2 = sender.clone();
107 let receiver2 = receiver.clone();
108
109 sender.send(1);
110 sender2.send(2);
111
112 assert_that!(receiver2.receive(), some(eq(1)));
113 assert_that!(receiver.receive(), some(eq(2)));
114 assert_that!(receiver2.receive(), none());
115 }
116}