1use crate::{HeadTailBuffer, OutputChunk};
2use tokio::sync::broadcast;
3use tokio::sync::broadcast::error::{RecvError, TryRecvError};
4
5pub 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 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 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 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}