use std::{sync::Arc, time::Duration};
use crate::pp_log::{PpLog, pp_error, pp_info, pp_warn};
use crossbeam_channel::{Receiver, Sender, unbounded};
use crate::{element::ElementType, error::Error, graph::ElementId};
#[derive(Debug)]
#[non_exhaustive]
pub enum BusEvent {
Eos {
element_type: ElementType,
name: Arc<str>,
},
Error {
element_type: ElementType,
name: Arc<str>,
error: Error,
},
Dropped {
element_type: ElementType,
name: Arc<str>,
},
Seeked {
element_type: ElementType,
name: Arc<str>,
requested: Duration,
landed: Duration,
},
}
#[derive(Clone)]
pub struct Bus {
tx: Sender<BusMessage>,
element_id: Option<ElementId>,
}
pub struct BusReceiver {
rx: Receiver<BusMessage>,
}
#[derive(Debug)]
pub struct BusMessage {
pub element_id: Option<ElementId>,
pub event: BusEvent,
}
impl Bus {
pub fn new() -> (Bus, BusReceiver) {
let (tx, rx) = unbounded();
(
Bus {
tx,
element_id: None,
},
BusReceiver { rx },
)
}
pub(crate) fn for_element(&self, element_id: ElementId) -> Bus {
Bus {
tx: self.tx.clone(),
element_id: Some(element_id),
}
}
pub fn post(&self, pp_log: &PpLog, event: BusEvent) {
match &event {
BusEvent::Eos { .. } => {
pp_info!(pp_log: pp_log, "event=eos phase=reported")
}
BusEvent::Error { error, .. } => pp_error!(pp_log: pp_log, "{error}"),
BusEvent::Dropped { .. } => {
pp_warn!(pp_log: pp_log, "dropped a buffer (queue full)")
}
BusEvent::Seeked {
requested, landed, ..
} => pp_info!(pp_log: pp_log, "seeked: requested {requested:.2?}, landed {landed:.2?}"),
}
let _ = self.tx.send(BusMessage {
element_id: self.element_id,
event,
});
}
}
impl BusReceiver {
pub fn recv(&self) -> Option<BusEvent> {
self.recv_message().map(|message| message.event)
}
pub fn try_recv(&self) -> Option<BusEvent> {
self.try_recv_message().map(|message| message.event)
}
pub fn iter(&self) -> impl Iterator<Item = BusEvent> + '_ {
self.iter_with_ids().map(|message| message.event)
}
pub fn recv_message(&self) -> Option<BusMessage> {
self.rx.recv().ok()
}
pub fn try_recv_message(&self) -> Option<BusMessage> {
self.rx.try_recv().ok()
}
pub fn iter_with_ids(&self) -> impl Iterator<Item = BusMessage> + '_ {
self.rx.iter()
}
pub fn log_events(&self) {
for event in self.iter() {
match event {
BusEvent::Error { name, error, .. } => eprintln!("[{name}] error: {error}"),
BusEvent::Eos { name, .. } => println!("[{name}] eos"),
BusEvent::Dropped { name, .. } => {
eprintln!("[{name}] dropped a buffer (queue full)")
}
BusEvent::Seeked {
name,
requested,
landed,
..
} => println!("[{name}] seeked: requested {requested:.2?}, landed {landed:.2?}"),
}
}
}
}