pub mod av1;
pub mod h264;
pub mod h265;
pub mod opus;
pub mod vp8;
pub mod vp9;
#[cfg(test)]
mod bitstream_test;
use bytes::Bytes;
use hang::catalog::VideoConfig;
use crate::Result;
#[derive(Clone, Debug)]
pub struct Frame {
pub timestamp_us: u64,
pub payload: Bytes,
}
pub trait Bridge: Send {
fn push(&mut self, frame: Frame) -> Result<()>;
fn abort(self: Box<Self>, err: moq_net::Error);
}
pub(crate) trait DeferredImport: Send + Sized {
fn create(track: moq_net::track::Producer, reserved: moq_mux::catalog::Reserved) -> moq_mux::Result<Self>;
fn decode(&mut self, frame: Bytes, pts: moq_net::Timestamp) -> moq_mux::Result<()>;
fn abort(self, err: moq_net::Error);
}
impl DeferredImport for moq_mux::codec::vp8::Import {
fn create(track: moq_net::track::Producer, reserved: moq_mux::catalog::Reserved) -> moq_mux::Result<Self> {
Self::new(track, reserved, Default::default())
}
fn decode(&mut self, frame: Bytes, pts: moq_net::Timestamp) -> moq_mux::Result<()> {
moq_mux::codec::vp8::Import::decode(self, frame, Some(pts))
}
fn abort(self, err: moq_net::Error) {
moq_mux::codec::vp8::Import::abort(self, err);
}
}
impl DeferredImport for moq_mux::codec::vp9::Import {
fn create(track: moq_net::track::Producer, reserved: moq_mux::catalog::Reserved) -> moq_mux::Result<Self> {
Self::new(track, reserved, Default::default())
}
fn decode(&mut self, frame: Bytes, pts: moq_net::Timestamp) -> moq_mux::Result<()> {
moq_mux::codec::vp9::Import::decode(self, frame, Some(pts))
}
fn abort(self, err: moq_net::Error) {
moq_mux::codec::vp9::Import::abort(self, err);
}
}
struct PendingVideo {
track: moq_net::track::Producer,
catalog: moq_mux::catalog::Producer,
}
enum DeferredState<I> {
Pending(Box<PendingVideo>),
Active(Box<I>),
Failed(Box<moq_net::track::Producer>),
Poisoned,
}
pub(crate) struct DeferredVideo<I> {
state: DeferredState<I>,
}
impl<I: DeferredImport> DeferredVideo<I> {
pub fn new(
mut broadcast: moq_net::broadcast::Producer,
catalog: moq_mux::catalog::Producer,
suffix: &str,
) -> Result<Self> {
let track = broadcast.unique_track(suffix, catalog.track_info())?;
Ok(Self {
state: DeferredState::Pending(Box::new(PendingVideo { track, catalog })),
})
}
pub fn decode(&mut self, frame: Bytes, pts: moq_net::Timestamp) -> Result<()> {
if let DeferredState::Active(import) = &mut self.state {
return import.decode(frame, pts).map_err(Into::into);
}
let DeferredState::Pending(pending) = std::mem::replace(&mut self.state, DeferredState::Poisoned) else {
return Err(crate::Error::Other(anyhow::anyhow!(
"video bridge initialization already failed"
)));
};
let reserved = pending.catalog.reserve();
let abort = pending.track.clone();
let import = match I::create(pending.track, reserved) {
Ok(import) => import,
Err(err) => {
self.state = DeferredState::Failed(Box::new(abort));
return Err(err.into());
}
};
self.state = DeferredState::Active(Box::new(import));
let DeferredState::Active(import) = &mut self.state else {
unreachable!();
};
import.decode(frame, pts).map_err(Into::into)
}
pub fn abort(self, err: moq_net::Error) {
match self.state {
DeferredState::Pending(pending) => {
let _ = pending.track.abort(err);
}
DeferredState::Active(import) => import.abort(err),
DeferredState::Failed(track) => {
let _ = track.abort(err);
}
DeferredState::Poisoned => {}
}
}
}
#[derive(Clone, Debug)]
pub struct PacketizedFrame {
pub timestamp_us: u64,
pub payload: Bytes,
}
pub struct Track {
consumer: moq_mux::container::Consumer<moq_mux::catalog::hang::Container>,
convert: TrackConvert,
}
enum TrackConvert {
Passthrough,
LengthPrefixed { length_size: usize, keyframe_prefix: Bytes },
}
impl Track {
pub fn opus(track: moq_net::track::Subscriber) -> Self {
let container = moq_mux::catalog::hang::Container::Legacy;
let consumer = moq_mux::container::Consumer::new(track, container);
Self {
consumer,
convert: TrackConvert::Passthrough,
}
}
pub fn video(track: moq_net::track::Subscriber, config: &VideoConfig) -> Result<Self> {
let container: moq_mux::catalog::hang::Container = (&config.container).try_into()?;
let consumer = moq_mux::container::Consumer::new(track, container);
let convert = match &config.codec {
hang::catalog::VideoCodec::VP8 => TrackConvert::Passthrough,
hang::catalog::VideoCodec::VP9(_) => TrackConvert::Passthrough,
hang::catalog::VideoCodec::AV1(_) => TrackConvert::Passthrough,
hang::catalog::VideoCodec::H264(_) => h264_convert(config)?,
hang::catalog::VideoCodec::H265(_) => h265_convert(config)?,
other => return Err(crate::Error::UnsupportedCodec(format!("{other:?}"))),
};
Ok(Self { consumer, convert })
}
pub async fn next(&mut self) -> Result<Option<PacketizedFrame>> {
loop {
let Some(frame) = self.consumer.read().await? else {
return Ok(None);
};
let payload = match &self.convert {
TrackConvert::Passthrough => frame.payload,
TrackConvert::LengthPrefixed {
length_size,
keyframe_prefix,
} => {
let prefix = frame.keyframe.then(|| keyframe_prefix.as_ref());
moq_mux::codec::annexb::from_length_prefixed(&frame.payload, *length_size, prefix)
.map_err(|err| crate::Error::Other(anyhow::anyhow!("annexb: {err}")))?
}
};
if payload.is_empty() {
continue;
}
return Ok(Some(PacketizedFrame {
timestamp_us: frame.timestamp.as_micros() as u64,
payload,
}));
}
}
}
fn h264_convert(config: &VideoConfig) -> Result<TrackConvert> {
let Some(avcc) = config.description.as_ref().filter(|d| !d.is_empty()) else {
return Ok(TrackConvert::Passthrough);
};
let params = moq_mux::codec::h264::Avcc::parse(avcc)
.map_err(|err| crate::Error::Other(anyhow::anyhow!("avcc parse: {err}")))?;
if params.sps.is_empty() || params.pps.is_empty() {
return Err(crate::Error::Other(anyhow::anyhow!(
"avc1 avcC is missing parameter sets (sps={}, pps={})",
params.sps.len(),
params.pps.len()
)));
}
let keyframe_prefix = moq_mux::codec::annexb::build_prefix(params.sps.iter().chain(params.pps.iter()));
Ok(TrackConvert::LengthPrefixed {
length_size: params.length_size,
keyframe_prefix,
})
}
fn h265_convert(config: &VideoConfig) -> Result<TrackConvert> {
let Some(hvcc) = config.description.as_ref().filter(|d| !d.is_empty()) else {
return Ok(TrackConvert::Passthrough);
};
let params = moq_mux::codec::h265::Hvcc::parse(hvcc)
.map_err(|err| crate::Error::Other(anyhow::anyhow!("hvcc parse: {err}")))?;
if params.vps.is_empty() || params.sps.is_empty() || params.pps.is_empty() {
return Err(crate::Error::Other(anyhow::anyhow!(
"hvc1 hvcC is missing parameter sets (vps={}, sps={}, pps={})",
params.vps.len(),
params.sps.len(),
params.pps.len()
)));
}
let keyframe_prefix =
moq_mux::codec::annexb::build_prefix(params.vps.iter().chain(params.sps.iter()).chain(params.pps.iter()));
Ok(TrackConvert::LengthPrefixed {
length_size: params.length_size,
keyframe_prefix,
})
}
#[cfg(test)]
mod tests {
use hang::catalog::{H264, H265, VideoConfig};
use super::*;
fn config(codec: impl Into<hang::catalog::VideoCodec>, description: Option<Bytes>) -> VideoConfig {
let mut config = VideoConfig::new(codec);
config.description = description;
config
}
fn h264(inline: bool) -> H264 {
H264 {
inline,
profile: 0x42,
constraints: 0,
level: 0x1f,
}
}
fn h265(in_band: bool) -> H265 {
H265 {
in_band,
profile_space: 0,
profile_idc: 1,
profile_compatibility_flags: [0; 4],
tier_flag: false,
level_idc: 0x5d,
constraint_flags: [0; 6],
}
}
fn build_avcc(sps: &[u8], pps: &[u8]) -> Bytes {
let mut v = vec![1, sps[1], sps[2], sps[3], 0xff, 0xe1];
v.extend_from_slice(&(sps.len() as u16).to_be_bytes());
v.extend_from_slice(sps);
v.push(1);
v.extend_from_slice(&(pps.len() as u16).to_be_bytes());
v.extend_from_slice(pps);
Bytes::from(v)
}
fn build_hvcc(vps: &[u8], sps: &[u8], pps: &[u8]) -> Bytes {
let mut v = vec![0u8; 21];
v.push(0xff); v.push(3); for (nal_type, nal) in [(32u8, vps), (33, sps), (34, pps)] {
v.push(nal_type); v.extend_from_slice(&1u16.to_be_bytes()); v.extend_from_slice(&(nal.len() as u16).to_be_bytes());
v.extend_from_slice(nal);
}
Bytes::from(v)
}
#[test]
fn h264_avc3_passthrough() {
let cfg = config(h264(true), None);
assert!(matches!(h264_convert(&cfg).unwrap(), TrackConvert::Passthrough));
}
#[test]
fn h264_avc1_length_prefixed() {
let sps: &[u8] = &[0x67, 0x42, 0xc0, 0x1f, 0xde];
let pps: &[u8] = &[0x68, 0xce, 0x3c, 0x80];
let cfg = config(h264(false), Some(build_avcc(sps, pps)));
let TrackConvert::LengthPrefixed {
length_size,
keyframe_prefix,
} = h264_convert(&cfg).unwrap()
else {
panic!("expected LengthPrefixed");
};
assert_eq!(length_size, 4);
assert!(keyframe_prefix.starts_with(&[0, 0, 0, 1]), "Annex-B start code");
assert!(keyframe_prefix.windows(sps.len()).any(|w| w == sps), "SPS in prefix");
assert!(keyframe_prefix.windows(pps.len()).any(|w| w == pps), "PPS in prefix");
}
#[test]
fn h265_hev1_passthrough() {
let cfg = config(h265(true), None);
assert!(matches!(h265_convert(&cfg).unwrap(), TrackConvert::Passthrough));
}
#[test]
fn h265_hvc1_length_prefixed() {
let vps: &[u8] = &[0x40, 0x01, 0x0c, 0x01];
let sps: &[u8] = &[0x42, 0x01, 0x01, 0x01];
let pps: &[u8] = &[0x44, 0x01, 0xc0, 0xf7];
let cfg = config(h265(false), Some(build_hvcc(vps, sps, pps)));
let TrackConvert::LengthPrefixed {
length_size,
keyframe_prefix,
} = h265_convert(&cfg).unwrap()
else {
panic!("expected LengthPrefixed");
};
assert_eq!(length_size, 4);
let v = keyframe_prefix.windows(vps.len()).position(|w| w == vps).expect("VPS");
let s = keyframe_prefix.windows(sps.len()).position(|w| w == sps).expect("SPS");
let p = keyframe_prefix.windows(pps.len()).position(|w| w == pps).expect("PPS");
assert!(v < s && s < p, "VPS < SPS < PPS order in prefix");
}
#[test]
fn h264_avc1_missing_param_sets_errors() {
let avcc = Bytes::from(vec![1, 0x42, 0, 0x1f, 0xff, 0xe0, 0x00]);
let cfg = config(h264(false), Some(avcc));
assert!(h264_convert(&cfg).is_err(), "missing SPS/PPS must error");
}
}