use std::time::Duration;
use hang::catalog::{AudioCodecKind, VideoCodecKind};
use moq_mux::catalog::{self, CatalogFormat, Stream};
use moq_mux::select;
use tokio::io::AsyncWriteExt;
#[derive(Clone, Copy)]
pub enum SubscribeFormat {
Fmp4,
Mkv,
H264,
H265,
Ts,
Flv,
}
#[derive(usage::ValueEnum, Clone, Copy)]
pub enum CatalogFormatArg {
Hang,
#[usage(name = "hangz")]
HangZ,
Msf,
}
impl From<CatalogFormatArg> for CatalogFormat {
fn from(format: CatalogFormatArg) -> Self {
match format {
CatalogFormatArg::Hang => Self::Hang,
CatalogFormatArg::HangZ => Self::HangZ,
CatalogFormatArg::Msf => Self::Msf,
}
}
}
#[derive(usage::ValueEnum, Clone, Copy)]
pub enum VideoCodecArg {
H264,
H265,
Vp8,
Vp9,
Av1,
}
impl From<VideoCodecArg> for VideoCodecKind {
fn from(value: VideoCodecArg) -> Self {
match value {
VideoCodecArg::H264 => Self::H264,
VideoCodecArg::H265 => Self::H265,
VideoCodecArg::Vp8 => Self::VP8,
VideoCodecArg::Vp9 => Self::VP9,
VideoCodecArg::Av1 => Self::AV1,
}
}
}
#[derive(usage::ValueEnum, Clone, Copy)]
pub enum AudioCodecArg {
Aac,
Opus,
Pcm,
}
impl From<AudioCodecArg> for AudioCodecKind {
fn from(value: AudioCodecArg) -> Self {
match value {
AudioCodecArg::Aac => Self::AAC,
AudioCodecArg::Opus => Self::Opus,
AudioCodecArg::Pcm => Self::Pcm,
}
}
}
#[derive(usage::Args, Clone, Default)]
#[usage(unknown_flags = "error", args_override_self = false)]
pub struct SelectArgs {
#[usage(long)]
pub video_name: Option<String>,
#[usage(long, value_enum)]
pub video_codec: Option<VideoCodecArg>,
#[usage(long)]
pub audio_name: Option<String>,
#[usage(long, value_enum)]
pub audio_codec: Option<AudioCodecArg>,
}
impl SelectArgs {
pub(crate) fn selection(&self, force: Option<VideoCodecKind>) -> select::Broadcast {
let mut video = select::Video::default();
if let Some(name) = &self.video_name {
video = video.name(name);
}
if let Some(codec) = force.or_else(|| self.video_codec.map(Into::into)) {
video = video.codec(codec);
}
let mut audio = select::Audio::default();
if let Some(name) = &self.audio_name {
audio = audio.name(name);
}
if let Some(codec) = self.audio_codec {
audio = audio.codec(codec.into());
}
select::Broadcast::default().video(video).audio(audio)
}
}
#[derive(Clone)]
pub struct SubscribeArgs {
pub format: SubscribeFormat,
pub max_age: Duration,
pub fragment_duration: Option<Duration>,
pub mux_rate: Option<u64>,
pub catalog: Option<CatalogFormatArg>,
pub select: SelectArgs,
}
impl SubscribeArgs {
pub fn catalog_format(&self, broadcast: &str) -> CatalogFormat {
self.catalog
.map(Into::into)
.or_else(|| CatalogFormat::detect(broadcast))
.unwrap_or_default()
}
fn format_codec(&self) -> Option<VideoCodecKind> {
match self.format {
SubscribeFormat::H264 => Some(VideoCodecKind::H264),
SubscribeFormat::H265 => Some(VideoCodecKind::H265),
SubscribeFormat::Fmp4 | SubscribeFormat::Mkv | SubscribeFormat::Ts | SubscribeFormat::Flv => None,
}
}
fn selection(&self) -> anyhow::Result<select::Broadcast> {
let user_codec = self.select.video_codec.map(VideoCodecKind::from);
let codec = match (self.format_codec(), user_codec) {
(Some(fmt), Some(user)) if fmt != user => {
anyhow::bail!(
"the output format implies video codec {fmt:?}, but --video-codec {user:?} was passed; \
remove --video-codec or pick a matching format"
);
}
(Some(fmt), _) => Some(fmt),
(None, user) => user,
};
Ok(self.select.selection(codec))
}
}
pub struct Subscribe {
source: moq_mux::Source,
catalog: CatalogFormat,
args: SubscribeArgs,
}
impl Subscribe {
pub fn new(source: moq_mux::Source, catalog: CatalogFormat, args: SubscribeArgs) -> Self {
Self { source, catalog, args }
}
async fn stream(&self) -> anyhow::Result<catalog::Select<catalog::Consumer>> {
let consumer = self.source.catalog(self.catalog).await?;
Ok(consumer.select(self.args.selection()?))
}
pub async fn run(self) -> anyhow::Result<()> {
match self.args.format {
SubscribeFormat::Fmp4 => self.run_fmp4().await,
SubscribeFormat::Mkv => self.run_mkv().await,
SubscribeFormat::H264 => self.run_h264().await,
SubscribeFormat::H265 => self.run_h265().await,
SubscribeFormat::Ts => self.run_ts().await,
SubscribeFormat::Flv => self.run_flv().await,
}
}
async fn run_fmp4(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let stream = self.stream().await?;
let mut fmp4 = moq_mux::container::fmp4::Export::new(self.source, stream)
.with_max_age(self.args.max_age)
.with_fragment_duration(self.args.fragment_duration);
while let Some(chunk) = fmp4.next().await? {
stdout.write_all(&chunk).await?;
stdout.flush().await?;
}
Ok(())
}
async fn run_mkv(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let stream = self.stream().await?;
let mut mkv = moq_mux::container::mkv::Export::new(self.source, stream)
.with_max_age(self.args.max_age)
.with_fragment_duration(self.args.fragment_duration);
while let Some(chunk) = mkv.next().await? {
stdout.write_all(&chunk).await?;
stdout.flush().await?;
}
Ok(())
}
async fn run_h264(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let stream = self.stream().await?;
let mut h264 = moq_mux::codec::h264::Export::new(self.source, stream).with_max_age(self.args.max_age);
while let Some(chunk) = h264.next().await? {
stdout.write_all(&chunk).await?;
stdout.flush().await?;
}
Ok(())
}
async fn run_h265(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let stream = self.stream().await?;
let mut h265 = moq_mux::codec::h265::Export::new(self.source, stream).with_max_age(self.args.max_age);
while let Some(chunk) = h265.next().await? {
stdout.write_all(&chunk).await?;
stdout.flush().await?;
}
Ok(())
}
async fn run_ts(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let mut ts = moq_mux::container::ts::Export::with_ts(self.source, self.catalog)
.await?
.with_max_age(self.args.max_age);
if let Some(mux_rate) = self.args.mux_rate {
ts = ts.with_mux_rate(mux_rate);
}
let mut delivery = Delivery::new(self.args.max_age);
loop {
let mut waited = false;
let frame = hang::moq_net::kio::wait(|waiter| match ts.poll_next(waiter) {
std::task::Poll::Pending => {
waited = true;
std::task::Poll::Pending
}
ready => ready,
})
.await?;
let Some(frame) = frame else { break };
delivery.update(&frame, ts.discontinuity());
delivery.deliver(&frame, waited, &mut stdout).await?;
}
Ok(())
}
async fn run_flv(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let mut flv = moq_mux::container::flv::Export::with_catalog_format(self.source, self.catalog)
.await?
.with_max_age(self.args.max_age);
while let Some(chunk) = flv.next().await? {
stdout.write_all(&chunk).await?;
stdout.flush().await?;
}
Ok(())
}
}
struct Delivery {
discontinuity: u64,
pacer: moq_mux::Pacer,
lead: Duration,
arrived: tokio::time::Instant,
}
impl Delivery {
fn new(lead: Duration) -> Self {
Self {
pacer: moq_mux::Pacer::default().with_lead(lead),
discontinuity: 0,
lead,
arrived: tokio::time::Instant::now(),
}
}
fn update(&mut self, frame: &moq_mux::container::Frame, discontinuity: u64) {
if discontinuity == self.discontinuity {
return;
}
self.discontinuity = discontinuity;
self.arrived = tokio::time::Instant::now();
self.pacer.hurry(frame.timestamp, self.arrived.into_std());
}
async fn deliver(
&mut self,
frame: &moq_mux::container::Frame,
waited: bool,
out: &mut (impl tokio::io::AsyncWrite + Unpin),
) -> anyhow::Result<()> {
let now = tokio::time::Instant::now();
if waited {
self.arrived = now;
}
let budget = self.lead + self.pacer.slack();
let mut send_at = self.pacer.pace(frame.timestamp, now.into_std());
if send_at.saturating_duration_since(self.arrived.into_std()) > budget {
send_at = self.pacer.hurry(frame.timestamp, now.into_std());
self.arrived = now;
}
let send_at = tokio::time::Instant::from_std(send_at);
tokio::time::sleep_until(send_at).await;
if send_at > now {
self.arrived = send_at;
}
out.write_all(&frame.payload).await?;
out.flush().await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use hang::moq_net::{Timescale, Timestamp};
fn frame(value: u64, scale: Timescale) -> moq_mux::container::Frame {
moq_mux::container::Frame {
timestamp: Timestamp::new(value, scale).unwrap(),
duration: None,
payload: bytes::Bytes::from_static(&[0x47; 188]),
keyframe: false,
}
}
#[tokio::test(start_paused = true)]
async fn ts_frames_are_paced_on_the_media_clock() {
let mut delivery = Delivery::new(Duration::from_millis(500));
let mut out = Vec::new();
let start = tokio::time::Instant::now();
delivery
.deliver(&frame(0, Timescale::MICRO), true, &mut out)
.await
.unwrap();
assert_eq!(start.elapsed(), Duration::ZERO);
delivery
.deliver(&frame(25_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
assert_eq!(start.elapsed(), Duration::from_millis(25));
delivery
.deliver(&frame(3_600, Timescale::new(90_000).unwrap()), true, &mut out)
.await
.unwrap();
assert_eq!(start.elapsed(), Duration::from_millis(40));
assert_eq!(out.len(), 3 * 188, "every payload was written");
}
#[tokio::test(start_paused = true)]
async fn a_sink_that_cannot_keep_up_sheds_the_lag() {
let mut delivery = Delivery::new(Duration::from_millis(500));
let mut out = Vec::new();
let start = tokio::time::Instant::now();
delivery
.deliver(&frame(0, Timescale::MICRO), true, &mut out)
.await
.unwrap();
tokio::time::advance(Duration::from_secs(2)).await;
delivery
.deliver(&frame(400_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
assert_eq!(
start.elapsed(),
Duration::from_secs(2),
"an overdue frame writes at once: the overshoot hurries and makes it the live edge"
);
delivery
.deliver(&frame(600_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
assert_eq!(start.elapsed(), Duration::from_secs(2) + Duration::from_millis(200));
for slot in 1..=3u64 {
tokio::time::advance(Duration::from_secs(2)).await;
let before = start.elapsed();
delivery
.deliver(&frame(600_000 + slot * 25_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
assert_eq!(start.elapsed(), before, "stall {slot} must shed, not pace");
}
}
#[tokio::test(start_paused = true)]
async fn pacing_resumes_after_a_hurry_without_a_wait() {
let mut delivery = Delivery::new(Duration::from_millis(500));
let mut out = Vec::new();
let start = tokio::time::Instant::now();
delivery
.deliver(&frame(0, Timescale::MICRO), true, &mut out)
.await
.unwrap();
tokio::time::advance(Duration::from_secs(2)).await;
delivery
.deliver(&frame(400_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
let hurried = start.elapsed();
assert_eq!(hurried, Duration::from_secs(2), "the overshoot sheds");
for slot in 1..=8u64 {
delivery
.deliver(&frame(400_000 + slot * 25_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
assert_eq!(
start.elapsed(),
hurried + Duration::from_millis(slot * 25),
"slot {slot} must be paced, not shed"
);
}
}
#[tokio::test(start_paused = true)]
async fn a_buffered_producer_keeps_pacing() {
let mut delivery = Delivery::new(Duration::from_millis(500));
let mut out = Vec::new();
let start = tokio::time::Instant::now();
delivery
.deliver(&frame(0, Timescale::MICRO), true, &mut out)
.await
.unwrap();
for slot in 1..=40u64 {
delivery
.deliver(&frame(slot * 25_000, Timescale::MICRO), false, &mut out)
.await
.unwrap();
assert_eq!(
start.elapsed(),
Duration::from_millis(slot * 25),
"slot {slot} must be paced, not shed"
);
}
}
#[tokio::test(start_paused = true)]
async fn a_rewind_re_anchors_the_pacer() {
let mut delivery = Delivery::new(Duration::from_millis(500));
let mut out = tokio::io::sink();
delivery
.deliver(&frame(10_000_000, Timescale::MICRO), true, &mut out)
.await
.unwrap();
let first = frame(0, Timescale::MICRO);
let now = tokio::time::Instant::now();
delivery.update(&first, 1);
delivery.deliver(&first, false, &mut out).await.unwrap();
assert_eq!(now.elapsed(), Duration::ZERO);
let next = frame(40_000, Timescale::MICRO);
delivery.update(&next, 1);
delivery.deliver(&next, false, &mut out).await.unwrap();
assert_eq!(now.elapsed(), Duration::from_millis(40));
}
}