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, Clock, timestamp};
use super::window::Event;
const MAX_VIDEO_FRAMES: usize = 30;
const AUDIO_BUFFER_MAX: Duration = Duration::from_secs(1);
const AUDIO_DRAIN_MAX: Duration = Duration::from_secs(4);
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) audio_clock: Arc<Mutex<Option<Clock>>>,
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() => {
playback.ended(joined(result.expect("guarded by 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;
}
}
}
}
}
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::Config::new();
decode.start = moq_video::decode::Start::Latest;
decode.latency_max = Some(self.args.latency_max);
match moq_video::decode::Consumer::new(&rendition, &config, &name, decode).await {
Ok(consumer) => {
tracing::info!(track = name, decoder = consumer.name(), "playing video rendition");
let video = self.video.clone();
let drained = self.drained.clone();
let proxy = self.proxy.clone();
tasks
.spawn(async move { (Kind::Video, play_video(consumer, 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 mut decode = moq_audio::decode::Config::new();
decode.start = moq_audio::decode::Start::Latest;
decode.latency_max = Some(self.args.latency_max);
decode.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 clock = self.audio_clock.clone();
let proxy = self.proxy.clone();
tasks.spawn(async move { (Kind::Audio, play_audio(consumer, clock, proxy).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,
video: Arc<Mutex<VecDeque<moq_video::Frame>>>,
drained: Arc<tokio::sync::Notify>,
proxy: EventLoopProxy<Event>,
) -> anyhow::Result<()> {
while let Some(frame) = consumer.read().await? {
while video.lock().unwrap().len() >= MAX_VIDEO_FRAMES {
drained.notified().await;
}
video.lock().unwrap().push_back(frame);
let _ = proxy.send_event(Event::Wake);
}
Ok(())
}
async fn play_audio(
mut consumer: moq_audio::decode::Consumer,
clock: Arc<Mutex<Option<Clock>>>,
proxy: EventLoopProxy<Event>,
) -> anyhow::Result<()> {
let sample_rate = consumer.sample_rate();
let channels = consumer.channels();
let engine = moq_audio::playback::Engine::open(Default::default()).await?;
let input = moq_audio::playback::Input {
format: moq_audio::Format::F32,
sample_rate,
channels,
};
let mut sink = engine.sink(input.clone())?;
let stride = channels as usize * size_of::<f32>();
let chunk = (sample_rate as usize * stride).max(stride);
let fill_max = (consumer.latency_max().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);
*clock.lock().unwrap() = None;
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(AUDIO_BUFFER_MAX) {
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(AUDIO_BUFFER_MAX) {
tokio::time::sleep(excess).await;
}
sink.write(part)?;
}
let previous = clock.lock().unwrap().replace(Clock {
media: timing.end.saturating_sub(sink.buffered()),
wall: Instant::now(),
});
if previous.is_none() {
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(AUDIO_DRAIN_MAX, drain).await;
Ok(())
}