use anyhow::Context;
use hang::moq_net;
use moq_mux::container::{flv, fmp4, ts};
#[derive(Clone, Copy)]
pub enum PublishFormat {
Avc3,
Fmp4,
Ts,
Flv,
}
#[cfg(feature = "capture")]
#[derive(usage::ValueEnum, Clone, Copy, Default)]
pub enum VideoCodec {
#[default]
H264,
H265,
}
#[cfg(feature = "capture")]
impl From<VideoCodec> for moq_video::encode::Codec {
fn from(codec: VideoCodec) -> Self {
match codec {
VideoCodec::H264 => moq_video::encode::Codec::H264,
VideoCodec::H265 => moq_video::encode::Codec::H265,
}
}
}
#[cfg(feature = "capture")]
#[derive(usage::Args, Clone)]
#[usage(unknown_flags = "error", args_override_self = false)]
#[usage(group("video-source"))]
#[usage(group("audio-source"))]
pub struct CaptureArgs {
#[usage(long, group = "video-source")]
pub camera: Option<Option<String>>,
#[usage(long, group = "video-source", alias = "screen")]
pub display: Option<Option<String>>,
#[usage(long, group = "video-source")]
pub window: Option<String>,
#[usage(long, group = "video-source")]
pub app: Option<String>,
#[usage(long)]
pub no_cursor: bool,
#[usage(long)]
pub width: Option<u32>,
#[usage(long)]
pub height: Option<u32>,
#[usage(long)]
pub fps: Option<u32>,
#[usage(long)]
pub bitrate: Option<u64>,
#[usage(long, value_enum, default = "h264")]
pub codec: VideoCodec,
#[usage(long, conflicts = "--software")]
pub hardware: bool,
#[usage(long)]
pub software: bool,
#[usage(long, group = "audio-source")]
pub microphone: Option<Option<String>>,
#[usage(long, group = "audio-source")]
pub system_audio: bool,
#[usage(long)]
pub audio_bitrate: Option<u32>,
#[usage(long, conflicts("--no-audio", "--camera", "--display", "--window", "--app"))]
pub no_video: bool,
#[usage(long, conflicts("--microphone", "--system-audio"))]
pub no_audio: bool,
}
fn log_stats(latest: Option<&ts::Stats>, previous: Option<&ts::Stats>) {
let Some(latest) = latest else {
return;
};
for (pid, stream) in &latest.streams {
if previous.and_then(|previous| previous.streams.get(pid)) == Some(stream) {
continue;
}
tracing::info!(
pid = *pid,
track = stream.track,
resyncs = stream.resyncs,
discarded = stream.discarded,
unconfirmed = stream.unconfirmed,
"audio frame sync lost"
);
}
}
enum PublishDecoder {
Avc3 {
split: Box<moq_mux::codec::h264::Split>,
import: Box<moq_mux::codec::h264::Import>,
},
Fmp4(Box<fmp4::Import>),
Ts(Box<ts::Import<ts::Ext>>),
Flv(Box<flv::Import>),
}
impl PublishDecoder {
fn decode_chunk(&mut self, chunk: &[u8]) -> anyhow::Result<()> {
match self {
Self::Avc3 { split, import } => {
let frames = split.decode(chunk, None)?;
import.decode(frames)?;
}
Self::Fmp4(d) => d.decode(chunk)?,
Self::Ts(d) => d.decode(chunk)?,
Self::Flv(d) => d.decode(chunk)?,
}
Ok(())
}
fn stats(&self) -> Option<ts::Stats> {
match self {
Self::Ts(d) => Some(d.stats()),
Self::Avc3 { .. } | Self::Fmp4(_) | Self::Flv(_) => None,
}
}
fn finish(&mut self) -> anyhow::Result<()> {
match self {
Self::Avc3 { split, import } => {
let tail = split.flush(None)?;
import.decode(tail)?;
import.finish()?;
}
Self::Fmp4(d) => d.finish()?,
Self::Ts(d) => d.finish()?,
Self::Flv(d) => d.finish()?,
}
Ok(())
}
fn abort(self, err: moq_net::Error) {
match self {
Self::Avc3 { import, .. } => import.abort(err),
Self::Fmp4(d) => d.abort(err),
Self::Ts(d) => d.abort(err),
Self::Flv(d) => d.abort(err),
}
}
}
enum PublishCatalog {
Media(moq_mux::catalog::Producer),
Ts(moq_mux::catalog::Producer<ts::Ext>),
}
impl PublishCatalog {
fn finish(&mut self) -> anyhow::Result<()> {
match self {
Self::Media(catalog) => catalog.finish()?,
Self::Ts(catalog) => catalog.finish()?,
}
Ok(())
}
}
#[allow(clippy::large_enum_variant)]
enum Source {
Stream {
decoder: PublishDecoder,
catalog: PublishCatalog,
},
#[cfg(feature = "capture")]
Capture {
catalog: moq_mux::catalog::Producer,
video: Option<(moq_video::capture::Config, moq_video::encode::Options)>,
audio: Option<(moq_audio::capture::Config, moq_audio::encode::Options)>,
},
}
pub struct Publish {
source: Source,
#[cfg_attr(not(feature = "capture"), allow(dead_code))]
broadcast: moq_net::broadcast::Producer,
}
impl Publish {
pub fn new(
mut broadcast: moq_net::broadcast::Producer,
format: &PublishFormat,
config: moq_mux::catalog::Config,
) -> anyhow::Result<Self> {
if let PublishFormat::Ts = format {
let config = config.with_catalog(moq_mux::catalog::hang::Catalog::<ts::Ext>::default());
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config)?;
let ts = ts::Import::new(broadcast.clone(), catalog.reserve()).live();
return Ok(Self {
source: Source::Stream {
decoder: PublishDecoder::Ts(Box::new(ts)),
catalog: PublishCatalog::Ts(catalog),
},
broadcast,
});
}
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config)?;
let decoder = match format {
PublishFormat::Avc3 => {
let track = broadcast.unique_track(".avc3", catalog.track_info(hang::catalog::PRIORITY.video))?;
let import = moq_mux::codec::h264::Import::new(track, catalog.reserve(), Default::default())?;
let split = Box::new(moq_mux::codec::h264::Split::new());
PublishDecoder::Avc3 {
split,
import: Box::new(import),
}
}
PublishFormat::Fmp4 => {
let fmp4 = fmp4::Import::new(broadcast.clone(), catalog.reserve()).live();
PublishDecoder::Fmp4(Box::new(fmp4))
}
PublishFormat::Ts => unreachable!("TS is handled above with the mpegts catalog extension"),
PublishFormat::Flv => {
let flv = flv::Import::new(broadcast.clone(), catalog.reserve()).live();
PublishDecoder::Flv(Box::new(flv))
}
};
Ok(Self {
source: Source::Stream {
decoder,
catalog: PublishCatalog::Media(catalog),
},
broadcast,
})
}
#[cfg(feature = "capture")]
pub fn capture(
mut broadcast: moq_net::broadcast::Producer,
args: &CaptureArgs,
bandwidth: moq_net::bandwidth::Allocator,
max_age: Option<std::time::Duration>,
) -> anyhow::Result<Self> {
let config = moq_mux::catalog::Config::default().with_max_age(max_age);
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config)?;
let video = if args.no_video {
None
} else {
Some((args.video_config()?, args.video_encode(bandwidth.clone())))
};
let audio = (!args.no_audio).then(|| (args.audio_config(), args.audio_encode(bandwidth)));
anyhow::ensure!(video.is_some() || audio.is_some(), "nothing to capture");
Ok(Self {
source: Source::Capture { catalog, video, audio },
broadcast,
})
}
pub fn announce(&self) -> anyhow::Result<()> {
self.broadcast
.announce(Default::default())
.context("failed to announce broadcast")
}
pub async fn run(self) -> anyhow::Result<()> {
match self.source {
Source::Stream { decoder, catalog } => decode(decoder, catalog, tokio::io::stdin()).await,
#[cfg(feature = "capture")]
Source::Capture { catalog, video, audio } => {
let clock = catalog.clock();
let video_fut = {
let broadcast = self.broadcast.clone();
let catalog = catalog.clone();
async move {
match video {
Some((config, encode)) => {
moq_video::encode::publish_capture(broadcast, catalog, config, encode, clock)
.await
.map_err(anyhow::Error::from)
}
None => Ok(()),
}
}
};
let audio_fut = {
let broadcast = self.broadcast.clone();
async move {
match audio {
Some((config, encode)) => {
let mut options = moq_audio::encode::PublicationOptions::default();
options.capture = config;
options.encode = encode;
options.clock = clock;
moq_audio::encode::publish_capture(broadcast, catalog, options)
.await
.map_err(anyhow::Error::from)
}
None => Ok(()),
}
}
};
tokio::try_join!(video_fut, audio_fut)?;
Ok(())
}
}
}
}
async fn decode(
mut decoder: PublishDecoder,
mut catalog: PublishCatalog,
mut input: impl tokio::io::AsyncRead + Unpin,
) -> anyhow::Result<()> {
let mut buffer = bytes::BytesMut::new();
let mut reported = decoder.stats();
let result: anyhow::Result<()> = async {
loop {
buffer.clear();
let n = tokio::io::AsyncReadExt::read_buf(&mut input, &mut buffer).await?;
if n == 0 {
return Ok(()); }
decoder.decode_chunk(&buffer)?;
let latest = decoder.stats();
if latest != reported {
log_stats(latest.as_ref(), reported.as_ref());
reported = latest;
}
}
}
.await;
let outcome = result.and_then(|()| decoder.finish());
log_stats(decoder.stats().as_ref(), reported.as_ref());
match outcome {
Ok(()) => catalog.finish(),
Err(err) => {
decoder.abort(moq_net::Error::Transport(err.to_string()));
Err(err)
}
}
}
#[cfg(feature = "capture")]
impl CaptureArgs {
fn video_source(&self) -> moq_video::capture::Source {
use moq_video::capture::Source;
if let Some(id) = &self.window {
Source::Window(id.clone())
} else if let Some(id) = &self.app {
Source::App(id.clone())
} else if let Some(index) = &self.display {
Source::Display(index.clone())
} else {
Source::Camera(self.camera.clone().flatten())
}
}
fn video_config(&self) -> anyhow::Result<moq_video::capture::Config> {
let mut config = moq_video::capture::Config::default();
config.source = self.video_source();
config.width = self.width;
config.height = self.height;
config.framerate = self
.fps
.map(|fps| moq_video::Rate::new(fps, 1))
.transpose()
.map_err(anyhow::Error::from)?;
config.cursor = !self.no_cursor;
Ok(config)
}
fn video_encode(&self, bandwidth: moq_net::bandwidth::Allocator) -> moq_video::encode::Options {
let mut options = moq_video::encode::Options::default();
options.bitrate = self.bitrate.map(moq_net::bandwidth::Rate::from_bps);
options.codec = self.codec.into();
options.kind = if self.software {
moq_video::encode::Kind::Software
} else if self.hardware {
moq_video::encode::Kind::Hardware
} else {
moq_video::encode::Kind::Auto
};
options.bandwidth = bandwidth;
options
}
fn audio_source(&self) -> moq_audio::capture::Source {
use moq_audio::capture::Source;
if self.system_audio {
Source::System
} else {
Source::Microphone(self.microphone.clone().flatten())
}
}
fn audio_config(&self) -> moq_audio::capture::Config {
let mut config = moq_audio::capture::Config::default();
config.source = self.audio_source();
config
}
fn audio_encode(&self, bandwidth: moq_net::bandwidth::Allocator) -> moq_audio::encode::Options {
let mut options = moq_audio::encode::Options::default();
options.settings.bitrate = self
.audio_bitrate
.map(|bps| moq_net::bandwidth::Rate::from_bps(bps.into()));
options.bandwidth = bandwidth;
options
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use bytes::BytesMut;
use moq_mux::catalog::CatalogFormat;
use moq_mux::catalog::hang::{Catalog, Container};
use moq_mux::container::ts::{self as tscat, Export, Import};
use moq_mux::container::{Consumer, Frame, Producer};
use moq_net::Timestamp;
use super::*;
const BBB: &[u8] = include_bytes!("../../moq-mux/src/container/ts/test_data/bbb.ts");
const CUE: &[u8] = &[
0xfc, 0x30, 0x1b, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0xff, 0xf0, 0x0a, 0x05, 0x00, 0x00, 0x2b, 0xb4,
0x7f, 0xdf, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0xad, 0x25, 0xe8, 0x39,
];
const PES_PAYLOAD: &[u8] = &[0xDE, 0xAD, 0xBE, 0xEF, 0x01, 0x02];
const SECTION_PID: u16 = 0x102;
const VERBATIM_PES_PID: u16 = 0x104;
const VERBATIM_PES_STREAM_ID: u8 = 0xC0;
async fn drain(mut exporter: Export<tscat::Ext>) -> Vec<u8> {
let mut out = Vec::new();
while let Ok(res) = tokio::time::timeout(Duration::from_millis(500), exporter.next()).await {
match res.expect("exporter error") {
Some(frame) => out.extend_from_slice(&frame.payload),
None => break,
}
}
out
}
async fn settle() {
for _ in 0..10 {
tokio::task::yield_now().await;
}
}
async fn manufacture_input() -> Vec<u8> {
let origin = moq_tokio::origin::spawn();
let mut broadcast = origin.create_broadcast("cli").unwrap();
broadcast.announce(Default::default()).unwrap();
settle().await;
let config = moq_mux::catalog::Config::default().with_catalog(Catalog::<tscat::Ext>::default());
let mut catalog = moq_mux::catalog::Producer::new(&mut broadcast, config).unwrap();
let section = broadcast
.unique_track(".scte35", hang::container::track_info(hang::catalog::PRIORITY.text))
.unwrap();
let mut section_track = tscat::Track::new(SECTION_PID);
section_track.verbatim = Some(tscat::Verbatim::new(0x86, tscat::Framing::Section));
catalog
.modify()
.unwrap()
.ext
.mpegts
.tracks
.insert(section.name().to_string(), section_track);
let mut section_producer = Producer::new(section, Container::Legacy(moq_mux::container::Kind::Data));
section_producer
.write(Frame {
timestamp: Timestamp::from_millis(1410).unwrap(),
duration: None,
payload: bytes::Bytes::from_static(CUE),
keyframe: true,
})
.unwrap();
section_producer.cut(None).unwrap();
section_producer.finish().unwrap();
let pes = broadcast
.unique_track(".data", hang::container::track_info(hang::catalog::PRIORITY.text))
.unwrap();
let mut verbatim = tscat::Verbatim::new(0x06, tscat::Framing::Pes);
verbatim.stream_id = Some(VERBATIM_PES_STREAM_ID);
let mut pes_track = tscat::Track::new(VERBATIM_PES_PID);
pes_track.verbatim = Some(verbatim);
catalog
.modify()
.unwrap()
.ext
.mpegts
.tracks
.insert(pes.name().to_string(), pes_track);
let mut pes_producer = Producer::new(pes, Container::Legacy(moq_mux::container::Kind::Data));
pes_producer
.write(Frame {
timestamp: Timestamp::from_millis(1410).unwrap(),
duration: None,
payload: bytes::Bytes::from_static(PES_PAYLOAD),
keyframe: true,
})
.unwrap();
pes_producer.cut(None).unwrap();
pes_producer.finish().unwrap();
let mut import = Import::new(broadcast, catalog.reserve());
import.decode(&BytesMut::from(BBB)).unwrap();
import.finish().unwrap();
drain(
Export::with_ts(moq_mux::Source::new(origin.consume(), "cli"), CatalogFormat::Hang)
.await
.unwrap()
.with_max_age(RECORDING_MAX_AGE),
)
.await
}
const RECORDING_MAX_AGE: std::time::Duration = Duration::from_secs(30);
#[tokio::test(start_paused = true)]
async fn ts_verbatim_streams_round_trip_through_cli() {
ts_verbatim_round_trip(CatalogFormat::Hang).await;
}
#[tokio::test(start_paused = true)]
async fn ts_verbatim_streams_round_trip_through_msf() {
ts_verbatim_round_trip(CatalogFormat::Msf).await;
}
async fn ts_verbatim_round_trip(format: CatalogFormat) {
let input = manufacture_input().await;
let origin = moq_tokio::origin::spawn();
let broadcast = origin.create_broadcast("cli").unwrap();
broadcast.announce(Default::default()).unwrap();
settle().await;
let mut publish = Publish::new(broadcast, &PublishFormat::Ts, Default::default()).unwrap();
#[allow(irrefutable_let_patterns)]
let Source::Stream { decoder, .. } = &mut publish.source else {
panic!("expected a stream source");
};
decoder.decode_chunk(&input).unwrap();
decoder.finish().unwrap();
let output = drain(
Export::with_ts(moq_mux::Source::new(origin.consume(), "cli"), format)
.await
.unwrap()
.with_max_age(RECORDING_MAX_AGE),
)
.await;
let mut broadcast = moq_net::broadcast::Info::new().produce();
let consumer = broadcast.consume();
let config = moq_mux::catalog::Config::default().with_catalog(Catalog::<tscat::Ext>::default());
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config).unwrap();
let mut import = Import::new(broadcast, catalog.reserve());
import.decode(&BytesMut::from(&output[..])).unwrap();
import.finish().unwrap();
let snapshot = catalog.snapshot();
let (section_name, section) = snapshot
.ext
.mpegts
.tracks
.iter()
.find(|(_, t)| t.verbatim.as_ref().is_some_and(|v| v.stream_type == 0x86))
.expect("SCTE-35 section survived the round-trip");
assert_eq!(section.pid, SECTION_PID, "section PID preserved");
assert_eq!(
section.verbatim.as_ref().unwrap().framing,
tscat::Framing::Section,
"section framing preserved"
);
let section_name = section_name.clone();
let (pes_name, pes) = snapshot
.ext
.mpegts
.tracks
.iter()
.find(|(_, t)| t.verbatim.as_ref().is_some_and(|v| v.stream_type == 0x06))
.expect("verbatim PES survived the round-trip");
assert_eq!(pes.pid, VERBATIM_PES_PID, "verbatim PES PID preserved");
let pes_verbatim = pes.verbatim.as_ref().unwrap();
assert_eq!(pes_verbatim.framing, tscat::Framing::Pes, "PES framing preserved");
assert_eq!(
pes_verbatim.stream_id,
Some(VERBATIM_PES_STREAM_ID),
"PES stream_id preserved"
);
let pes_name = pes_name.clone();
assert_eq!(
read_frame(&consumer, §ion_name).await,
CUE,
"SCTE-35 section round-trips byte-for-byte"
);
assert_eq!(
read_frame(&consumer, &pes_name).await,
PES_PAYLOAD,
"verbatim PES payload round-trips byte-for-byte"
);
}
#[tokio::test]
async fn ts_import_publishes_on_the_broadcast_clock() {
let ago = Duration::from_secs(60);
let clock = moq_mux::Clock::at(std::time::Instant::now() - ago, std::time::SystemTime::now() - ago).unwrap();
let broadcast = moq_net::broadcast::Info::new().produce();
let consumer = broadcast.consume();
let config = moq_mux::catalog::Config::default().with_clock(clock);
let mut publish = Publish::new(broadcast, &PublishFormat::Ts, config).unwrap();
#[allow(irrefutable_let_patterns)]
let Source::Stream { decoder, .. } = &mut publish.source else {
panic!("expected a stream source");
};
let before = clock.now();
decoder.decode_chunk(BBB).unwrap();
let after = clock.now();
decoder.finish().unwrap();
let catalog = hang::catalog::Catalog::<()>::subscribe(&consumer)
.await
.unwrap()
.next()
.await
.unwrap()
.expect("a catalog");
assert_eq!(
catalog.clock,
Some(clock.wall()),
"the advertised clock is the one stamped on"
);
let (name, config) = catalog.video.renditions.iter().next().expect("a video rendition");
let track = consumer.track(name).unwrap().subscribe(None).await.unwrap();
let container = moq_mux::catalog::hang::Container::try_from(config).unwrap();
let first = Consumer::new(track, container)
.read()
.await
.unwrap()
.expect("a video frame")
.timestamp;
let skew = Duration::from_secs(2).as_micros();
assert!(
before.as_micros() - skew <= first.as_micros() && first.as_micros() <= after.as_micros() + skew,
"the first frame is live on arrival: {first:?} not in {before:?}..={after:?}"
);
}
#[tokio::test(start_paused = true)]
async fn eof_finishes_the_catalog_with_its_renditions() {
let broadcast = moq_net::broadcast::Info::new().produce();
let consumer = broadcast.consume();
let publish = Publish::new(broadcast, &PublishFormat::Ts, Default::default()).unwrap();
let mut catalogs = hang::catalog::Catalog::<()>::subscribe(&consumer).await.unwrap();
#[allow(irrefutable_let_patterns)]
let Source::Stream { decoder, catalog } = publish.source else {
panic!("expected a stream source");
};
decode(decoder, catalog, BBB).await.unwrap();
drop(publish.broadcast);
let mut last = None;
loop {
let next = tokio::time::timeout(Duration::from_secs(1), catalogs.next())
.await
.expect("the catalog track ends");
match next.expect("the catalog ends cleanly") {
Some(catalog) => last = Some(catalog),
None => break,
}
}
let last = last.expect("a catalog");
assert_eq!(last.video.renditions.len(), 1, "the video rendition is still listed");
assert_eq!(last.audio.renditions.len(), 1, "the audio rendition is still listed");
}
async fn read_frame(consumer: &moq_net::broadcast::Consumer, name: &str) -> Vec<u8> {
let track = consumer.track(name).unwrap().subscribe(None).await.unwrap();
let mut reader = Consumer::new(track, Container::Legacy(moq_mux::container::Kind::Data));
let frame = tokio::time::timeout(Duration::from_secs(1), reader.read())
.await
.expect("verbatim read timed out")
.unwrap()
.expect("a published verbatim frame");
frame.payload.to_vec()
}
}