use std::time::Duration;
use clap::ValueEnum;
use hang::moq_net;
use moq_mux::catalog::CatalogFormat;
use tokio::io::AsyncWriteExt;
#[derive(Clone, Copy)]
pub enum SubscribeFormat {
Fmp4,
Mkv,
Ts,
Flv,
}
#[derive(ValueEnum, Clone, Copy)]
pub enum CatalogFormatArg {
Hang,
#[value(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(Clone)]
pub struct SubscribeArgs {
pub format: SubscribeFormat,
pub max_latency: Duration,
pub fragment_duration: Option<Duration>,
pub catalog: Option<CatalogFormatArg>,
}
impl SubscribeArgs {
pub fn catalog_format(&self, broadcast: &str) -> CatalogFormat {
self.catalog
.map(Into::into)
.or_else(|| CatalogFormat::detect(broadcast))
.unwrap_or_default()
}
}
pub struct Subscribe {
broadcast: moq_net::BroadcastConsumer,
catalog: CatalogFormat,
args: SubscribeArgs,
}
impl Subscribe {
pub fn new(broadcast: moq_net::BroadcastConsumer, catalog: CatalogFormat, args: SubscribeArgs) -> Self {
Self {
broadcast,
catalog,
args,
}
}
pub async fn run(self) -> anyhow::Result<()> {
match self.args.format {
SubscribeFormat::Fmp4 => self.run_fmp4().await,
SubscribeFormat::Mkv => self.run_mkv().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 catalog = moq_mux::catalog::Consumer::<()>::new(&self.broadcast, self.catalog)?;
let mut fmp4 = moq_mux::container::fmp4::Export::new(self.broadcast, catalog)
.with_latency(self.args.max_latency)
.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 catalog = moq_mux::catalog::Consumer::<()>::new(&self.broadcast, self.catalog)?;
let mut mkv = moq_mux::container::mkv::Export::new(self.broadcast, catalog)
.with_latency(self.args.max_latency)
.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_ts(self) -> anyhow::Result<()> {
let mut stdout = tokio::io::stdout();
let mut ts =
moq_mux::container::ts::Export::with_ts(self.broadcast, self.catalog)?.with_latency(self.args.max_latency);
while let Some(frame) = ts.next().await? {
stdout.write_all(&frame.payload).await?;
stdout.flush().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.broadcast, self.catalog)?
.with_latency(self.args.max_latency);
while let Some(chunk) = flv.next().await? {
stdout.write_all(&chunk).await?;
stdout.flush().await?;
}
Ok(())
}
}