use std::collections::VecDeque;
use bytes::Bytes;
use super::decoder::{Config, Decoder};
use crate::resample::{Resampler, remix, validate_remix};
use crate::{Activity, Error, Format, Frame, Layout};
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[non_exhaustive]
pub enum Start {
#[default]
Oldest,
Latest,
}
#[derive(Clone, Debug, Default)]
#[non_exhaustive]
pub struct Output {
pub format: Format,
pub sample_rate: Option<u32>,
pub layout: Option<Layout>,
}
#[derive(Clone, Debug, Default)]
#[non_exhaustive]
pub struct Options {
pub decoder: Config,
pub output: Output,
pub max_age: std::time::Duration,
pub start: Start,
}
impl Options {
pub fn new() -> Self {
Self::default()
}
}
pub struct Consumer {
decoder: Decoder,
track: moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
resampler: Option<Resampler>,
options: Options,
max_age: std::time::Duration,
resolved_sample_rate: u32,
resolved_layout: Layout,
next_start: Option<moq_net::Timestamp>,
ready: VecDeque<Frame>,
spans: VecDeque<ActivitySpan>,
trailing: Activity,
epoch: Option<moq_net::Timestamp>,
delay_trimmed: usize,
frames_decoded: usize,
end: Option<moq_net::Timestamp>,
terminal_start: Option<moq_net::Timestamp>,
discontinuity: u64,
}
struct ActivitySpan {
end: moq_net::Timestamp,
activity: Activity,
}
impl Consumer {
pub async fn new(
broadcast: &moq_net::broadcast::Consumer,
catalog: &hang::catalog::AudioConfig,
name: impl Into<String>,
options: Options,
) -> Result<Self, Error> {
let decoder = Decoder::new(catalog, &options.decoder)?;
let sample_rate = options.output.sample_rate.unwrap_or_else(|| decoder.sample_rate());
let layout = options.output.layout.unwrap_or_else(|| decoder.layout());
validate_remix(decoder.layout(), layout)?;
let resampler = if sample_rate == decoder.sample_rate() {
None
} else {
let chunk_frames = (decoder.sample_rate() as usize * 20) / 1000;
Some(Resampler::new(
decoder.sample_rate(),
sample_rate,
decoder.layout().channels(),
chunk_frames,
)?)
};
let name = name.into();
let track = broadcast.track(&name)?;
let mut subscriber = track
.subscribe(
moq_net::track::Subscription::default()
.with_priority(hang::catalog::PRIORITY.audio)
.with_max_age(options.max_age),
)
.await?;
if options.start == Start::Latest
&& let Some(live_edge) = track.latest()
{
subscriber.set_groups(live_edge..);
}
let track = subscriber;
let max_age = options.max_age.min(track.info().max_age);
let container = moq_mux::catalog::hang::Container::try_from(catalog)?;
let track = moq_mux::container::Consumer::new(track, container);
Ok(Self {
decoder,
track,
resampler,
options,
max_age,
resolved_sample_rate: sample_rate,
resolved_layout: layout,
next_start: None,
ready: VecDeque::new(),
spans: VecDeque::new(),
trailing: Activity::Active,
epoch: None,
delay_trimmed: 0,
frames_decoded: 0,
end: None,
terminal_start: None,
discontinuity: 0,
})
}
pub fn options(&self) -> &Options {
&self.options
}
pub fn max_age(&self) -> std::time::Duration {
self.max_age
}
pub fn sample_rate(&self) -> u32 {
self.resolved_sample_rate
}
pub fn layout(&self) -> Layout {
self.resolved_layout
}
pub async fn read(&mut self) -> Result<Option<Frame>, Error> {
loop {
if let Some(frame) = self.ready.pop_front() {
return Ok(Some(frame));
}
let mux_frame = self.track.read().await?;
self.apply_discontinuity()?;
let Some(mux_frame) = mux_frame else {
return self.flush();
};
if let Some(end) = self.track.end()
&& self.end != Some(end)
{
self.end = Some(end);
self.frames_decoded = 0;
self.terminal_start = None;
}
if self.end.is_none()
&& self
.next_start
.is_some_and(|next| discontinuous(next, mux_frame.timestamp))
&& let Some(frame) = self.gap()?
{
self.ready.push_back(frame);
}
let rate = self.decoder.sample_rate();
let epoch = *self.epoch.get_or_insert(mux_frame.timestamp);
let delay = self.decoder.delay_remaining();
let decoded = self.decoder.decode(&mux_frame.payload)?;
let trimmed = delay - self.decoder.delay_remaining();
self.delay_trimmed += trimmed;
let activity = decoded.activity;
let mut decoded = decoded.samples;
if let Some(end) = self.end {
let terminal_start = *self
.terminal_start
.get_or_insert(rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch));
let total = frames_between(terminal_start, end, rate)?;
let remaining = total.saturating_sub(self.frames_decoded);
decoded.truncate(remaining.saturating_mul(self.decoder.layout().channels() as usize));
}
let frames = decoded.len() / self.decoder.layout().channels() as usize;
let decoded_at = if let Some(terminal_start) = self.terminal_start {
advance(terminal_start, self.frames_decoded, rate)?
} else {
rewind(mux_frame.timestamp, self.delay_trimmed, rate)?.max(epoch)
};
if self.end.is_some() {
self.frames_decoded += frames;
}
self.next_start = Some(advance(mux_frame.timestamp, frames + trimmed, rate)?);
if decoded.is_empty() {
continue;
}
let (pcm, timestamp) = match self.resampler.as_mut() {
Some(r) => {
let held = if r.pending_frames() == 0 {
decoded_at
} else {
r.held_at().unwrap_or(decoded_at)
};
let skipped = r.skipped();
let pcm = r.process(&decoded, decoded_at)?;
(pcm, rewind(held, skipped, self.resolved_sample_rate)?)
}
None => (decoded, decoded_at),
};
let decoded_end = advance(decoded_at, frames, rate)?;
let resampled = self.resampler.is_some();
if resampled {
self.spans.push_back(ActivitySpan {
end: decoded_end,
activity,
});
}
if pcm.is_empty() {
continue;
}
let activity = if resampled {
self.activity_at(timestamp)
} else {
activity
};
let frame = self.frame(pcm, timestamp, activity)?;
self.ready.push_back(frame);
}
}
fn apply_discontinuity(&mut self) -> Result<(), Error> {
let discontinuity = self.track.discontinuity();
if discontinuity == self.discontinuity {
return Ok(());
}
self.discontinuity = discontinuity;
self.next_start = None;
self.spans.clear();
self.trailing = Activity::Active;
self.frames_decoded = 0;
self.end = None;
self.terminal_start = None;
self.epoch = None;
self.delay_trimmed = 0;
self.decoder.reapply_delay();
Ok(())
}
fn gap(&mut self) -> Result<Option<Frame>, Error> {
self.decoder.reset_prediction()?;
let mut frame = None;
if let Some(resampler) = self.resampler.as_mut() {
let held = resampler.held_at();
let skipped = resampler.skipped();
let pcm = resampler.drain()?;
frame = self.tail(pcm, held, skipped)?;
}
self.next_start = None;
self.spans.clear();
self.trailing = Activity::Active;
self.epoch = None;
self.delay_trimmed = 0;
Ok(frame)
}
fn flush(&mut self) -> Result<Option<Frame>, Error> {
let Some(resampler) = self.resampler.take() else {
return Ok(None);
};
let held = resampler.held_at();
let skipped = resampler.skipped();
self.tail(resampler.flush()?, held, skipped)
}
fn tail(
&mut self,
pcm: Vec<f32>,
held: Option<moq_net::Timestamp>,
skipped: usize,
) -> Result<Option<Frame>, Error> {
let Some(held) = held.filter(|_| !pcm.is_empty()) else {
return Ok(None);
};
let timestamp = rewind(held, skipped, self.resolved_sample_rate)?;
let activity = self.activity_at(timestamp);
Ok(Some(self.frame(pcm, timestamp, activity)?))
}
fn activity_at(&mut self, timestamp: moq_net::Timestamp) -> Activity {
while let Some(span) = self.spans.front().filter(|span| span.end <= timestamp) {
self.trailing = span.activity;
self.spans.pop_front();
}
self.spans.front().map_or(self.trailing, |span| span.activity)
}
fn frame(&self, pcm: Vec<f32>, timestamp: moq_net::Timestamp, activity: Activity) -> Result<Frame, Error> {
let pcm = if self.decoder.layout() == self.resolved_layout {
pcm
} else {
remix(&pcm, self.decoder.layout(), self.resolved_layout)?
};
let bytes = self
.options
.output
.format
.from_interleaved_f32(&pcm, self.resolved_layout.channels())?;
Ok(Frame {
timestamp,
data: Bytes::from(bytes),
activity,
})
}
}
fn discontinuous(expected: moq_net::Timestamp, timestamp: moq_net::Timestamp) -> bool {
let scale = expected.scale().max(timestamp.scale());
let quantum = scale.min(moq_net::Timescale::default());
let tolerance = (scale.as_u64() as u128).div_ceil(quantum.as_u64() as u128) + 1;
expected.as_scale(scale).abs_diff(timestamp.as_scale(scale)) > tolerance
}
fn advance(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
if frames == 0 {
return Ok(timestamp);
}
let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
Ok(timestamp.checked_add(offset)?)
}
fn frames_between(start: moq_net::Timestamp, end: moq_net::Timestamp, sample_rate: u32) -> Result<usize, Error> {
let duration = end.checked_sub(start)?;
let frames = (std::time::Duration::from(duration).as_nanos() * sample_rate as u128 + 500_000_000) / 1_000_000_000;
usize::try_from(frames).map_err(|_| Error::Unsupported("audio duration does not fit in memory".into()))
}
fn rewind(timestamp: moq_net::Timestamp, frames: usize, sample_rate: u32) -> Result<moq_net::Timestamp, Error> {
if frames == 0 {
return Ok(timestamp);
}
let offset = moq_net::Timestamp::from_scale(frames as u64, sample_rate as u64)?.convert(timestamp.scale())?;
Ok(timestamp
.checked_sub(offset)
.unwrap_or(moq_net::Timestamp::new(0, timestamp.scale())?))
}
#[cfg(test)]
mod tests {
use moq_net::Timestamp;
use super::*;
use crate::encode::{Encoder, Input, Options as EncodeOptions, Producer, Settings};
use crate::{Format, Layout};
#[tokio::test]
async fn remixes_mono_stream_to_stereo_output() {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, moq_mux::catalog::Config::default()).unwrap();
let subscriber = broadcast.consume();
let input = Input {
format: Format::F32,
sample_rate: 48_000,
layout: Layout::Mono,
};
let options = EncodeOptions {
track: Some("audio".to_string()),
settings: Settings::new(48_000, Layout::Mono),
..EncodeOptions::default()
};
let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
let catalog = Encoder::new(&Settings::new(input.sample_rate, input.layout))
.unwrap()
.catalog();
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
output: Output {
layout: Some(Layout::Stereo),
..Output::default()
},
..Options::new()
},
)
.await
.unwrap();
let samples = vec![0.1f32; 960];
let mut data = Vec::with_capacity(samples.len() * size_of::<f32>());
for sample in samples {
data.extend_from_slice(&sample.to_le_bytes());
}
producer.write(&Frame::new(data.into(), Timestamp::ZERO)).unwrap();
let frame = consumer.read().await.unwrap().expect("decoded frame");
let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
assert_eq!(samples.len(), (960 - 312) * 2);
for pair in samples.as_chunks::<2>().0.iter() {
assert_eq!(pair[0], pair[1]);
}
}
#[tokio::test]
async fn opus_timestamps_follow_the_48k_clock() {
use crate::decode::decoder::tests::{opus_catalog, opus_packets};
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let catalog = opus_catalog(moq_mux::codec::opus::Config::new(44_100, 1).with_pre_skip(312));
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Options::new())
.await
.unwrap();
assert_eq!(consumer.sample_rate(), 48_000);
for (packet, payload) in opus_packets(3).into_iter().enumerate() {
producer
.write(moq_mux::container::Frame {
timestamp: Timestamp::from_micros(packet as u64 * 20_000).unwrap(),
duration: None,
payload,
keyframe: packet == 0,
})
.unwrap();
}
for (micros, frames) in [(0, 960 - 312), (13_500, 960), (33_500, 960)] {
let frame = consumer.read().await.unwrap().expect("decoded frame");
assert_eq!(frame.timestamp.as_micros(), micros);
assert_eq!(frame.data.len() / size_of::<f32>(), frames);
}
}
#[tokio::test]
async fn resampled_timestamps_follow_the_samples() {
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
output: Output {
sample_rate: Some(48_000),
..Output::default()
},
max_age: std::time::Duration::from_secs(1),
..Options::new()
},
)
.await
.unwrap();
const FRAMES: u64 = 1024;
let payload: Bytes = vec![0u8; FRAMES as usize * size_of::<f32>()].into();
for packet in 0..2 {
producer
.write(moq_mux::container::Frame {
timestamp: moq_net::Timestamp::from_scale(packet * FRAMES, 44_100).unwrap(),
duration: None,
payload: payload.clone(),
keyframe: true,
})
.unwrap();
}
let first = consumer.read().await.unwrap().expect("decoded frame");
assert_eq!(first.timestamp.as_micros(), 0);
let second = consumer.read().await.unwrap().expect("decoded frame");
let first_frames = (first.data.len() / size_of::<f32>()) as u128;
let ends_at = first_frames * 1_000_000 / 48_000;
let gap = second.timestamp.as_micros().abs_diff(ends_at);
assert!(gap < 100, "expected the frames to meet, got a {gap} us gap");
}
#[tokio::test]
async fn resampled_tail_survives_the_end_of_the_track() {
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 44_100, 1);
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
output: Output {
sample_rate: Some(48_000),
..Output::default()
},
..Options::new()
},
)
.await
.unwrap();
const FRAMES: usize = 1024;
let payload: Bytes = vec![0u8; FRAMES * size_of::<f32>()].into();
producer
.write(moq_mux::container::Frame {
timestamp: moq_net::Timestamp::ZERO,
duration: None,
payload,
keyframe: true,
})
.unwrap();
producer.finish().unwrap();
let first = consumer.read().await.unwrap().expect("decoded frame");
let first_frames = first.data.len() / size_of::<f32>();
let tail = consumer.read().await.unwrap().expect("flushed tail");
let tail_frames = tail.data.len() / size_of::<f32>();
assert!((215..=230).contains(&tail_frames), "unexpected tail: {tail_frames}");
let ends_at = (first_frames as u128) * 1_000_000 / 48_000;
let gap = tail.timestamp.as_micros().abs_diff(ends_at);
assert!(gap < 100, "expected the tail to meet the body, got a {gap} us gap");
let total = first_frames + tail_frames;
assert!((1105..=1120).contains(&total), "unexpected total: {total}");
assert!(consumer.read().await.unwrap().is_none());
}
#[tokio::test]
async fn resampling_keeps_the_activity_boundary_on_its_source() {
let mut encoder = Encoder::new(&Settings {
dtx: true,
bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
frame_duration: std::time::Duration::from_millis(10),
..Settings::new(48_000, Layout::Mono)
})
.unwrap();
let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Opus, 48_000, 1);
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
output: Output {
sample_rate: Some(44_100),
..Output::default()
},
max_age: std::time::Duration::from_secs(1),
..Options::new()
},
)
.await
.unwrap();
let active = vec![0.5; encoder.frame_size()];
let silence = vec![0.0; encoder.frame_size()];
let mut first_dtx = None;
for index in 0..40u64 {
let packet = encoder.encode(if index == 0 { &active } else { &silence }).unwrap();
let timestamp = Timestamp::from_scale(index * encoder.frame_size() as u64, 48_000).unwrap();
if first_dtx.is_none() && packet.activity.is_dtx() {
first_dtx = Some(timestamp);
}
producer
.write(moq_mux::container::Frame {
timestamp,
payload: packet.payload,
keyframe: true,
duration: None,
})
.unwrap();
producer.cut(None).unwrap();
}
producer.finish().unwrap();
let expected = first_dtx.expect("silence should enter Opus DTX");
let mut actual = None;
while let Some(frame) = consumer.read().await.unwrap() {
assert!(!frame.data.is_empty(), "read returned a frame with no samples");
if frame.activity.is_dtx() {
actual = Some(frame.timestamp);
break;
}
}
let actual = actual.expect("consumer should report Opus DTX");
let delay = actual.as_micros() as i128 - expected.as_micros() as i128;
let chunk_us = 20_000i128;
assert!(
(0..chunk_us).contains(&delay),
"DTX label landed {delay} us from its source, outside [0, {chunk_us})"
);
}
async fn pcm_gaps(rate: u32, out_rate: u32, frames: usize, stamps: &[Timestamp]) -> Vec<(u128, usize)> {
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, rate, 1);
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
output: Output {
sample_rate: Some(out_rate),
..Output::default()
},
max_age: std::time::Duration::from_secs(1),
..Options::new()
},
)
.await
.unwrap();
let payload: Bytes = vec![0u8; frames * size_of::<f32>()].into();
for stamp in stamps {
producer
.write(moq_mux::container::Frame {
timestamp: *stamp,
duration: None,
payload: payload.clone(),
keyframe: true,
})
.unwrap();
}
producer.finish().unwrap();
let mut read = Vec::new();
while let Some(frame) = consumer.read().await.unwrap() {
read.push((frame.timestamp.as_micros(), frame.data.len() / size_of::<f32>()));
}
read
}
#[tokio::test]
async fn a_missing_packet_leaves_a_hole() {
const FRAMES: usize = 1024;
let stamps = [
Timestamp::from_scale(0, 44_100).unwrap(),
Timestamp::from_scale(2 * FRAMES as u64, 44_100).unwrap(),
];
let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
assert_eq!(read.len(), 4, "unexpected frames: {read:?}");
let before: usize = read[..2].iter().map(|(_, frames)| frames).sum();
assert!((1105..=1120).contains(&before), "unexpected pre-gap audio: {before}");
assert_eq!(read[2].0, stamps[1].as_micros());
let ends_at = read[1].0 + (read[1].1 as u128) * 1_000_000 / 48_000;
let hole = read[2].0 - ends_at;
assert!((23_100..=23_350).contains(&hole), "unexpected hole: {hole} us");
}
#[tokio::test]
async fn a_jump_inside_the_slack_leaves_the_held_samples_alone() {
const FRAMES: usize = 441;
let stamps = [
Timestamp::from_micros(0).unwrap(),
Timestamp::from_micros(11_000).unwrap(),
];
let read = pcm_gaps(44_100, 48_000, FRAMES, &stamps).await;
assert_eq!(read.len(), 2, "unexpected frames: {read:?}");
assert_eq!(read[0].0, 0, "held samples moved with the jump: {read:?}");
}
#[tokio::test]
async fn a_jump_after_a_full_chunk_uses_the_new_packet_timestamp() {
let stamps = [
Timestamp::from_micros(0).unwrap(),
Timestamp::from_micros(21_000).unwrap(),
];
let read = pcm_gaps(44_100, 48_000, 882, &stamps).await;
let mut r = crate::resample::Resampler::new(44_100, 48_000, 1, 882).unwrap();
r.process(&[0.25; 882], stamps[0]).unwrap();
let expected = rewind(stamps[1], r.skipped(), 48_000).unwrap().as_micros();
assert_eq!(read[1].0, expected);
}
#[tokio::test]
async fn a_terminal_jump_leaves_the_held_samples_alone() {
let mut encoder = Encoder::new(&Settings {
dtx: true,
bitrate: Some(moq_net::bandwidth::Rate::from_bps(24_000)),
..Settings::new(48_000, Layout::Mono)
})
.unwrap();
let catalog = encoder.catalog();
let active = encoder.encode(&vec![0.5f32; encoder.frame_size()]).unwrap();
assert!(active.activity.is_active());
let silence = vec![0.0f32; encoder.frame_size()];
let dtx = (0..200)
.map(|_| encoder.encode(&silence).unwrap())
.find(|packet| packet.activity.is_dtx())
.expect("silence should enter Opus DTX");
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
output: Output {
sample_rate: Some(44_100),
..Output::default()
},
..Options::new()
},
)
.await
.unwrap();
let write = |producer: &mut moq_mux::container::Producer<_>, frames: u64, payload: Bytes, keyframe: bool| {
producer
.write(moq_mux::container::Frame {
timestamp: Timestamp::from_scale(frames, 48_000).unwrap(),
duration: None,
payload,
keyframe,
})
.unwrap();
};
write(&mut producer, 0, active.payload, true);
write(&mut producer, 3 * 48_000, Bytes::new(), false);
write(&mut producer, 48_000, dtx.payload, false);
producer.finish().unwrap();
let frame = consumer.read().await.unwrap().expect("decoded frame");
assert_eq!(frame.timestamp.as_micros(), 0, "held samples moved with the jump");
assert!(
frame.activity.is_active(),
"held samples took the terminal packet's label"
);
}
#[tokio::test]
async fn millisecond_stamps_are_not_a_gap() {
const FRAMES: u64 = 1024;
const PACKETS: u64 = 32;
let stamps: Vec<_> = (0..PACKETS)
.map(|packet| Timestamp::from_millis(packet * FRAMES * 1_000 / 44_100).unwrap())
.collect();
let read = pcm_gaps(44_100, 48_000, FRAMES as usize, &stamps).await;
assert_eq!(read.len(), stamps.len() + 1, "unexpected frames: {read:?}");
for pair in read.windows(2) {
let ends_at = pair[0].0 + (pair[0].1 as u128) * 1_000_000 / 48_000;
assert!(
pair[1].0.abs_diff(ends_at) <= 1_100,
"frames at {} and {} do not meet",
pair[0].0,
pair[1].0
);
}
}
#[tokio::test]
async fn a_lost_opus_packet_shorter_than_its_neighbour_is_a_gap() {
let input = Input {
format: Format::F32,
sample_rate: 48_000,
layout: Layout::Mono,
};
let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
let catalog = encoder.catalog();
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
max_age: std::time::Duration::from_secs(1),
..Options::new()
},
)
.await
.unwrap();
let pcm = vec![0.25f32; encoder.frame_size()];
for timestamp in [
Timestamp::from_micros(0).unwrap(),
Timestamp::from_micros(22_500).unwrap(),
Timestamp::from_micros(42_500).unwrap(),
] {
producer
.write(moq_mux::container::Frame {
timestamp,
duration: None,
payload: encoder.encode(&pcm).unwrap().payload,
keyframe: true,
})
.unwrap();
producer.cut(None).unwrap();
}
let first = consumer.read().await.unwrap().expect("decoded frame");
let frames = first.data.len() / size_of::<f32>();
assert!(frames < 960, "the pre-skip should be trimmed, got {frames} frames");
let second = consumer.read().await.unwrap().expect("decoded frame");
assert_eq!(second.timestamp.as_micros(), 22_500);
assert_eq!(second.data.len() / size_of::<f32>(), 960, "pre-skip was reapplied");
let third = consumer.read().await.unwrap().expect("decoded frame after gap");
let second_frames = second.data.len() / size_of::<f32>();
assert_eq!(
third.timestamp,
advance(second.timestamp, second_frames, 48_000).unwrap()
);
}
#[tokio::test]
async fn max_age_is_clamped_to_publisher_retention() {
let broadcast = moq_net::broadcast::Info::new().produce();
let info = hang::container::track_info(hang::catalog::PRIORITY.audio)
.with_max_age(std::time::Duration::from_millis(100));
let _track = broadcast.create_track("audio", info).unwrap();
let subscriber = broadcast.consume();
let catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
let consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
max_age: std::time::Duration::from_millis(500),
..Options::new()
},
)
.await
.unwrap();
assert_eq!(consumer.max_age(), std::time::Duration::from_millis(100));
}
#[tokio::test]
async fn opus_pre_skip_does_not_leave_a_timestamp_hole() {
let input = Input {
format: Format::F32,
sample_rate: 48_000,
layout: Layout::Mono,
};
let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
let catalog = encoder.catalog();
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
max_age: std::time::Duration::from_secs(1),
..Options::new()
},
)
.await
.unwrap();
let pcm = vec![0.25f32; encoder.frame_size()];
for packet in 0..2 {
producer
.write(moq_mux::container::Frame {
timestamp: Timestamp::from_scale(packet * encoder.frame_size() as u64, 48_000).unwrap(),
duration: None,
payload: encoder.encode(&pcm).unwrap().payload,
keyframe: true,
})
.unwrap();
producer.cut(None).unwrap();
}
let first = consumer.read().await.unwrap().expect("first decoded frame");
let second = consumer.read().await.unwrap().expect("second decoded frame");
let first_frames = first.data.len() / size_of::<f32>();
let expected = advance(first.timestamp, first_frames, 48_000).unwrap();
assert_eq!(second.timestamp, expected);
}
#[tokio::test]
async fn a_playhead_event_reapplies_opus_pre_skip() {
let input = Input {
format: Format::F32,
sample_rate: 48_000,
layout: Layout::Mono,
};
let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
let catalog = encoder.catalog();
let frame_size = encoder.frame_size();
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let subscriber = broadcast.consume();
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Legacy(moq_mux::container::Kind::Audio),
);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
max_age: std::time::Duration::from_secs(1),
..Options::new()
},
)
.await
.unwrap();
let pcm = vec![0.25f32; frame_size];
let write = |producer: &mut moq_mux::container::Producer<_>, packet: u64, payload: bytes::Bytes| {
producer
.write(moq_mux::container::Frame {
timestamp: Timestamp::from_scale(packet * frame_size as u64, 48_000).unwrap(),
duration: None,
payload,
keyframe: true,
})
.unwrap();
producer.cut(None).unwrap();
};
write(&mut producer, 0, encoder.encode(&pcm).unwrap().payload);
write(&mut producer, 1, encoder.encode(&pcm).unwrap().payload);
producer.discontinuity().unwrap();
write(&mut producer, 2, encoder.encode(&pcm).unwrap().payload);
producer.finish().unwrap();
let first = consumer.read().await.unwrap().expect("first decoded frame");
let _second = consumer.read().await.unwrap().expect("second decoded frame");
let resumed = consumer.read().await.unwrap().expect("resumed decoded frame");
let first_frames = first.data.len() / size_of::<f32>();
let resumed_frames = resumed.data.len() / size_of::<f32>();
assert!(first_frames < frame_size, "the first epoch trims pre-skip");
assert_eq!(
resumed_frames, first_frames,
"a playhead event reapplies pre-skip without flushing the decoder"
);
}
#[tokio::test]
async fn reads_the_container_the_catalog_declares() {
let broadcast = moq_net::broadcast::Info::new().produce();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let observed = track.clone();
let subscriber = broadcast.consume();
let mut catalog = hang::catalog::AudioConfig::new(hang::catalog::AudioCodec::Pcm, 48_000, 1);
catalog.container = hang::catalog::Container::Loc;
let mut producer = moq_mux::container::Producer::new(
track,
moq_mux::catalog::hang::Container::Loc(moq_mux::container::Kind::Audio),
);
let max_age = std::time::Duration::from_millis(250);
let mut consumer = Consumer::new(
&subscriber,
&catalog,
"audio",
Options {
max_age,
..Options::new()
},
)
.await
.unwrap();
assert_eq!(observed.subscription().unwrap().max_age, max_age);
let samples = [0.25f32, -0.5, 0.75, -1.0];
let payload: Vec<u8> = samples.iter().flat_map(|sample| sample.to_le_bytes()).collect();
producer
.write(moq_mux::container::Frame {
timestamp: Timestamp::ZERO,
duration: None,
payload: payload.into(),
keyframe: true,
})
.unwrap();
let frame = consumer.read().await.unwrap().expect("decoded frame");
assert_eq!(
Format::F32.as_interleaved_f32(&frame.data, 1).unwrap().as_ref(),
samples
);
}
#[tokio::test]
async fn decodes_a_cmaf_framed_track() {
let input = Input {
format: Format::F32,
sample_rate: 48_000,
layout: Layout::Stereo,
};
let mut encoder = Encoder::new(&Settings::new(input.sample_rate, input.layout)).unwrap();
let mut catalog = encoder.catalog();
let pcm = vec![0.0f32; encoder.frame_size() * encoder.codec_channels() as usize];
let packet = encoder.encode(&pcm).unwrap();
let muxer = moq_mux::container::fmp4::Muxer::audio(&catalog).unwrap();
let init = muxer.init().unwrap().expect("an out-of-band codec has an init segment");
catalog.container = hang::catalog::Container::Cmaf { init };
let broadcast = moq_net::broadcast::Info::new().produce();
let subscriber = broadcast.consume();
let track = broadcast
.create_track("audio", hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
let container = moq_mux::catalog::hang::Container::try_from(&catalog).unwrap();
let mut producer = moq_mux::container::Producer::new(track, container);
let mut consumer = Consumer::new(&subscriber, &catalog, "audio", Options::new())
.await
.unwrap();
producer
.write(moq_mux::container::Frame {
timestamp: Timestamp::ZERO,
payload: packet.payload,
keyframe: true,
duration: None,
})
.unwrap();
producer.cut(None).unwrap();
let frame = consumer.read().await.unwrap().expect("decoded frame");
assert_eq!(frame.timestamp.as_micros(), 0);
let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
assert_eq!(samples.len(), (960 - 312) * 2);
}
}