use std::{
sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
},
thread::{self, JoinHandle},
time::Duration,
};
use crate::pp_log::{PpLog, pp_info, pp_trace};
use crate::{
bus::{Bus, BusEvent, BusReceiver},
clock::Clock,
control::{ControlMsg, ControlReceiver, ControlSender},
element::{Context, SourceElement},
error::Result,
graph::{GraphSnapshot, NodeInfo, PipelineGraph, log_topology},
playback_clock::PlaybackClock,
};
use super::{PipelineBuilder, builder::SourceEntry};
pub struct Pipeline {
pub(super) id: Arc<str>,
pub(super) pp_log: PpLog,
pub(super) sources: Mutex<Option<Vec<SourceEntry>>>,
pub(super) bus: Mutex<Option<Bus>>,
pub(super) control_txs: Vec<ControlSender>,
pub(super) control_rxs: Mutex<Option<Vec<ControlReceiver>>>,
pub(super) clock: Arc<Clock>,
pub(super) playback_clock: Arc<PlaybackClock>,
pub(super) bus_rx: BusReceiver,
pub(super) running: Arc<AtomicUsize>,
pub(super) workers: Mutex<Vec<JoinHandle<()>>>,
pub(super) graph: PipelineGraph,
}
impl Pipeline {
pub fn new<S: SourceElement + 'static>(
id: impl Into<String>,
source: S,
wire: impl FnOnce(&mut S, &Arc<Context>) -> Result<()>,
) -> Result<Arc<Self>> {
Ok(PipelineBuilder::new(id).add_source(source, wire)?.build())
}
pub fn id(&self) -> &str {
&self.id
}
pub fn bus(&self) -> &BusReceiver {
&self.bus_rx
}
pub fn graph(&self) -> GraphSnapshot {
self.graph.snapshot()
}
pub fn elements(&self) -> Vec<NodeInfo> {
self.graph().nodes
}
pub fn topology(&self) -> String {
self.graph().topology()
}
pub fn clock(&self) -> &Arc<Clock> {
&self.clock
}
pub fn playback_clock(&self) -> &Arc<PlaybackClock> {
&self.playback_clock
}
pub fn run(&self) {
let Some(sources) = self.sources.lock().unwrap().take() else {
return;
};
let Some(bus) = self.bus.lock().unwrap().take() else {
return;
};
let Some(control_rxs) = self.control_rxs.lock().unwrap().take() else {
return;
};
if crate::log::enabled(crate::log::Level::Info) {
log_topology(&self.pp_log, "run", &self.graph());
}
self.running.store(sources.len(), Ordering::Release);
for ((source_id, source), control_rx) in sources.into_iter().zip(control_rxs) {
let bus = bus.for_element(source_id);
let running = Arc::clone(&self.running);
let handle = thread::Builder::new()
.name("pipeline:source".into())
.spawn(move || {
let mut source = source;
let control_rx = control_rx;
let _running = RunningSourceGuard::new(running);
let source_name = source.name();
let source_type = source.element_type();
let outcome = if let Err(error) = source.run(&control_rx, &bus) {
bus.post(
source.pp_log(),
BusEvent::Error {
element_type: source_type,
name: source_name.clone(),
error,
},
);
"error"
} else {
"ok"
};
pp_info!(pp_log: source.pp_log(), "finished outcome={outcome}");
})
.expect("failed to spawn pipeline source thread");
self.workers.lock().unwrap().push(handle);
}
}
pub fn pause(&self) {
if self.running.load(Ordering::Acquire) == 0 {
return;
}
let msg = ControlMsg::Pause;
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=requested"
);
self.clock.interrupt();
self.clock.pause();
for control_tx in &self.control_txs {
control_tx.send(msg);
}
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=completed outcome=ok"
);
}
pub fn resume(&self) {
if self.running.load(Ordering::Acquire) == 0 {
return;
}
let msg = ControlMsg::Resume;
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=requested"
);
self.clock.resume();
for control_tx in &self.control_txs {
control_tx.send(msg);
}
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=completed outcome=ok"
);
}
pub fn stop(&self) {
if self.running.load(Ordering::Acquire) == 0 {
return;
}
let msg = ControlMsg::Stop;
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=requested"
);
self.clock.interrupt();
for control_tx in &self.control_txs {
control_tx.send(msg);
}
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=completed outcome=ok"
);
}
pub fn seek(&self, target: Duration) {
if self.running.load(Ordering::Acquire) == 0 {
return;
}
let msg = ControlMsg::Seek(target);
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=requested"
);
self.clock.interrupt();
self.playback_clock.reset_for_seek();
for control_tx in &self.control_txs {
control_tx.send(msg);
}
pp_trace!(
pp_log: &self.pp_log,
"event=control control={msg:?} phase=completed outcome=ok"
);
}
}
struct RunningSourceGuard {
running: Arc<AtomicUsize>,
}
impl RunningSourceGuard {
fn new(running: Arc<AtomicUsize>) -> Self {
Self { running }
}
}
impl Drop for RunningSourceGuard {
fn drop(&mut self) {
self.running.fetch_sub(1, Ordering::AcqRel);
}
}
impl Drop for Pipeline {
fn drop(&mut self) {
self.stop();
let current_thread = thread::current().id();
let workers = self
.workers
.get_mut()
.unwrap_or_else(|poisoned| poisoned.into_inner());
for worker in workers.drain(..) {
if worker.thread().id() != current_thread {
let _ = worker.join();
}
}
}
}