Skip to main content

onlyne_client/backend/
outcome.rs

1use super::*;
2
3/// The queue behind one backend's outcome stream, and the flag its consumer
4/// sleeps on. Every handle shares it.
5#[derive(Default)]
6struct OutcomeQueue {
7    items: parking_lot::Mutex<VecDeque<SessionOutcome>>,
8    arrival: parking_lot::Condvar,
9}
10
11/// Producer half of a backend's outcome stream, held by the backend.
12#[derive(Clone, Default)]
13pub struct OutcomeSink {
14    queue: Arc<OutcomeQueue>,
15}
16
17impl OutcomeSink {
18    /// Record one terminal fact for the next consumer that asks.
19    pub fn push(&self, outcome: SessionOutcome) {
20        self.queue.items.lock().push_back(outcome);
21        self.queue.arrival.notify_all();
22    }
23}
24
25/// Receiver side of a backend's own outcome stream.
26///
27/// Cloning hands out another view of the same queue, and a view takes each fact
28/// at most once, so exactly one consumer settles a given task however many
29/// drains are alive. That is why this is a queue and not a channel: a backend
30/// outlives its first drain, and a receiver handed to a consumer that then died
31/// would leave the stream un-drainable.
32#[derive(Clone, Default)]
33pub struct OutcomeFeed {
34    queue: Arc<OutcomeQueue>,
35}
36
37impl OutcomeFeed {
38    /// The paired ends of one outcome stream.
39    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    /// The next fact, without waiting.
50    pub fn try_recv(&self) -> Option<SessionOutcome> {
51        self.queue.items.lock().pop_front()
52    }
53
54    /// The next fact, waiting at most `timeout`.
55    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}