Skip to main content

ante_exec/
receiver.rs

1use crate::{HeadTailBuffer, OutputChunk};
2use tokio::sync::broadcast;
3use tokio::sync::broadcast::error::{RecvError, TryRecvError};
4
5/// Ergonomic wrapper around a broadcast receiver of process output.
6pub struct OutputReceiver {
7    rx: broadcast::Receiver<OutputChunk>,
8}
9
10impl OutputReceiver {
11    pub fn new(rx: broadcast::Receiver<OutputChunk>) -> Self {
12        Self { rx }
13    }
14
15    /// Receive the next available chunk, skipping lagged messages.
16    pub async fn recv(&mut self) -> Option<OutputChunk> {
17        loop {
18            match self.rx.recv().await {
19                Ok(chunk) => return Some(chunk),
20                Err(RecvError::Lagged(_)) => continue,
21                Err(RecvError::Closed) => return None,
22            }
23        }
24    }
25
26    /// Drain all immediately available chunks without blocking.
27    pub fn try_drain(&mut self) -> Vec<OutputChunk> {
28        let mut drained = Vec::new();
29        self.drain_with(|chunk| drained.push(chunk));
30        drained
31    }
32
33    pub(crate) fn drain_with<F>(&mut self, mut on_chunk: F)
34    where
35        F: FnMut(OutputChunk),
36    {
37        loop {
38            match self.rx.try_recv() {
39                Ok(chunk) => on_chunk(chunk),
40                Err(TryRecvError::Lagged(_)) => continue,
41                Err(TryRecvError::Empty | TryRecvError::Closed) => return,
42            }
43        }
44    }
45
46    /// Drain all immediately available chunks into a buffer, merging streams.
47    pub fn drain_into(&mut self, buf: &mut HeadTailBuffer) {
48        self.drain_with(|chunk| buf.push_chunk(chunk.data));
49    }
50}
51
52impl From<broadcast::Receiver<OutputChunk>> for OutputReceiver {
53    fn from(rx: broadcast::Receiver<OutputChunk>) -> Self {
54        Self::new(rx)
55    }
56}