use std::sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicUsize},
};
use crate::{
bus::{Bus, BusReceiver},
clock::Clock,
control::{self, ControlReceiver, ControlSender},
element::{Context, SourceElement, element_pp_log, pipeline_pp_log},
error::Result,
graph::{ElementId, PipelineGraph},
playback_clock::PlaybackClock,
};
use super::Pipeline;
pub(super) type SourceEntry = (ElementId, Box<dyn SourceElement>);
pub struct PipelineBuilder {
id: Arc<str>,
bus: Bus,
bus_rx: BusReceiver,
clock: Arc<Clock>,
playback_clock: Arc<PlaybackClock>,
graph: PipelineGraph,
sources: Vec<SourceEntry>,
control_pairs: Vec<(ControlSender, ControlReceiver)>,
}
impl PipelineBuilder {
pub fn new(id: impl Into<String>) -> Self {
let id: Arc<str> = id.into().into();
let (bus, bus_rx) = Bus::new();
let clock = Arc::new(Clock::new());
Self {
id,
bus,
bus_rx,
playback_clock: Arc::new(PlaybackClock::new(clock.clone())),
clock,
graph: PipelineGraph::new(),
sources: Vec::new(),
control_pairs: Vec::new(),
}
}
pub fn add_source<S: SourceElement + 'static>(
mut self,
mut source: S,
wire: impl FnOnce(&mut S, &Arc<Context>) -> Result<()>,
) -> Result<Self> {
*source.pp_log_mut() =
element_pp_log(source.element_type(), &source.name(), Some(&self.id));
let source_id = self.graph.add_source(source.element_type(), source.name());
let context = Arc::new(Context {
bus: self.bus.clone(),
pipeline_id: self.id.clone(),
graph: self.graph.clone(),
clock: self.clock.clone(),
playback_clock: self.playback_clock.clone(),
source_id,
});
wire(&mut source, &context)?;
self.sources.push((source_id, Box::new(source)));
self.control_pairs.push(control::channel());
Ok(self)
}
pub fn build(self) -> Arc<Pipeline> {
assert!(
!self.sources.is_empty(),
"PipelineBuilder::build called with no sources added"
);
let (control_txs, control_rxs): (Vec<_>, Vec<_>) = self.control_pairs.into_iter().unzip();
let pp_log = pipeline_pp_log(&self.id);
Arc::new(Pipeline {
id: self.id,
pp_log,
sources: Mutex::new(Some(self.sources)),
bus: Mutex::new(Some(self.bus)),
control_txs,
control_rxs: Mutex::new(Some(control_rxs)),
clock: self.clock,
playback_clock: self.playback_clock,
bus_rx: self.bus_rx,
running: Arc::new(AtomicUsize::new(0)),
paused: AtomicBool::new(false),
workers: Mutex::new(Vec::new()),
graph: self.graph,
})
}
}