use std::time::{Duration, Instant};
use crate::pp_log::pp_trace;
use crossbeam_channel::{Receiver, Sender, unbounded};
use crate::{
bus::{Bus, BusEvent},
element::SourceElement,
error::Result,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ControlMsg {
Pause,
Resume,
Stop,
Seek(Duration),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RequestKind {
Control(ControlMsg),
Finish,
}
pub(crate) struct Request {
pub(crate) kind: RequestKind,
pub(crate) ack: Sender<()>,
}
#[derive(Clone)]
pub struct ControlSender {
tx: Sender<Request>,
}
#[derive(Clone)]
pub struct ControlReceiver {
pub(crate) rx: Receiver<Request>,
}
pub fn channel() -> (ControlSender, ControlReceiver) {
let (tx, rx) = unbounded();
(ControlSender { tx }, ControlReceiver { rx })
}
impl ControlSender {
pub fn send(&self, msg: ControlMsg) {
self.send_request(RequestKind::Control(msg));
}
pub(crate) fn finish(&self) {
self.send_request(RequestKind::Finish);
}
fn send_request(&self, kind: RequestKind) {
let (ack_tx, ack_rx) = crossbeam_channel::bounded(0);
if self.tx.send(Request { kind, ack: ack_tx }).is_ok() {
let _ = ack_rx.recv();
}
}
}
impl ControlReceiver {
pub(crate) fn try_recv(&self) -> Option<(RequestKind, Sender<()>)> {
self.rx.try_recv().ok().map(|r| (r.kind, r.ack))
}
pub(crate) fn recv(&self) -> Option<(RequestKind, Sender<()>)> {
self.rx.recv().ok().map(|r| (r.kind, r.ack))
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct ControlOutcome {
pub stopped: bool,
pub paused_for: Duration,
}
pub fn drain_control<S: SourceElement>(
control: &ControlReceiver,
source: &mut S,
bus: &Bus,
) -> Result<ControlOutcome> {
let mut paused_for = Duration::ZERO;
while let Some((request, ack)) = control.try_recv() {
let RequestKind::Control(msg) = request else {
apply_finish(source, bus, &ack);
return Ok(ControlOutcome {
stopped: true,
paused_for,
});
};
if msg == ControlMsg::Pause {
let pause_start = Instant::now();
apply_one(source, bus, msg, &ack)?;
let stopped = wait_out_pause(control, source, bus)?;
paused_for += pause_start.elapsed();
if stopped {
return Ok(ControlOutcome {
stopped: true,
paused_for,
});
}
continue;
}
if apply_one(source, bus, msg, &ack)? {
return Ok(ControlOutcome {
stopped: true,
paused_for,
});
}
}
Ok(ControlOutcome {
stopped: false,
paused_for,
})
}
pub(crate) fn apply_finish<S: SourceElement>(source: &mut S, bus: &Bus, ack: &Sender<()>) {
pp_trace!(
pp_log: source.pp_log(),
"event=finish phase=received"
);
let pp_log = source.pp_log().clone();
let element_type = source.element_type();
let name = source.name();
for pad in source.src_pads() {
if let Err(error) = pad.push_eos(&pp_log) {
bus.post(
&pp_log,
BusEvent::Error {
element_type,
name: name.clone(),
error,
},
);
}
}
let _ = ack.send(());
pp_trace!(
pp_log: source.pp_log(),
"event=finish phase=completed outcome=ok"
);
}
pub(crate) fn apply_one<S: SourceElement>(
source: &mut S,
bus: &Bus,
msg: ControlMsg,
ack: &Sender<()>,
) -> Result<bool> {
let is_stop = apply_one_unacked(source, bus, msg)?;
let _ = ack.send(());
Ok(is_stop)
}
pub(crate) fn apply_one_unacked<S: SourceElement>(
source: &mut S,
bus: &Bus,
msg: ControlMsg,
) -> Result<bool> {
pp_trace!(
pp_log: source.pp_log(),
"event=control control={msg:?} phase=received"
);
let result: Result<bool> = (|| {
apply_seek(source, bus, msg)?;
for pad in source.src_pads() {
pad.control(msg)?;
}
Ok(msg == ControlMsg::Stop)
})();
match &result {
Ok(_) => pp_trace!(
pp_log: source.pp_log(),
"event=control control={msg:?} phase=completed outcome=ok"
),
Err(error) => pp_trace!(
pp_log: source.pp_log(),
"event=control control={msg:?} phase=completed outcome=error error={error}"
),
}
result
}
pub(crate) fn wait_out_pause<S: SourceElement>(
control: &ControlReceiver,
source: &mut S,
bus: &Bus,
) -> Result<bool> {
loop {
let Some((request, ack)) = control.recv() else {
return Ok(true); };
let RequestKind::Control(msg) = request else {
apply_finish(source, bus, &ack);
return Ok(true);
};
if apply_one(source, bus, msg, &ack)? {
return Ok(true);
}
if msg == ControlMsg::Resume {
return Ok(false);
}
}
}
fn apply_seek<S: SourceElement>(source: &mut S, bus: &Bus, msg: ControlMsg) -> Result<()> {
if let ControlMsg::Seek(target) = msg {
let landed = source.seek(target)?;
bus.post(
source.pp_log(),
BusEvent::Seeked {
element_type: source.element_type(),
name: source.name(),
requested: target,
landed,
},
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use std::{sync::Arc, thread};
use crate::pp_log::PpLog;
use super::*;
use crate::{
buffer::MediaBuffer,
element::{Element, ElementType, Sink, Source, element_pp_log},
pad::SrcPad,
};
struct DummySource {
pp_log: PpLog,
pad: SrcPad,
}
impl DummySource {
fn new() -> Self {
Self {
pp_log: element_pp_log(ElementType::Other, "dummy", None),
pad: SrcPad::new("dummy_src"),
}
}
}
impl Element for DummySource {
fn name(&self) -> Arc<str> {
"dummy".into()
}
fn element_type(&self) -> ElementType {
ElementType::Other
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Source for DummySource {
fn src_pads(&mut self) -> &mut [SrcPad] {
std::slice::from_mut(&mut self.pad)
}
}
impl SourceElement for DummySource {
fn run(&mut self, _control: &ControlReceiver, _bus: &Bus) -> Result<()> {
unreachable!("not exercised by these tests")
}
fn seek(&mut self, target: Duration) -> Result<Duration> {
Ok(target)
}
}
struct SlowPauseSink {
pp_log: PpLog,
pause_delay: Duration,
}
impl Element for SlowPauseSink {
fn name(&self) -> Arc<str> {
"slow-pause".into()
}
fn element_type(&self) -> ElementType {
ElementType::Other
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Sink for SlowPauseSink {
fn consume(&mut self, _buf: MediaBuffer) -> Result<()> {
Ok(())
}
fn control(&mut self, msg: ControlMsg) -> Result<()> {
if msg == ControlMsg::Pause {
thread::sleep(self.pause_delay);
}
Ok(())
}
}
#[test]
fn wait_out_pause_treats_a_dropped_sender_as_stop() {
let (tx, rx) = channel();
drop(tx);
let (bus, _bus_rx) = Bus::new();
let mut source = DummySource::new();
let stopped = wait_out_pause(&rx, &mut source, &bus)
.expect("no real seek/push happens on this path, so this can't fail");
assert!(
stopped,
"a dropped ControlSender must be treated the same as an explicit Stop"
);
}
#[test]
fn wait_out_pause_blocks_until_resume_then_returns_false() {
let (tx, rx) = channel();
let (bus, _bus_rx) = Bus::new();
let mut source = DummySource::new();
let worker = thread::spawn(move || wait_out_pause(&rx, &mut source, &bus));
tx.send(ControlMsg::Pause);
tx.send(ControlMsg::Resume);
let stopped = worker
.join()
.expect("worker must not panic")
.expect("no real seek/push happens on this path, so this can't fail");
assert!(
!stopped,
"Resume must unblock wait_out_pause with Ok(false)"
);
}
#[test]
fn drain_control_counts_the_pause_cascade_as_paused_time() {
let pause_delay = Duration::from_millis(80);
let (tx, rx) = channel();
let controller = thread::spawn(move || {
tx.send(ControlMsg::Pause);
tx.send(ControlMsg::Resume);
});
let (bus, _bus_rx) = Bus::new();
let mut source = DummySource::new();
source.pad.link(Box::new(SlowPauseSink {
pause_delay,
pp_log: element_pp_log(ElementType::Other, "slow-pause", None),
}));
let outcome = loop {
let outcome = drain_control(&rx, &mut source, &bus)
.expect("the synthetic control cascade cannot fail");
if outcome.paused_for > Duration::ZERO {
break outcome;
}
thread::yield_now();
};
controller.join().expect("controller must not panic");
assert!(!outcome.stopped);
assert!(
outcome.paused_for >= Duration::from_millis(60),
"the {:?} Pause cascade was omitted from paused_for: {:?}",
pause_delay,
outcome.paused_for
);
}
}