use std::time::Duration;
use crate::container::Frame;
use super::Fragment;
use super::export::{apply_codec_durations, infer_missing_duration};
pub(super) struct Pending {
frame: Frame,
group_start: bool,
}
pub struct Fragmenter {
pub(super) track_id: u32,
pub(super) timescale: moq_net::Timescale,
pub(super) default_frame: Duration,
pub(super) is_video: bool,
pub(super) opus: bool,
pub(super) infer_missing: bool,
pub(super) pending: Option<Pending>,
pub(super) dts: Option<u64>,
pub(super) sequence: u32,
}
impl Fragmenter {
pub fn push(&mut self, frame: Frame) -> crate::Result<Vec<Fragment>> {
self.push_inner(frame, false)
}
pub fn push_group(&mut self, frame: Frame) -> crate::Result<Vec<Fragment>> {
self.push_inner(frame, true)
}
fn push_inner(&mut self, mut frame: Frame, group_start: bool) -> crate::Result<Vec<Fragment>> {
apply_codec_durations(std::slice::from_mut(&mut frame), self.opus);
if !stated(&frame) && !self.infer_missing {
return Err(super::Error::MissingVideoDuration.into());
}
let mut next = self.snapshot();
let mut fragments = Vec::new();
if let Some(mut pending) = next.pending.take() {
let successor = (!group_start).then_some(&frame);
infer_missing_duration(&mut pending.frame, successor, next.default_frame, next.timescale)?;
fragments.push(next.emit(pending.frame, pending.group_start)?);
}
if stated(&frame) {
fragments.push(next.emit(frame, group_start)?);
} else {
next.pending = Some(Pending { frame, group_start });
}
*self = next;
Ok(fragments)
}
pub fn flush(mut self) -> crate::Result<Option<Fragment>> {
let Some(mut pending) = self.pending.take() else {
return Ok(None);
};
infer_missing_duration(&mut pending.frame, None, self.default_frame, self.timescale)?;
Ok(Some(self.emit(pending.frame, pending.group_start)?))
}
fn emit(&mut self, frame: Frame, group_start: bool) -> crate::Result<Fragment> {
let pts = super::base_ticks(&frame, self.timescale)?;
let dts = match self.dts {
Some(dts) if !group_start => dts,
_ => pts,
};
let info = super::FragmentInfo {
track_id: self.track_id,
timescale: self.timescale,
sequence_number: self.sequence,
};
let ticks = frame
.duration
.map(|duration| super::trun_duration(duration, self.timescale))
.transpose()?
.unwrap_or(0);
let next_dts = dts.checked_add(u64::from(ticks)).ok_or(super::Error::PtsOverflow)?;
let data = super::encode_at(info, dts, std::slice::from_ref(&frame))?;
self.sequence = self.sequence.wrapping_add(1);
self.dts = Some(next_dts);
Ok(Fragment {
data,
init: false,
independent: !self.is_video || frame.keyframe,
duration: f64::from(ticks) / self.timescale.as_u64() as f64,
})
}
fn snapshot(&self) -> Self {
Self {
track_id: self.track_id,
timescale: self.timescale,
default_frame: self.default_frame,
is_video: self.is_video,
opus: self.opus,
infer_missing: self.infer_missing,
pending: self.pending.as_ref().map(|pending| Pending {
frame: pending.frame.clone(),
group_start: pending.group_start,
}),
dts: self.dts,
sequence: self.sequence,
}
}
}
fn stated(frame: &Frame) -> bool {
frame.duration.is_some_and(|duration| !duration.is_zero())
}
#[cfg(test)]
mod tests {
use bytes::Bytes;
use hang::catalog::{AudioConfig, VideoCodec, VideoConfig};
use moq_net::Timestamp;
use super::super::Muxer;
use super::*;
fn video_muxer() -> Muxer {
let mut config = VideoConfig::new(VideoCodec::VP8);
config.framerate = Some(30.0);
Muxer::video(&config).unwrap()
}
fn infer_video_durations() -> super::super::fragment::Config {
super::super::fragment::Config {
missing_duration: super::super::fragment::MissingDuration::InferFromPresentationTime,
}
}
fn tick_frame(pts: u64, keyframe: bool) -> Frame {
Frame {
timestamp: Timestamp::from_scale(pts, 30_000).unwrap(),
payload: Bytes::from_static(&[0xDE, 0xAD]),
keyframe,
duration: Some(Timestamp::from_scale(1_000, 30_000).unwrap()),
}
}
fn untimed_frame(pts: u64, keyframe: bool) -> Frame {
Frame {
duration: None,
..tick_frame(pts, keyframe)
}
}
fn sequence(fragment: &Fragment) -> u32 {
use mp4_atom::DecodeMaybe;
let mut cursor = std::io::Cursor::new(fragment.data.as_ref());
while let Some(atom) = mp4_atom::Any::decode_maybe(&mut cursor).unwrap() {
if let mp4_atom::Any::Moof(moof) = atom {
return moof.mfhd.sequence_number;
}
}
panic!("no moof");
}
fn one(mut fragments: Vec<Fragment>) -> Fragment {
assert_eq!(fragments.len(), 1);
fragments.pop().unwrap()
}
#[test]
fn stated_durations_emit_with_no_lookahead() {
let muxer = video_muxer();
let timescale = moq_net::Timescale::new(30_000).unwrap();
let mut fragmenter = muxer.fragmenter(Default::default());
let input = [tick_frame(0, true), tick_frame(3_000, false), tick_frame(1_000, false)];
let fragments: Vec<Fragment> = input
.iter()
.map(|frame| one(fragmenter.push(frame.clone()).unwrap()))
.collect();
let timelines: Vec<_> = fragments.iter().map(|f| super::super::timeline(&f.data)).collect();
assert_eq!(
timelines,
vec![(0, vec![0]), (1_000, vec![2_000]), (2_000, vec![-1_000])],
"tfdt advances one frame period while cts carries the reorder"
);
for (fragment, expected) in fragments.iter().zip(&input) {
let decoded = super::super::decode(fragment.data.clone(), timescale).unwrap();
assert_eq!(decoded.len(), 1);
assert_eq!(decoded[0].timestamp, expected.timestamp, "pts survives the reorder");
}
assert!(
muxer.fragmenter(Default::default()).flush().unwrap().is_none(),
"nothing pushed, nothing pending"
);
assert!(fragmenter.flush().unwrap().is_none(), "every frame was already emitted");
}
#[test]
fn a_pending_frame_is_timed_by_its_real_successor() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
let input = [
untimed_frame(0, true),
untimed_frame(1_500, false),
untimed_frame(3_000, false),
];
assert!(
fragmenter.push(input[0].clone()).unwrap().is_empty(),
"no duration and no successor yet"
);
let first = one(fragmenter.push(input[1].clone()).unwrap());
let second = one(fragmenter.push(input[2].clone()).unwrap());
let last = fragmenter.flush().unwrap().expect("the final pending frame");
let durations: Vec<_> = [&first, &second, &last]
.iter()
.map(|f| super::super::sample_durations(&f.data))
.collect();
assert_eq!(
durations,
vec![vec![Some(1_500)], vec![Some(1_500)], vec![Some(1_000)]],
"timed by the successor; only the flushed tail takes the catalog cadence"
);
let tfdts: Vec<u64> = [&first, &second, &last]
.iter()
.map(|f| super::super::timeline(&f.data).0)
.collect();
assert_eq!(tfdts, vec![0, 1_500, 3_000]);
}
#[test]
fn a_successor_with_another_scale_times_the_pending_frame() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
assert!(fragmenter.push(untimed_frame(0, true)).unwrap().is_empty());
let successor = Frame {
timestamp: Timestamp::from_micros(50_000).unwrap(),
duration: None,
..tick_frame(0, false)
};
let first = one(fragmenter.push(successor).unwrap());
assert_eq!(super::super::sample_durations(&first.data), vec![Some(1_500)]);
}
#[test]
fn a_coarse_predecessor_does_not_quantize_the_successor_gap() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
let first = Frame {
timestamp: Timestamp::from_secs(1).unwrap(),
duration: None,
..tick_frame(0, true)
};
let successor = Frame {
timestamp: Timestamp::from_scale(3, 2).unwrap(),
duration: None,
..tick_frame(0, false)
};
assert!(fragmenter.push(first).unwrap().is_empty());
let first = one(fragmenter.push(successor).unwrap());
assert_eq!(super::super::sample_durations(&first.data), vec![Some(15_000)]);
}
#[test]
fn a_simple_gap_does_not_require_a_representable_lcm() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
let start_scale = 4_000_000_007;
let end_scale = 4_000_000_009;
let first = Frame {
timestamp: Timestamp::from_scale(0, start_scale).unwrap(),
duration: None,
..tick_frame(0, true)
};
let successor = Frame {
timestamp: Timestamp::from_scale(end_scale, end_scale).unwrap(),
duration: None,
..tick_frame(0, false)
};
assert!(fragmenter.push(first).unwrap().is_empty());
let first = one(fragmenter.push(successor).unwrap());
assert_eq!(super::super::sample_durations(&first.data), vec![Some(30_000)]);
}
#[test]
fn a_coarse_tail_timestamp_keeps_the_catalog_fallback() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
let tail = Frame {
timestamp: Timestamp::from_secs(1).unwrap(),
duration: None,
..tick_frame(0, true)
};
assert!(fragmenter.push(tail).unwrap().is_empty());
let tail = fragmenter.flush().unwrap().expect("the pending tail");
assert_eq!(super::super::sample_durations(&tail.data), vec![Some(1_000)]);
}
#[test]
fn a_stated_frame_emits_with_the_pending_fragment() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
assert!(fragmenter.push(untimed_frame(0, true)).unwrap().is_empty());
let stated = Frame {
duration: Some(Timestamp::from_scale(1_000, 30_000).unwrap()),
..tick_frame(1_500, false)
};
let fragments = fragmenter.push(stated).unwrap();
assert_eq!(fragments.len(), 2);
let first = &fragments[0];
let second = &fragments[1];
assert_eq!(super::super::sample_durations(&first.data), vec![Some(1_500)]);
assert_eq!(super::super::sample_durations(&second.data), vec![Some(1_000)]);
assert!(fragmenter.flush().unwrap().is_none());
}
#[test]
fn durationless_reordered_video_is_rejected_by_default() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(Default::default());
let input = [
untimed_frame(0, true),
untimed_frame(3_000, false),
untimed_frame(1_000, false),
];
let err = fragmenter.push(input[0].clone()).unwrap_err();
assert!(matches!(
err,
crate::Error::Cmaf(super::super::Error::MissingVideoDuration)
));
}
#[test]
fn a_coarse_timescale_rejects_a_sub_tick_duration() {
let muxer = video_muxer().with_timescale(moq_net::Timescale::SECOND).unwrap();
let mut fragmenter = muxer.fragmenter(Default::default());
let err = fragmenter.push(tick_frame(0, true)).unwrap_err();
assert!(matches!(
err,
crate::Error::Cmaf(super::super::Error::SampleDurationTooSmall(1))
));
}
#[test]
fn fragmenter_rejects_an_inexact_sample_duration() {
let muxer = video_muxer().with_timescale(moq_net::Timescale::MILLI).unwrap();
let mut fragmenter = muxer.fragmenter(Default::default());
let input_scale = moq_net::Timescale::new(24).unwrap();
let frame = Frame {
timestamp: Timestamp::new(0, input_scale).unwrap(),
duration: Some(Timestamp::new(1, input_scale).unwrap()),
..tick_frame(0, true)
};
let err = fragmenter.push(frame).unwrap_err();
assert!(matches!(
err,
crate::Error::Cmaf(super::super::Error::SampleDurationInexact(1_000))
));
}
#[test]
fn fragments_carry_the_part_metadata() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
let mut fragments = Vec::new();
for (pts, keyframe) in [(0u64, true), (1_500, false), (3_000, false)] {
fragments.extend(fragmenter.push(untimed_frame(pts, keyframe)).unwrap());
}
fragments.extend(fragmenter.flush().unwrap());
assert!(!fragments.iter().any(|f| f.init), "these are media fragments");
assert_eq!(
fragments.iter().map(|f| f.independent).collect::<Vec<_>>(),
vec![true, false, false],
"video is independent only at a GOP boundary"
);
let durations: Vec<_> = fragments.iter().map(|f| (f.duration * 1e6).round() as u64).collect();
assert_eq!(durations, vec![50_000, 50_000, 33_333], "microseconds");
}
#[test]
fn audio_fragments_are_always_independent() {
let config = AudioConfig::new(hang::catalog::AudioCodec::Opus, 48_000, 2);
let muxer = Muxer::audio(&config).unwrap();
let mut fragmenter = muxer.fragmenter(Default::default());
let packet = Bytes::from_static(&[0x78, 0x00, 0x00, 0x00]);
for micros in [0u64, 20_000] {
let frame = Frame {
timestamp: Timestamp::from_micros(micros).unwrap(),
payload: packet.clone(),
keyframe: false,
duration: None,
};
let fragment = one(fragmenter.push(frame).unwrap());
assert!(fragment.independent, "audio fragments are always independent");
assert!((fragment.duration - 0.02).abs() < 1e-9, "the 20 ms TOC duration");
}
}
#[test]
fn a_new_group_does_not_time_the_pending_frame() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(infer_video_durations());
let paused_until = 2_405 * 30_000;
assert!(fragmenter.push(untimed_frame(0, true)).unwrap().is_empty());
let first = one(fragmenter.push(untimed_frame(1_500, false)).unwrap());
let second = one(fragmenter.push_group(untimed_frame(paused_until, true)).unwrap());
let third = fragmenter.flush().unwrap().expect("the keyframe itself");
assert_eq!(super::super::sample_durations(&first.data), vec![Some(1_500)]);
assert_eq!(
super::super::sample_durations(&second.data),
vec![Some(1_000)],
"the pause is a discontinuity, not a 2405 second sample"
);
assert_eq!(
super::super::timeline(&third.data).0,
paused_until,
"the new group re-anchors at its keyframe's presentation time"
);
}
#[test]
fn a_mid_group_keyframe_does_not_reanchor_the_timeline() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(Default::default());
let first = one(fragmenter.push(tick_frame(0, true)).unwrap());
let second = one(fragmenter.push(tick_frame(1_000, false)).unwrap());
let sync = one(fragmenter.push(tick_frame(5_000, true)).unwrap());
let next_group = one(fragmenter.push_group(tick_frame(10_000, true)).unwrap());
assert_eq!(super::super::timeline(&first.data).0, 0);
assert_eq!(super::super::timeline(&second.data).0, 1_000);
assert_eq!(super::super::timeline(&sync.data).0, 2_000);
assert_eq!(super::super::timeline(&next_group.data).0, 10_000);
assert!(
sync.independent,
"the mid-group sync sample stays independently decodable"
);
}
#[test]
fn fragments_number_consecutively() {
let muxer = video_muxer();
let mut fragmenter = muxer.fragmenter(Default::default());
let sequences: Vec<u32> = (0..3)
.map(|i| {
let fragment = one(fragmenter.push(tick_frame(i * 1_000, i == 0)).unwrap());
sequence(&fragment)
})
.collect();
assert_eq!(sequences, vec![0, 1, 2]);
}
#[test]
fn one_frame_matches_muxer_fragment() {
let muxer = video_muxer();
let frame = tick_frame(5_000, true);
let batch = muxer.fragment(0, std::slice::from_ref(&frame)).unwrap();
let pushed = one(muxer.fragmenter(Default::default()).push(frame).unwrap());
assert_eq!(pushed.data, batch);
}
}