use std::task::Poll;
use std::time::Duration;
use bytes::Bytes;
use hang::catalog::{AudioConfig, VideoCodec, VideoConfig};
use crate::catalog::hang::Container as HangContainer;
use crate::codec::h264::Avc1;
use crate::codec::h265::Hvc1;
use crate::container::{Consumer, Frame};
pub(crate) enum VideoTransform {
Avc1(Avc1),
Hvc1(Hvc1),
}
impl VideoTransform {
pub(crate) fn codec_private(&self) -> Option<&Bytes> {
match self {
VideoTransform::Avc1(t) => t.avcc(),
VideoTransform::Hvc1(t) => t.hvcc(),
}
}
pub(crate) fn transform(&mut self, payload: Bytes) -> crate::Result<Option<Bytes>> {
match self {
VideoTransform::Avc1(t) => Ok(t.transform(payload)?),
VideoTransform::Hvc1(t) => Ok(t.transform(payload)?),
}
}
}
enum SourceState {
Requesting(kio::Pending<moq_net::origin::Requesting>, String),
Subscribing(kio::Pending<moq_net::track::Subscribing>),
Active(Box<Consumer<HangContainer>>),
}
pub(crate) struct ExportSource {
state: SourceState,
media: Option<HangContainer>,
latency: Duration,
transform: Option<VideoTransform>,
description: Option<Bytes>,
video_codec: Option<VideoCodec>,
video_dimensions: Option<(u32, u32)>,
}
impl ExportSource {
pub fn for_video(
source: &crate::Source,
name: &str,
config: &VideoConfig,
latency: Duration,
) -> Result<Option<Self>, crate::Error> {
Self::video(source, name, config, latency, build_video_transform(config))
}
pub fn for_video_raw(
source: &crate::Source,
name: &str,
config: &VideoConfig,
latency: Duration,
) -> Result<Option<Self>, crate::Error> {
Self::video(source, name, config, latency, None)
}
fn video(
source: &crate::Source,
name: &str,
config: &VideoConfig,
latency: Duration,
transform: Option<VideoTransform>,
) -> Result<Option<Self>, crate::Error> {
let media: HangContainer = (&config.container).try_into()?;
let description = config.description.as_ref().filter(|b| !b.is_empty()).cloned();
let Some(request) = source.request(config.broadcast.as_ref()) else {
return Ok(None);
};
let mut source = Self {
state: SourceState::Requesting(request, name.to_string()),
media: Some(media),
latency,
transform,
description,
video_codec: Some(config.codec.clone()),
video_dimensions: catalog_dimensions(config),
};
source.resolve_video_dimensions(&[])?;
Ok(Some(source))
}
pub fn for_audio(
source: &crate::Source,
name: &str,
config: &AudioConfig,
latency: Duration,
) -> Result<Option<Self>, crate::Error> {
let media: HangContainer = (&config.container).try_into()?;
let description = config.description.as_ref().filter(|b| !b.is_empty()).cloned();
let Some(request) = source.request(config.broadcast.as_ref()) else {
return Ok(None);
};
Ok(Some(Self {
state: SourceState::Requesting(request, name.to_string()),
media: Some(media),
latency,
transform: None,
description,
video_codec: None,
video_dimensions: None,
}))
}
pub fn for_stream(source: &crate::Source, name: &str, latency: Duration) -> Result<Self, crate::Error> {
let request = source.request(None).expect("the catalog broadcast is always valid");
Ok(Self {
state: SourceState::Requesting(request, name.to_string()),
media: Some(HangContainer::Legacy),
latency,
transform: None,
description: None,
video_codec: None,
video_dimensions: None,
})
}
pub fn description(&self) -> Option<&Bytes> {
self.description.as_ref()
}
pub fn header_ready(&self) -> bool {
self.transform.is_none() || self.description.is_some()
}
pub fn video_config(&self, config: &VideoConfig) -> Option<VideoConfig> {
if catalog_dimensions(config).is_some() {
return Some(config.clone());
}
let (width, height) = self.video_dimensions?;
let mut config = config.clone();
config.coded_width = Some(width);
config.coded_height = Some(height);
Some(config)
}
pub fn video_geometry_ready(&self, config: &VideoConfig) -> bool {
!matches!(
config.codec,
VideoCodec::H264(_) | VideoCodec::H265(_) | VideoCodec::VP8 | VideoCodec::VP9(_) | VideoCodec::AV1(_)
) || self.video_config(config).is_some()
}
pub fn poll_read(&mut self, waiter: &kio::Waiter) -> Poll<crate::Result<Option<Frame>>> {
if matches!(self.state, SourceState::Requesting(..)) {
let (broadcast, name) = {
let SourceState::Requesting(pending, name) = &self.state else {
unreachable!("just matched Requesting");
};
match pending.poll_ok(waiter) {
Poll::Ready(Ok(broadcast)) => (broadcast, name.clone()),
Poll::Ready(Err(e)) => return Poll::Ready(Err(e.into())),
Poll::Pending => return Poll::Pending,
}
};
self.state = SourceState::Subscribing(broadcast.track(&name)?.subscribe(None));
}
if matches!(self.state, SourceState::Subscribing(_)) {
let track = {
let SourceState::Subscribing(pending) = &self.state else {
unreachable!("just matched Subscribing");
};
match pending.poll_ok(waiter) {
Poll::Ready(Ok(track)) => track,
Poll::Ready(Err(e)) => return Poll::Ready(Err(e.into())),
Poll::Pending => return Poll::Pending,
}
};
let media = self
.media
.take()
.expect("media present until the subscription resolves");
self.state = SourceState::Active(Box::new(Consumer::new(track, media).with_latency(self.latency)));
}
loop {
let frame = {
let SourceState::Active(consumer) = &mut self.state else {
unreachable!("subscription resolved into an Active consumer");
};
match consumer.poll_read(waiter) {
Poll::Ready(Ok(Some(f))) => f,
Poll::Ready(Ok(None)) => return Poll::Ready(Ok(None)),
Poll::Ready(Err(e)) => return Poll::Ready(Err(e)),
Poll::Pending => return Poll::Pending,
}
};
let Some(transform) = self.transform.as_mut() else {
self.resolve_video_dimensions(&frame.payload)?;
return Poll::Ready(Ok(Some(frame)));
};
match transform.transform(frame.payload.clone())? {
None => {
self.refresh_description();
self.resolve_video_dimensions(&frame.payload)?;
continue;
}
Some(payload) => {
self.refresh_description();
self.resolve_video_dimensions(&payload)?;
return Poll::Ready(Ok(Some(Frame { payload, ..frame })));
}
}
}
}
fn refresh_description(&mut self) {
if let Some(transform) = self.transform.as_ref()
&& let Some(d) = transform.codec_private()
&& self.description.as_ref() != Some(d)
{
self.description = Some(d.clone());
}
}
fn resolve_video_dimensions(&mut self, payload: &[u8]) -> crate::Result<()> {
if self.video_dimensions.is_some() {
return Ok(());
}
let Some(codec) = self.video_codec.as_ref() else {
return Ok(());
};
self.video_dimensions = codec_dimensions(codec, self.description.as_deref(), payload)?;
Ok(())
}
}
pub(crate) fn catalog_dimensions(config: &VideoConfig) -> Option<(u32, u32)> {
let dimensions = (config.coded_width?, config.coded_height?);
(dimensions.0 > 0 && dimensions.1 > 0).then_some(dimensions)
}
pub(crate) fn codec_dimensions(
codec: &VideoCodec,
description: Option<&[u8]>,
payload: &[u8],
) -> crate::Result<Option<(u32, u32)>> {
let dimensions = match codec {
VideoCodec::H264(_) => match description {
Some(description) => catalog_dimensions(&crate::codec::h264::config(description)?),
None => None,
},
VideoCodec::H265(_) => match description {
Some(description) => catalog_dimensions(&crate::codec::h265::config(description)?),
None => None,
},
VideoCodec::VP8 if !payload.is_empty() => crate::codec::vp8::FrameHeader::parse(payload)?
.dimensions
.map(|(width, height)| (u32::from(width), u32::from(height))),
VideoCodec::VP9(_) if !payload.is_empty() => crate::codec::vp9::config_from_keyframe(payload)?
.as_ref()
.and_then(catalog_dimensions),
VideoCodec::AV1(_) if !payload.is_empty() => crate::codec::av1::dimensions(payload)?,
_ => None,
};
Ok(dimensions.filter(|(width, height)| *width > 0 && *height > 0))
}
pub(crate) fn build_video_transform(config: &VideoConfig) -> Option<VideoTransform> {
let needs_transform = config.description.as_ref().map(|d| d.is_empty()).unwrap_or(true);
if !needs_transform {
return None;
}
match &config.codec {
VideoCodec::H264(_) => Some(VideoTransform::Avc1(Avc1::new())),
VideoCodec::H265(_) => Some(VideoTransform::Hvc1(Hvc1::new())),
_ => None,
}
}