use crate::error::PacketSinkError;
use ffmpeg_sys_next::{av_rescale_q_rnd, AVRational, AVRounding};
pub(crate) enum Timeline {
Unanchored { storage: Vec<i64> },
Anchored { offsets: Vec<i64> },
}
impl Timeline {
pub(crate) fn new(stream_count: usize) -> Self {
Timeline::Unanchored {
storage: Vec::with_capacity(stream_count),
}
}
pub(crate) fn ensure_anchored(
&mut self,
dts0: i64,
tb0: AVRational,
stream_time_bases: &[AVRational],
) -> Result<(), PacketSinkError> {
let storage = match self {
Timeline::Anchored { .. } => return Ok(()),
Timeline::Unanchored { storage } => storage,
};
debug_assert!(storage.is_empty());
for (stream_index, tb) in stream_time_bases.iter().enumerate() {
let offset =
unsafe { av_rescale_q_rnd(dts0, tb0, *tb, AVRounding::AV_ROUND_NEAR_INF) };
if offset == i64::MIN {
storage.clear();
return Err(PacketSinkError::TimestampOverflow { stream_index });
}
storage.push(offset);
}
let offsets = std::mem::take(storage);
*self = Timeline::Anchored { offsets };
Ok(())
}
pub(crate) fn offset(&self, stream_index: usize) -> i64 {
match self {
Timeline::Anchored { offsets } => offsets[stream_index],
Timeline::Unanchored { .. } => unreachable!("timeline queried before anchoring"),
}
}
}
pub(crate) struct StreamTimeline {
last_dts: Option<i64>,
pending_pts: Vec<i64>,
}
const PENDING_PTS_CAPACITY: usize = 16;
impl StreamTimeline {
pub(crate) fn new() -> Self {
Self {
last_dts: None,
pending_pts: Vec::with_capacity(PENDING_PTS_CAPACITY),
}
}
pub(crate) fn observe(
&mut self,
stream_index: usize,
pts: i64,
dts: i64,
) -> Result<(), PacketSinkError> {
if pts < dts {
return Err(PacketSinkError::PtsBeforeDts {
stream_index,
pts,
dts,
});
}
if let Some(prev) = self.last_dts {
if dts <= prev {
return Err(PacketSinkError::NonMonotonicDts {
stream_index,
prev,
current: dts,
});
}
}
let mut duplicate = false;
self.pending_pts.retain(|&p| {
duplicate |= p == pts;
p > dts
});
if duplicate {
return Err(PacketSinkError::DuplicatePts { stream_index, pts });
}
self.pending_pts.push(pts);
self.last_dts = Some(dts);
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn anchor_is_all_or_nothing_and_idempotent() {
let tb25 = AVRational { num: 1, den: 25 };
let tb_audio = AVRational { num: 1, den: 44100 };
let mut tl = Timeline::new(2);
tl.ensure_anchored(5, tb25, &[tb25, tb_audio]).unwrap();
assert_eq!(tl.offset(0), 5);
assert_eq!(tl.offset(1), 8820);
tl.ensure_anchored(100, tb25, &[tb25, tb_audio]).unwrap();
assert_eq!(tl.offset(0), 5);
}
#[test]
fn anchor_overflow_leaves_the_timeline_unanchored_and_reusable() {
let tb1 = AVRational { num: 1, den: 1 };
let tb90k = AVRational { num: 1, den: 90000 };
let mut tl = Timeline::new(2);
assert!(matches!(
tl.ensure_anchored(i64::MAX / 2, tb1, &[tb1, tb90k]),
Err(PacketSinkError::TimestampOverflow { stream_index: 1 })
));
assert!(matches!(tl, Timeline::Unanchored { .. }));
tl.ensure_anchored(5, tb1, &[tb1, tb90k]).unwrap();
assert_eq!(tl.offset(0), 5);
assert_eq!(tl.offset(1), 450_000);
}
#[test]
fn observe_rejects_the_boundary_duplicate() {
let mut st = StreamTimeline::new();
st.observe(0, 3, 0).unwrap();
assert!(matches!(
st.observe(0, 3, 3),
Err(PacketSinkError::DuplicatePts { pts: 3, .. })
));
}
}