use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use anyhow::Context;
use hang::moq_net;
use moq_mux::catalog::{self, Stream};
use tokio::time::Instant;
use super::args::Args;
use super::output::{Output, Sink, Speaker};
use super::playback::{Kind, Playback, joined};
use super::source::subscribe;
use super::timeline::{AudioTimeline, Presentation, timestamp};
use super::video::Video;
use super::window::Event;
const AUDIO_BUFFER_MIN: Duration = Duration::from_millis(50);
const AUDIO_CHUNK: Duration = Duration::from_millis(20);
const AUDIO_DRAIN_GRACE: Duration = Duration::from_secs(1);
pub(super) struct Media<O: Output> {
pub(super) origin: moq_net::origin::Consumer,
pub(super) broadcast: String,
pub(super) args: Args,
pub(super) video: Arc<Mutex<VecDeque<moq_video::Frame>>>,
pub(super) presentation: Arc<Mutex<Presentation>>,
pub(super) drained: Arc<tokio::sync::Notify>,
pub(super) output: O,
}
impl<O: Output> Media<O> {
pub(super) async fn run(self) {
let output = self.output.clone();
let event = match self.play().await {
Ok(()) => Event::Ended,
Err(err) => Event::Failed(format!("{err:#}")),
};
output.send(event);
}
async fn play(self) -> anyhow::Result<()> {
let source = subscribe(self.origin.clone(), &self.broadcast).await?;
let broadcast = source
.broadcast()
.await
.context("failed to subscribe to the broadcast")?;
let catalog = catalog::Consumer::<()>::new(&broadcast, self.args.catalog_format(&self.broadcast))
.await
.context("failed to subscribe to the catalog")?;
let mut catalogs = catalog.select(self.args.select.selection(None));
let depth = self.args.delay.into_std().max(AUDIO_BUFFER_MIN);
let mut tasks = tokio::task::JoinSet::new();
let mut playback = Playback::default();
let mut speaker = None;
let mut tails = tokio::task::JoinSet::new();
loop {
if playback.done() {
tails.join_all().await;
return Ok(());
}
if playback.pending().is_none() {
tokio::select! {
result = tasks.join_next(), if !tasks.is_empty() => {
let ended = joined(result.expect("guarded by is_empty"))?.map(|(kind, sink)| {
if kind == Kind::Audio {
self.presentation.lock().unwrap().restarted();
}
if let Some(sink) = sink {
tails.spawn(drain(sink, depth));
}
kind
});
playback.ended(ended);
}
_ = tails.join_next(), if !tails.is_empty() => {}
snapshot = catalogs.next(), if playback.following() => {
match snapshot.context("failed to read the catalog")? {
Some(snapshot) => playback.received(snapshot),
None => {
anyhow::ensure!(playback.played, "the catalog contains no playable audio or video renditions");
playback.catalog_ended = true;
}
}
}
}
}
if !playback.playing(Kind::Audio) && tails.is_empty() && speaker.take().is_some() {
self.presentation.lock().unwrap().stopped();
self.drained.notify_one();
}
let Some(snapshot) = playback.pending().cloned() else {
continue;
};
let mut rejected = Vec::new();
if playback.wants(Kind::Video) {
playback.read(Kind::Video);
for (name, config) in snapshot.video.renditions {
let rendition = match source.resolve(config.broadcast.as_ref()).await {
Ok(rendition) => rendition,
Err(err) => {
tracing::warn!(track = name, %err, "cannot resolve video rendition");
rejected.push(format!("video `{name}`: {err}"));
continue;
}
};
let max_age = self.args.delay.into_std();
let opened = async {
let decoder = moq_video::decode::Sink::open(&config, &Default::default()).await?;
let track = rendition.track(&name)?;
let mut subscriber = track
.subscribe(
moq_net::track::Subscription::default()
.with_priority(hang::catalog::PRIORITY.video)
.with_max_age(max_age),
)
.await?;
if let Some(latest) = track.latest() {
subscriber.set_groups(latest..);
}
let format = catalog::hang::Container::try_from(&config)?;
Ok::<_, anyhow::Error>((moq_mux::container::Consumer::new(subscriber, format), decoder))
}
.await;
match opened {
Ok((track, decoder)) => {
tracing::info!(track = name, decoder = decoder.name(), "playing video rendition");
let video = Video {
presentation: self.presentation.clone(),
frames: self.video.clone(),
changed: self.drained.clone(),
output: self.output.clone(),
max_age,
};
tasks.spawn(async move { (Kind::Video, video.run(track, decoder).await.map(|()| None)) });
playback.started(Kind::Video);
break;
}
Err(err) => {
tracing::warn!(track = name, %err, "cannot play video rendition");
rejected.push(format!("video `{name}`: {err}"));
}
}
}
}
if playback.wants(Kind::Audio) {
playback.read(Kind::Audio);
for (name, config) in snapshot.audio.renditions {
let rendition = match source.resolve(config.broadcast.as_ref()).await {
Ok(rendition) => rendition,
Err(err) => {
tracing::warn!(track = name, %err, "cannot resolve audio rendition");
rejected.push(format!("audio `{name}`: {err}"));
continue;
}
};
let mut decode = moq_audio::decode::Options::new();
decode.start = moq_audio::decode::Start::Latest;
decode.max_age = depth;
decode.output.format = moq_audio::Format::F32;
match moq_audio::decode::Consumer::new(&rendition, &config, &name, decode).await {
Ok(consumer) => {
tracing::info!(track = name, decoder = consumer.name(), "playing audio rendition");
if speaker.is_none() {
speaker = Some(self.output.speaker().await?);
}
let audio = AudioPlayback {
speaker: speaker.clone().expect("opened above"),
presentation: self.presentation.clone(),
depth,
changed: self.drained.clone(),
output: self.output.clone(),
};
tasks.spawn(async move { (Kind::Audio, play_audio(consumer, audio).await.map(Some)) });
playback.started(Kind::Audio);
break;
}
Err(err) => {
tracing::warn!(track = name, %err, "cannot play audio rendition");
rejected.push(format!("audio `{name}`: {err}"));
}
}
}
}
anyhow::ensure!(
!tasks.is_empty() || rejected.is_empty(),
"no playable rendition in the catalog: {}",
rejected.join("; ")
);
}
}
}
struct AudioPlayback<O: Output> {
changed: Arc<tokio::sync::Notify>,
speaker: O::Speaker,
presentation: Arc<Mutex<Presentation>>,
depth: Duration,
output: O,
}
async fn play_audio<O: Output>(
mut consumer: moq_audio::decode::Consumer,
playback: AudioPlayback<O>,
) -> anyhow::Result<<O::Speaker as Speaker>::Sink> {
let AudioPlayback {
changed,
speaker,
presentation,
depth,
output,
} = playback;
let sample_rate = consumer.sample_rate();
let layout = consumer.layout();
let channels = layout.channels();
let mut input = moq_audio::playback::Input::default();
input.format = moq_audio::Format::F32;
input.sample_rate = sample_rate;
input.layout = layout;
input.latency = depth;
let mut sink = speaker.sink(input.clone())?;
let stride = channels as usize * size_of::<f32>();
let chunk = ((AUDIO_CHUNK.as_secs_f64() * sample_rate as f64) as usize * stride).max(stride);
let fill_max = (consumer.max_age().as_secs_f64() * sample_rate as f64) as u64;
let silence = vec![0u8; chunk];
let mut timeline = AudioTimeline::default();
let mut dropping = false;
loop {
let frame = match consumer.read().await {
Ok(Some(frame)) => frame,
Ok(None) => break,
Err(err @ moq_audio::Error::Decode(_)) => {
if dropping {
tracing::debug!(%err, "dropping an audio frame");
} else {
tracing::warn!(%err, "dropping an audio frame");
dropping = true;
}
continue;
}
Err(err) => return Err(err.into()),
};
dropping = false;
let samples = frame.data.len() / size_of::<f32>() / channels as usize;
let start = timestamp(frame.timestamp);
let timing = timeline.push(start, samples, sample_rate, fill_max);
if timing.reset_sink {
drop(sink);
presentation.lock().unwrap().restarted();
sink = speaker.sink(input.clone())?;
}
if timing.silence > 0 {
let mut remaining = usize::try_from(timing.silence)
.unwrap_or(usize::MAX / stride)
.saturating_mul(stride);
while remaining > 0 {
if let Some(excess) = sink.buffered().checked_sub(depth) {
tokio::time::sleep(excess).await;
}
let part = remaining.min(silence.len());
sink.write(&silence[..part])?;
remaining -= part;
}
}
for part in frame.data.chunks(chunk) {
if let Some(excess) = sink.buffered().checked_sub(depth) {
tokio::time::sleep(excess).await;
}
sink.write(part)?;
}
let moved = presentation
.lock()
.unwrap()
.audio(timing.end, sink.buffered(), Instant::now().into_std());
if moved {
changed.notify_one();
output.send(Event::Wake);
}
}
Ok(sink)
}
async fn drain(sink: impl Sink, latency: Duration) {
let drain = async {
while let Some(remaining) = sink.buffered().checked_sub(Duration::from_millis(10)) {
tokio::time::sleep(remaining.max(Duration::from_millis(10))).await;
}
};
let _ = tokio::time::timeout(latency + AUDIO_CHUNK + AUDIO_DRAIN_GRACE, drain).await;
}
#[cfg(test)]
mod tests {
use bytes::Bytes;
use hang::catalog::{AudioCodec, AudioConfig};
use moq_mux::catalog::hang::Container;
use super::*;
use crate::play::fake::Recorder;
const SAMPLE_RATE: u32 = 48_000;
const PACKET: u64 = 960;
const PACKET_DURATION: Duration = Duration::from_millis(20);
fn rendition(
broadcast: &moq_net::broadcast::Producer,
catalog: &catalog::Producer,
name: &str,
) -> moq_mux::container::Producer<Container, AudioConfig> {
let track = broadcast
.create_track(name, hang::container::track_info(hang::catalog::PRIORITY.audio))
.unwrap();
catalog
.audio(
track,
Container::Legacy(moq_mux::container::Kind::Audio),
AudioConfig::new(AudioCodec::Pcm, SAMPLE_RATE, 1),
)
.unwrap()
}
fn packet(index: u64, sample: f32) -> moq_mux::container::Frame {
let payload: Vec<u8> = std::iter::repeat_n(sample.to_le_bytes(), PACKET as usize)
.flatten()
.collect();
moq_mux::container::Frame {
timestamp: moq_net::Timestamp::from_scale(index * PACKET, SAMPLE_RATE as u64).unwrap(),
duration: None,
payload: Bytes::from(payload),
keyframe: true,
}
}
fn media(origin: &moq_net::origin::Producer, delay: Duration, output: Recorder) -> Media<Recorder> {
Media {
origin: origin.consume(),
broadcast: "room".to_string(),
args: Args {
catalog_format: None,
delay: delay.into(),
select: Default::default(),
},
video: Default::default(),
presentation: Arc::new(Mutex::new(Presentation::new(delay))),
drained: Default::default(),
output,
}
}
#[tokio::test]
async fn an_audio_rendition_switch_leaves_no_gap() {
tokio::time::pause();
const OLD: f32 = 0.25;
const NEW: f32 = 0.5;
let delay = Duration::from_millis(500);
let origin = moq_tokio::origin::spawn();
let mut broadcast = origin.create_broadcast("room").unwrap();
broadcast.announce(Default::default()).unwrap();
let mut catalog = catalog::Producer::new(&mut broadcast, Default::default()).unwrap();
let recorder = Recorder::default();
let player = tokio::spawn(media(&origin, delay, recorder.clone()).run());
let mut old = rendition(&broadcast, &catalog, "old");
let start = Instant::now();
let mut index = 0;
while index < 50 {
old.write(packet(index, OLD)).unwrap();
index += 1;
tokio::time::sleep_until(start + PACKET_DURATION * index as u32).await;
}
let mut new = rendition(&broadcast, &catalog, "new");
new.write(packet(index, NEW)).unwrap();
index += 1;
old.finish().unwrap();
drop(old);
while index < 100 {
tokio::time::sleep_until(start + PACKET_DURATION * index as u32).await;
new.write(packet(index, NEW)).unwrap();
index += 1;
}
new.finish().unwrap();
drop(new);
catalog.finish().unwrap();
player.await.unwrap();
match recorder.events().pop() {
Some(Event::Ended) => {}
Some(Event::Failed(err)) => panic!("playback failed: {err}"),
_ => panic!("playback never ended"),
}
let played = recorder.played();
let old_end = played.iter().filter(|p| p.sample == OLD).map(|p| p.to).max().unwrap();
let new_start = played.iter().filter(|p| p.sample == NEW).map(|p| p.from).min().unwrap();
let gap = new_start.saturating_duration_since(old_end);
assert!(gap < AUDIO_CHUNK, "the switch went silent for {gap:?}");
}
#[tokio::test]
async fn a_wide_delay_observes_the_whole_tune_in_burst() {
tokio::time::pause();
let delay = Duration::from_secs(2);
let origin = moq_tokio::origin::spawn();
let broadcast = origin.create_broadcast("room").unwrap();
let track = broadcast
.create_track("video", hang::container::track_info(hang::catalog::PRIORITY.video))
.unwrap();
let mut producer = moq_mux::container::Producer::new(track, Container::Legacy(moq_mux::container::Kind::Data));
let mut config = moq_video::encode::Config::new(64, 64, moq_video::Rate::new(30, 1).unwrap());
config.kind = moq_video::encode::Kind::Software;
config.gop = moq_video::encode::Gop::Keyframe { interval: 120 };
let mut encoder = moq_video::encode::Encoder::new(&config).unwrap();
for index in 0..=60 {
let surface = moq_video::Surface::rgba(&vec![128; 64 * 64 * 4], moq_video::Size::new(64, 64)).unwrap();
let frame = moq_video::Frame::new(surface, moq_net::Timestamp::from_millis(index * 33).unwrap());
for encoded in encoder.encode(&frame).unwrap() {
producer
.write(moq_mux::container::Frame {
timestamp: encoded.timestamp,
duration: None,
payload: encoded.payload,
keyframe: index == 0,
})
.unwrap();
}
}
producer.finish().unwrap();
let catalog = hang::catalog::VideoConfig::new(hang::catalog::H264 {
inline: true,
profile: 0x42,
constraints: 0,
level: 30,
});
let mut options = moq_video::decode::Options::new();
options.decoder.kind = moq_video::decode::Kind::Software;
options.max_age = delay;
let decoder = moq_video::decode::Sink::open(&catalog, &options.decoder).await.unwrap();
let subscriber = broadcast
.consume()
.track("video")
.unwrap()
.subscribe(moq_net::track::Subscription::default().with_max_age(delay))
.await
.unwrap();
let track = moq_mux::container::Consumer::new(subscriber, Container::try_from(&catalog).unwrap());
let recorder = Recorder::default();
let media = media(&origin, delay, recorder.clone());
let now = Instant::now();
let last = moq_net::Timestamp::from_millis(60 * 33).unwrap();
let task = tokio::spawn(
Video {
presentation: media.presentation.clone(),
frames: media.video.clone(),
changed: media.drained.clone(),
output: recorder.clone(),
max_age: delay,
}
.run(track, decoder),
);
loop {
tokio::task::yield_now().await;
if media.video.lock().unwrap().len() == 30
|| media.presentation.lock().unwrap().due(last) == Some((now + delay).into_std())
{
break;
}
}
let due = media.presentation.lock().unwrap().due(last);
task.abort();
let _ = task.await;
assert_eq!(
due,
Some((now + delay).into_std()),
"the decoder did not observe the live edge"
);
assert!(recorder.present(&media).is_none(), "nothing is due before its delay");
}
}