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(clap::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(clap::Args, Clone)]
#[command(group = clap::ArgGroup::new("video-source").multiple(false))]
#[command(group = clap::ArgGroup::new("audio-source").multiple(false))]
pub struct CaptureArgs {
#[arg(long, num_args = 0..=1, group = "video-source")]
pub camera: Option<Option<String>>,
#[arg(long, num_args = 0..=1, group = "video-source", alias = "screen")]
pub display: Option<Option<String>>,
#[arg(long, group = "video-source")]
pub window: Option<String>,
#[arg(long, group = "video-source")]
pub app: Option<String>,
#[arg(long)]
pub no_cursor: bool,
#[arg(long)]
pub width: Option<u32>,
#[arg(long)]
pub height: Option<u32>,
#[arg(long)]
pub fps: Option<u32>,
#[arg(long)]
pub bitrate: Option<u64>,
#[arg(long, value_enum, default_value_t)]
pub codec: VideoCodec,
#[arg(long, conflicts_with = "software")]
pub hardware: bool,
#[arg(long)]
pub software: bool,
#[arg(long, num_args = 0..=1, group = "audio-source")]
pub microphone: Option<Option<String>>,
#[arg(long, group = "audio-source")]
pub system_audio: bool,
#[arg(long)]
pub audio_bitrate: Option<u32>,
#[arg(long, conflicts_with = "no_audio", conflicts_with = "video-source")]
pub no_video: bool,
#[arg(long, conflicts_with = "audio-source")]
pub no_audio: bool,
}
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 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),
}
}
}
#[allow(clippy::large_enum_variant)]
enum Source {
Stream(PublishDecoder),
#[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) -> anyhow::Result<Self> {
if let PublishFormat::Ts = format {
let catalog = moq_mux::catalog::Producer::with_catalog(
&mut broadcast,
moq_mux::catalog::hang::Catalog::<ts::Ext>::default(),
)?;
let ts = ts::Import::new(broadcast.clone(), catalog.reserve());
return Ok(Self {
source: Source::Stream(PublishDecoder::Ts(Box::new(ts))),
broadcast,
});
}
let catalog = moq_mux::catalog::Producer::new(&mut broadcast)?;
let source = match format {
PublishFormat::Avc3 => {
let track = moq_mux::import::unique_track(&mut broadcast, ".avc3")?;
let import = moq_mux::codec::h264::Import::new(track, catalog.reserve(), Default::default())?;
let split = Box::new(moq_mux::codec::h264::Split::new());
Source::Stream(PublishDecoder::Avc3 {
split,
import: Box::new(import),
})
}
PublishFormat::Fmp4 => {
let fmp4 = fmp4::Import::new(broadcast.clone(), catalog.reserve());
Source::Stream(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());
Source::Stream(PublishDecoder::Flv(Box::new(flv)))
}
};
Ok(Self { source, broadcast })
}
#[cfg(feature = "capture")]
pub fn capture(
mut broadcast: moq_net::broadcast::Producer,
args: &CaptureArgs,
bandwidth: Option<moq_net::bandwidth::Consumer>,
) -> anyhow::Result<Self> {
let catalog = moq_mux::catalog::Producer::new(&mut broadcast)?;
let video = (!args.no_video).then(|| (args.video_config(), args.video_encode(bandwidth)));
let audio = (!args.no_audio).then(|| (args.audio_config(), args.audio_encode()));
anyhow::ensure!(video.is_some() || audio.is_some(), "nothing to capture");
Ok(Self {
source: Source::Capture { catalog, video, audio },
broadcast,
})
}
pub async fn run(self) -> anyhow::Result<()> {
match self.source {
Source::Stream(mut decoder) => {
let mut stdin = tokio::io::stdin();
let mut buffer = bytes::BytesMut::new();
let result: anyhow::Result<()> = async {
loop {
buffer.clear();
let n = tokio::io::AsyncReadExt::read_buf(&mut stdin, &mut buffer).await?;
if n == 0 {
return Ok(()); }
decoder.decode_chunk(&buffer)?;
}
}
.await;
let outcome = result.and_then(|()| decoder.finish());
if let Err(err) = &outcome {
decoder.abort(moq_net::Error::Transport(err.to_string()));
}
outcome
}
#[cfg(feature = "capture")]
Source::Capture { catalog, video, audio } => {
let clock = moq_mux::Clock::new();
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)) => {
moq_audio::encode::publish_capture(broadcast, catalog, config, encode, clock)
.await
.map_err(anyhow::Error::from)
}
None => Ok(()),
}
}
};
tokio::try_join!(video_fut, audio_fut)?;
Ok(())
}
}
}
}
#[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) -> 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;
config.cursor = !self.no_cursor;
config
}
fn video_encode(&self, bandwidth: Option<moq_net::bandwidth::Consumer>) -> moq_video::encode::Options {
let mut options = moq_video::encode::Options::default();
options.bitrate = self.bitrate;
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) -> moq_audio::encode::Options {
let mut options = moq_audio::encode::Options::default();
options.bitrate = self.audio_bitrate;
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_net::Origin::random().produce();
let mut broadcast = origin
.create_broadcast("cli", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
settle().await;
let mut catalog =
moq_mux::catalog::Producer::with_catalog(&mut broadcast, Catalog::<tscat::Ext>::default()).unwrap();
let section = broadcast
.unique_track(".scte35", hang::container::track_info())
.unwrap();
let mut section_track = tscat::Track::new(SECTION_PID);
section_track.verbatim = Some(tscat::Verbatim::new(0x86, tscat::Framing::Section));
catalog
.lock()
.mpegts
.tracks
.insert(section.name().to_string(), section_track);
let mut section_producer = Producer::new(section, Container::Legacy);
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()).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.lock().mpegts.tracks.insert(pes.name().to_string(), pes_track);
let mut pes_producer = Producer::new(pes, Container::Legacy);
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_latency(Duration::ZERO),
)
.await
}
#[tokio::test(start_paused = true)]
async fn ts_verbatim_streams_round_trip_through_cli() {
let input = manufacture_input().await;
let origin = moq_net::Origin::random().produce();
let broadcast = origin
.create_broadcast("cli", moq_net::broadcast::Route::new().with_announce(true))
.unwrap();
settle().await;
let mut publish = Publish::new(broadcast, &PublishFormat::Ts).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"), CatalogFormat::Hang)
.await
.unwrap()
.with_latency(Duration::ZERO),
)
.await;
let mut broadcast = moq_net::broadcast::Info::new().produce();
let consumer = broadcast.consume();
let catalog =
moq_mux::catalog::Producer::with_catalog(&mut broadcast, Catalog::<tscat::Ext>::default()).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
.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
.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"
);
}
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).with_latency(Duration::ZERO);
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()
}
}