onlyne_client/backend/
outcome.rs1use super::*;
2
3#[derive(Default)]
6struct OutcomeQueue {
7 items: parking_lot::Mutex<VecDeque<SessionOutcome>>,
8 arrival: parking_lot::Condvar,
9}
10
11#[derive(Clone, Default)]
13pub struct OutcomeSink {
14 queue: Arc<OutcomeQueue>,
15}
16
17impl OutcomeSink {
18 pub fn push(&self, outcome: SessionOutcome) {
20 self.queue.items.lock().push_back(outcome);
21 self.queue.arrival.notify_all();
22 }
23}
24
25#[derive(Clone, Default)]
33pub struct OutcomeFeed {
34 queue: Arc<OutcomeQueue>,
35}
36
37impl OutcomeFeed {
38 pub fn channel() -> (OutcomeSink, Self) {
40 let queue = Arc::new(OutcomeQueue::default());
41 (
42 OutcomeSink {
43 queue: Arc::clone(&queue),
44 },
45 Self { queue },
46 )
47 }
48
49 pub fn try_recv(&self) -> Option<SessionOutcome> {
51 self.queue.items.lock().pop_front()
52 }
53
54 pub fn recv_timeout(&self, timeout: Duration) -> Option<SessionOutcome> {
56 let deadline = Instant::now() + timeout;
57 let mut items = self.queue.items.lock();
58 while items.is_empty() {
59 let left = deadline.saturating_duration_since(Instant::now());
60 if self.queue.arrival.wait_for(&mut items, left).timed_out() && items.is_empty() {
61 return None;
62 }
63 }
64 items.pop_front()
65 }
66}