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>,
}
impl ExportSource {
pub fn for_video(
source: &crate::Source,
name: &str,
config: &VideoConfig,
latency: Duration,
) -> Result<Self, crate::Error> {
let media: HangContainer = (&config.container).try_into()?;
let transform = build_video_transform(config);
let description = config.description.as_ref().filter(|b| !b.is_empty()).cloned();
Ok(Self {
state: SourceState::Requesting(source.request(config.broadcast.as_ref()), name.to_string()),
media: Some(media),
latency,
transform,
description,
})
}
pub fn for_video_raw(
source: &crate::Source,
name: &str,
config: &VideoConfig,
latency: Duration,
) -> Result<Self, crate::Error> {
let media: HangContainer = (&config.container).try_into()?;
let description = config.description.as_ref().filter(|b| !b.is_empty()).cloned();
Ok(Self {
state: SourceState::Requesting(source.request(config.broadcast.as_ref()), name.to_string()),
media: Some(media),
latency,
transform: None,
description,
})
}
pub fn for_audio(
source: &crate::Source,
name: &str,
config: &AudioConfig,
latency: Duration,
) -> Result<Self, crate::Error> {
let media: HangContainer = (&config.container).try_into()?;
let description = config.description.as_ref().filter(|b| !b.is_empty()).cloned();
Ok(Self {
state: SourceState::Requesting(source.request(config.broadcast.as_ref()), name.to_string()),
media: Some(media),
latency,
transform: None,
description,
})
}
pub fn for_stream(source: &crate::Source, name: &str, latency: Duration) -> Result<Self, crate::Error> {
Ok(Self {
state: SourceState::Requesting(source.request(None), name.to_string()),
media: Some(HangContainer::Legacy),
latency,
transform: None,
description: 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 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 {
return Poll::Ready(Ok(Some(frame)));
};
match transform.transform(frame.payload.clone())? {
None => {
self.refresh_description();
continue;
}
Some(payload) => {
self.refresh_description();
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());
}
}
}
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,
}
}