use std::{
collections::VecDeque,
sync::Arc,
thread,
time::{Duration, Instant},
};
use crate::pp_log::{PpLog, pp_info};
use ffmpeg_next as ffmpeg;
use thiserror::Error as ThisError;
use crate::{
buffer::MediaBuffer,
clock::Clock,
control::ControlMsg,
element::{Element, ElementType, Sink, Source, element_pp_log},
pad::SrcPad,
time::{InvalidTimeBase, MediaTimestamp, TimeBase},
};
#[derive(Debug, ThisError)]
pub enum PacerError {
#[error(
"invalid time base {numerator}/{denominator}: both numerator and denominator must be positive"
)]
InvalidTimeBase { numerator: i32, denominator: i32 },
#[error("pts {pts} is too far from this pacer's first pts {first_pts} to pace against")]
TimestampDeltaOverflow { pts: i64, first_pts: i64 },
}
fn nanoseconds() -> TimeBase {
TimeBase::new_unchecked(ffmpeg::Rational::new(1, 1_000_000_000))
}
const INTERRUPT_POLL_INTERVAL: Duration = Duration::from_millis(10);
pub struct Pacer {
pp_log: PpLog,
name: Arc<str>,
time_base: TimeBase,
clock: Arc<Clock>,
first_pts: Option<i64>,
interrupt_epoch: u64,
pending: VecDeque<MediaBuffer>,
pad: SrcPad,
}
impl Pacer {
pub fn new(
name: impl Into<String>,
time_base: ffmpeg::Rational,
clock: Arc<Clock>,
) -> Result<Self, PacerError> {
let name: Arc<str> = name.into().into();
let pp_log = element_pp_log(ElementType::Pacer, &name, None);
pp_info!(pp_log: &pp_log, "created: time_base={time_base}");
let pad = SrcPad::new(format!("{name}_src"));
let interrupt_epoch = clock.interrupt_epoch();
let time_base = TimeBase::try_new(time_base).map_err(
|InvalidTimeBase {
numerator,
denominator,
}| PacerError::InvalidTimeBase {
numerator,
denominator,
},
)?;
Ok(Self {
name,
pp_log,
time_base,
clock,
first_pts: None,
interrupt_epoch,
pending: VecDeque::new(),
pad,
})
}
fn wait_for(&mut self, pts: Option<i64>) -> Result<bool, PacerError> {
if self.clock.interrupt_epoch() != self.interrupt_epoch {
return Ok(false);
}
let Some(pts) = pts else { return Ok(true) };
let first_pts = *self.first_pts.get_or_insert(pts);
let elapsed_ticks = pts
.checked_sub(first_pts)
.ok_or(PacerError::TimestampDeltaOverflow { pts, first_pts })?;
if elapsed_ticks <= 0 {
return Ok(true);
}
let elapsed_ns = MediaTimestamp::new_unchecked(elapsed_ticks, self.time_base)
.rescale(nanoseconds())
.max(0) as u64;
let due = self.clock.start() + Duration::from_nanos(elapsed_ns);
loop {
if self.clock.interrupt_epoch() != self.interrupt_epoch {
return Ok(false);
}
let now = Instant::now();
if due <= now {
return Ok(true);
}
thread::sleep((due - now).min(INTERRUPT_POLL_INTERVAL));
}
}
}
impl Element for Pacer {
fn name(&self) -> Arc<str> {
self.name.clone()
}
fn element_type(&self) -> ElementType {
ElementType::Pacer
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Source for Pacer {
fn src_pads(&mut self) -> &mut [SrcPad] {
std::slice::from_mut(&mut self.pad)
}
}
impl Sink for Pacer {
fn consume(&mut self, buf: MediaBuffer) -> crate::error::Result<()> {
self.pending.push_back(buf);
while let Some(buf) = self.pending.pop_front() {
let ready = match &buf {
MediaBuffer::Packet(packet) => self.wait_for(packet.pts())?,
MediaBuffer::Video(frame) => self.wait_for(frame.pts())?,
MediaBuffer::Audio(frame) => self.wait_for(frame.pts())?,
MediaBuffer::Eos => true,
};
if !ready {
self.pending.push_front(buf);
return Ok(());
}
self.pad.push(buf)?;
}
Ok(())
}
fn control(&mut self, msg: ControlMsg) -> crate::error::Result<()> {
self.interrupt_epoch = self.clock.interrupt_epoch();
match msg {
ControlMsg::Seek(_) => {
self.pending.clear();
self.first_pts = None;
self.clock.reset();
}
ControlMsg::Stop => self.pending.clear(),
ControlMsg::Pause | ControlMsg::Resume => {}
}
self.pad.control(msg)
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::{sync::mpsc, time::Duration};
fn packet(pts: i64) -> MediaBuffer {
let mut packet = ffmpeg::Packet::empty();
packet.set_pts(Some(pts));
MediaBuffer::Packet(Arc::new(packet))
}
#[test]
fn long_wait_returns_promptly_when_control_interrupts_it() {
let clock = Arc::new(Clock::new());
let mut pacer = Pacer::new("pacer", ffmpeg::Rational::new(1, 1), clock.clone()).unwrap();
assert!(
pacer.wait_for(Some(0)).unwrap(),
"first pts should establish the anchor"
);
let (started_tx, started_rx) = mpsc::channel();
let worker = thread::spawn(move || {
started_tx.send(()).expect("test receiver alive");
pacer.wait_for(Some(60))
});
started_rx.recv().expect("paced wait should start");
thread::sleep(Duration::from_millis(20));
clock.interrupt();
assert!(
!worker
.join()
.expect("paced wait should return")
.expect("interrupted wait is Ok(false), not an error"),
"an interrupted paced wait must return before its due time"
);
}
#[test]
fn pause_retains_interrupted_buffer_but_seek_and_stop_discard_it() {
let clock = Arc::new(Clock::new());
let mut pacer = Pacer::new("pacer", ffmpeg::Rational::new(1, 1), clock.clone()).unwrap();
clock.interrupt();
pacer.consume(packet(0)).expect("interrupted consume");
assert_eq!(pacer.pending.len(), 1);
pacer.control(ControlMsg::Pause).expect("pause");
assert_eq!(pacer.pending.len(), 1, "pause must retain the buffer");
pacer
.control(ControlMsg::Seek(Duration::ZERO))
.expect("seek");
assert!(pacer.pending.is_empty(), "seek must discard stale data");
clock.interrupt();
pacer.consume(packet(1)).expect("interrupted consume");
assert_eq!(pacer.pending.len(), 1);
pacer.control(ControlMsg::Stop).expect("stop");
assert!(pacer.pending.is_empty(), "stop must abandon pending data");
}
#[test]
fn new_rejects_an_invalid_time_base() {
let clock = Arc::new(Clock::new());
for rational in [
ffmpeg::Rational::new(0, 1),
ffmpeg::Rational::new(1, 0),
ffmpeg::Rational::new(-1, 1),
ffmpeg::Rational::new(1, -1),
] {
assert!(
matches!(
Pacer::new("pacer", rational, clock.clone()),
Err(PacerError::InvalidTimeBase { .. })
),
"expected {rational} to be rejected"
);
}
}
#[test]
fn a_pathological_pts_jump_is_a_typed_error_not_silent_passthrough() {
let clock = Arc::new(Clock::new());
let mut pacer = Pacer::new("pacer", ffmpeg::Rational::new(1, 1), clock).unwrap();
assert!(pacer.consume(packet(-1)).is_ok(), "establishes first_pts");
let error = pacer
.consume(packet(i64::MAX))
.expect_err("pts far enough from first_pts to overflow the subtraction");
assert!(matches!(
error,
crate::Error::PacerError(PacerError::TimestampDeltaOverflow {
pts: i64::MAX,
first_pts: -1,
})
));
assert!(
pacer.pending.is_empty(),
"the overflowing buffer must not get stuck in `pending`"
);
assert!(pacer.wait_for(Some(0)).is_ok());
}
}