use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use anyhow::Context;
use hang::moq_net;
use moq_mux::catalog::{self, Stream};
use winit::event_loop::EventLoopProxy;
use super::args::Args;
use super::playback::{Kind, Playback, joined};
use super::source::subscribe;
use super::timeline::{AudioTimeline, Presentation, timestamp};
use super::window::Event;
const MAX_VIDEO_FRAMES: usize = 30;
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 {
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) proxy: EventLoopProxy<Event>,
}
impl Media {
pub(super) async fn run(self) {
let proxy = self.proxy.clone();
let event = match self.play().await {
Ok(()) => Event::Ended,
Err(err) => Event::Failed(format!("{err:#}")),
};
let _ = proxy.send_event(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 mut tasks = tokio::task::JoinSet::new();
let mut playback = Playback::default();
loop {
if playback.done() {
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"))?;
if ended == Some(Kind::Audio) {
self.presentation.lock().unwrap().stopped();
}
playback.ended(ended);
}
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;
}
}
}
}
}
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 mut decode = moq_video::decode::Options::new();
decode.start = moq_video::decode::Start::Latest;
decode.max_age = self.args.delay.into_std();
match moq_video::decode::Consumer::new(&rendition, &config, &name, decode).await {
Ok(consumer) => {
tracing::info!(track = name, decoder = consumer.name(), "playing video rendition");
let presentation = self.presentation.clone();
let video = self.video.clone();
let drained = self.drained.clone();
let proxy = self.proxy.clone();
tasks.spawn(async move {
(
Kind::Video,
play_video(consumer, presentation, video, drained, proxy).await,
)
});
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 depth = self.args.delay.into_std().max(AUDIO_BUFFER_MIN);
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, "playing audio rendition");
let audio = AudioPlayback {
presentation: self.presentation.clone(),
depth,
proxy: self.proxy.clone(),
};
tasks.spawn(async move { (Kind::Audio, play_audio(consumer, audio).await) });
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("; ")
);
}
}
}
async fn play_video(
mut consumer: moq_video::decode::Consumer,
presentation: Arc<Mutex<Presentation>>,
video: Arc<Mutex<VecDeque<moq_video::Frame>>>,
drained: Arc<tokio::sync::Notify>,
proxy: EventLoopProxy<Event>,
) -> anyhow::Result<()> {
while let Some(frame) = consumer.read().await? {
if presentation.lock().unwrap().video(frame.timestamp, Instant::now()) {
let _ = proxy.send_event(Event::Wake);
}
while video.lock().unwrap().len() >= MAX_VIDEO_FRAMES {
drained.notified().await;
}
video.lock().unwrap().push_back(frame);
let _ = proxy.send_event(Event::Wake);
}
Ok(())
}
struct AudioPlayback {
presentation: Arc<Mutex<Presentation>>,
depth: Duration,
proxy: EventLoopProxy<Event>,
}
async fn play_audio(mut consumer: moq_audio::decode::Consumer, playback: AudioPlayback) -> anyhow::Result<()> {
let AudioPlayback {
presentation,
depth,
proxy,
} = playback;
let sample_rate = consumer.sample_rate();
let layout = consumer.layout();
let channels = layout.channels();
let engine = moq_audio::playback::Engine::open(Default::default()).await?;
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 = engine.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 = engine.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());
let _ = 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;
}
let _ = sink.write(part)?;
}
let moved = presentation
.lock()
.unwrap()
.audio(timing.end, sink.buffered(), Instant::now());
if moved {
let _ = proxy.send_event(Event::Wake);
}
}
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(depth + AUDIO_CHUNK + AUDIO_DRAIN_GRACE, drain).await;
Ok(())
}