use std::ffi::{c_char, c_void};
use std::time::Duration;
use tokio::sync::oneshot;
use crate::ffi::OnStatus;
use crate::{Error, Id, NonZeroSlab, Shared, State, ffi};
#[repr(C)]
#[allow(non_camel_case_types)]
#[derive(Clone, Copy, Debug)]
pub enum moq_video_pixel_format {
MOQ_VIDEO_PIXEL_FORMAT_I420 = 0,
MOQ_VIDEO_PIXEL_FORMAT_RGBA = 1,
}
#[repr(C)]
#[allow(non_camel_case_types)]
#[derive(Clone, Copy, Debug)]
pub enum moq_video_codec {
MOQ_VIDEO_CODEC_H264 = 0,
MOQ_VIDEO_CODEC_H265 = 1,
}
#[repr(C)]
#[allow(non_camel_case_types)]
#[derive(Clone, Copy, Debug)]
pub enum moq_video_encoder_kind {
MOQ_VIDEO_ENCODER_KIND_AUTO = 0,
MOQ_VIDEO_ENCODER_KIND_HARDWARE = 1,
MOQ_VIDEO_ENCODER_KIND_SOFTWARE = 2,
MOQ_VIDEO_ENCODER_KIND_NAMED = 3,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_video_encoder_input {
pub format: u32,
pub width: u32,
pub height: u32,
pub framerate: u32,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_video_encoder_output {
pub codec: u32,
pub bitrate: u64,
pub gop: u32,
pub kind: u32,
pub encoder: *const c_char,
pub encoder_len: usize,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_video_encoder_frame {
pub timestamp_us: u64,
pub data: *const u8,
pub data_size: usize,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_video_decoder_output {
pub latency_max_ms: u64,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_video_frame {
pub timestamp_us: u64,
pub width: u32,
pub height: u32,
pub data: *const u8,
pub data_size: usize,
}
#[derive(Default)]
pub struct Video {
producers: NonZeroSlab<Shared<VideoEncoder>>,
consumer_tasks: NonZeroSlab<Option<VideoTaskEntry>>,
frames: NonZeroSlab<VideoFrame>,
}
fn block_on<T>(future: impl std::future::Future<Output = T>) -> T {
pollster::block_on(future)
}
pub(crate) struct VideoEncoder {
encoder: moq_video::encode::Sink,
producer: moq_video::encode::Producer<moq_mux::catalog::hang::Extra>,
format: moq_video_pixel_format,
size: moq_video::Size,
}
struct VideoFrame {
timestamp_us: u64,
width: u32,
height: u32,
data: bytes::Bytes,
}
fn finalize(
producer: moq_video::encode::Producer<moq_mux::catalog::hang::Extra>,
drained: Result<(), moq_video::Error>,
) -> Result<(), Error> {
match drained {
Ok(()) => Ok(producer.finish()?),
Err(err) => {
producer.abort(moq_net::Error::Transport(err.to_string()));
Err(err.into())
}
}
}
struct VideoTaskEntry {
close: Option<oneshot::Sender<()>>,
callback: OnStatus,
}
impl VideoEncoder {
fn publish_frame(&mut self, timestamp_us: u64, data: &[u8]) -> Result<(), Error> {
let size = self.size;
let surface = match self.format {
moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 => {
moq_video::Surface::I420(moq_video::I420::new(size.width, size.height, data.to_vec())?)
}
moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA => moq_video::Surface::rgba(data, size)?,
};
let frame = moq_video::Frame::new(surface, moq_net::Timestamp::from_micros(timestamp_us)?);
let encoded = block_on(self.encoder.encode(frame))?;
self.producer.publish(&encoded)?;
Ok(())
}
fn publish_cut(&mut self) {
self.encoder.keyframe();
}
fn publish_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
block_on(self.encoder.set_bitrate(bitrate))?;
Ok(())
}
fn publish_finish(self) -> Result<(), Error> {
let VideoEncoder {
encoder, mut producer, ..
} = self;
let drained = block_on(encoder.finish()).and_then(|encoded| producer.publish(&encoded));
finalize(producer, drained)
}
}
impl Video {
pub fn publish(
&mut self,
broadcast: &moq_net::broadcast::Producer,
catalog: moq_mux::catalog::Producer<moq_mux::catalog::hang::Extra>,
format: moq_video_pixel_format,
config: &moq_video::encode::Config,
rendition: hang::catalog::VideoConfig,
encoder: moq_video::encode::Sink,
) -> Result<Id, Error> {
let producer = moq_video::encode::Producer::new(broadcast.clone(), catalog, rendition)?;
self.producers.insert(Shared::new(VideoEncoder {
encoder,
producer,
format,
size: config.size(),
}))
}
pub(crate) fn producer(&self, id: Id) -> Result<Shared<VideoEncoder>, Error> {
self.producers.get(id).cloned().ok_or(Error::MediaNotFound)
}
pub(crate) fn remove(&mut self, id: Id) -> Result<Shared<VideoEncoder>, Error> {
self.producers.remove(id).ok_or(Error::MediaNotFound)
}
pub fn consume(
&mut self,
broadcast: &moq_net::broadcast::Consumer,
catalog: &hang::catalog::VideoConfig,
name: &str,
config: moq_video::decode::Config,
on_frame: OnStatus,
) -> Result<Id, Error> {
let broadcast = broadcast.clone();
let catalog = catalog.clone();
let name = name.to_string();
let channel = oneshot::channel();
let entry = VideoTaskEntry {
close: Some(channel.0),
callback: on_frame,
};
let id = self.consumer_tasks.insert(Some(entry))?;
tokio::spawn(async move {
let res = async move {
let consumer = moq_video::decode::Consumer::new(&broadcast, &catalog, name, config).await?;
Self::run(on_frame, consumer, channel.1).await
}
.await;
let entry = State::lock().video.consumer_tasks.remove(id).flatten();
if let Some(entry) = entry {
entry.callback.call(res);
}
});
Ok(id)
}
async fn run(
callback: OnStatus,
mut consumer: moq_video::decode::Consumer,
mut close: oneshot::Receiver<()>,
) -> Result<(), Error> {
loop {
let frame = tokio::select! {
biased;
_ = &mut close => return Ok(()),
frame = consumer.read() => match frame? {
Some(frame) => frame,
None => return Ok(()),
},
};
let size = frame.size();
let frame = VideoFrame {
timestamp_us: frame.timestamp.as_micros() as u64,
width: size.width,
height: size.height,
data: frame.surface.into_i420()?,
};
let frame_id = State::lock().video.frames.insert(frame)?;
callback.call(Ok(frame_id));
}
}
pub fn consume_close(&mut self, id: Id) -> Result<(), Error> {
self.consumer_tasks
.get_mut(id)
.and_then(|entry| entry.as_mut())
.ok_or(Error::TrackNotFound)?
.close
.take()
.ok_or(Error::TrackNotFound)?;
Ok(())
}
pub fn frame_info(&self, id: Id, dst: &mut moq_video_frame) -> Result<(), Error> {
let frame = self.frames.get(id).ok_or(Error::FrameNotFound)?;
*dst = moq_video_frame {
timestamp_us: frame.timestamp_us,
width: frame.width,
height: frame.height,
data: frame.data.as_ptr(),
data_size: frame.data.len(),
};
Ok(())
}
pub fn frame_free(&mut self, id: Id) -> Result<(), Error> {
self.frames.remove(id).ok_or(Error::FrameNotFound)?;
Ok(())
}
}
fn pixel_format_from_u32(value: u32) -> Result<moq_video_pixel_format, Error> {
Ok(match value {
v if v == moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32 => {
moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420
}
v if v == moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32 => {
moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA
}
_ => return Err(Error::InvalidCode),
})
}
fn codec_from_u32(value: u32) -> Result<moq_video::encode::Codec, Error> {
use moq_video::encode::Codec;
Ok(match value {
v if v == moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32 => Codec::H264,
v if v == moq_video_codec::MOQ_VIDEO_CODEC_H265 as u32 => Codec::H265,
_ => return Err(Error::InvalidCode),
})
}
unsafe fn encoder_kind(output: &moq_video_encoder_output) -> Result<moq_video::encode::Kind, Error> {
use moq_video::encode::Kind;
Ok(match output.kind {
v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_AUTO as u32 => Kind::Auto,
v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_HARDWARE as u32 => Kind::Hardware,
v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32 => Kind::Software,
v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_NAMED as u32 => {
Kind::Named(unsafe { ffi::parse_str(output.encoder, output.encoder_len)? }.to_string())
}
_ => return Err(Error::InvalidCode),
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_publish_video_raw(
broadcast: u32,
input: *const moq_video_encoder_input,
output: *const moq_video_encoder_output,
) -> i32 {
ffi::enter(move || {
let broadcast = ffi::parse_id(broadcast)?;
let raw_input = unsafe { input.as_ref() }.ok_or(Error::InvalidPointer)?;
let raw_output = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
let format = pixel_format_from_u32(raw_input.format)?;
let mut config = moq_video::encode::Config::new(raw_input.width, raw_input.height, raw_input.framerate);
config.codec = codec_from_u32(raw_output.codec)?;
config.kind = unsafe { encoder_kind(raw_output)? };
config.bitrate = (raw_output.bitrate != 0).then_some(raw_output.bitrate);
if raw_output.gop != 0 {
config.gop = raw_output.gop;
}
let rendition = block_on(config.probe())?;
let encoder = block_on(moq_video::encode::Sink::open(&config))?;
let mut state = State::lock();
let State { publish, video, .. } = &mut *state;
let (broadcast_producer, catalog) = publish.pair_mut(broadcast)?;
video.publish(broadcast_producer, catalog.clone(), format, &config, rendition, encoder)
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_publish_video_raw_frame(producer: u32, frame: *const moq_video_encoder_frame) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let frame = unsafe { frame.as_ref() }.ok_or(Error::InvalidPointer)?;
let data = unsafe { ffi::parse_slice(frame.data, frame.data_size)? };
let producer = State::lock().video.producer(producer)?;
producer
.lock()
.as_mut()
.ok_or(Error::MediaNotFound)?
.publish_frame(frame.timestamp_us, data)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_publish_video_raw_cut(producer: u32) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let producer = State::lock().video.producer(producer)?;
producer.lock().as_mut().ok_or(Error::MediaNotFound)?.publish_cut();
Ok(())
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_publish_video_raw_bitrate(producer: u32, bitrate: u64) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let producer = State::lock().video.producer(producer)?;
producer
.lock()
.as_mut()
.ok_or(Error::MediaNotFound)?
.publish_bitrate(bitrate)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_publish_video_raw_finish(producer: u32) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let producer = State::lock().video.remove(producer)?;
producer.take().ok_or(Error::MediaNotFound)?.publish_finish()
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_consume_video_raw(
catalog: u32,
index: u32,
output: *const moq_video_decoder_output,
on_frame: Option<extern "C" fn(user_data: *mut c_void, frame: i32)>,
user_data: *mut c_void,
) -> i32 {
ffi::enter(move || {
let catalog = ffi::parse_id(catalog)?;
let raw = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
let mut config = moq_video::decode::Config::new();
config.latency_max = if raw.latency_max_ms == 0 {
None
} else {
Some(Duration::from_millis(raw.latency_max_ms))
};
let on_frame = unsafe { OnStatus::new(user_data, on_frame) };
let mut state = State::lock();
let (broadcast, video_cfg, name) = state.consume.video_rendition(catalog, index as usize)?;
let State { video, .. } = &mut *state;
video.consume(&broadcast, &video_cfg, &name, config, on_frame)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_consume_video_raw_close(consumer: u32) -> i32 {
ffi::enter(move || {
let consumer = ffi::parse_id(consumer)?;
State::lock().video.consume_close(consumer)
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_consume_video_raw_frame(id: u32, dst: *mut moq_video_frame) -> i32 {
ffi::enter(move || {
let id = ffi::parse_id(id)?;
let dst = unsafe { dst.as_mut() }.ok_or(Error::InvalidPointer)?;
State::lock().video.frame_info(id, dst)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_consume_video_raw_frame_free(id: u32) -> i32 {
ffi::enter(move || {
let id = ffi::parse_id(id)?;
State::lock().video.frame_free(id)
})
}
#[cfg(test)]
mod tests {
use super::*;
async fn track_under_test() -> (
moq_video::encode::Producer<moq_mux::catalog::hang::Extra>,
moq_net::track::Subscriber,
) {
let mut broadcast = moq_net::broadcast::Info::new().produce();
let catalog =
moq_mux::catalog::Producer::with_catalog(&mut broadcast, moq_mux::catalog::hang::Catalog::default())
.unwrap();
let consumer = broadcast.consume();
let rendition = moq_video::encode::Config::new(320, 240, 30).probe().await.unwrap();
let producer = moq_video::encode::Producer::new(broadcast, catalog, rendition).unwrap();
let name = producer.demand().name().to_string();
let track = consumer.track(&name).unwrap().subscribe(None).await.unwrap();
(producer, track)
}
#[tokio::test]
async fn a_successful_drain_ends_the_track_cleanly() {
let (producer, mut track) = track_under_test().await;
finalize(producer, Ok(())).unwrap();
assert!(matches!(track.recv_group().await, Ok(None)), "expected a clean end");
}
#[tokio::test]
async fn a_failed_drain_aborts_the_track() {
let (producer, mut track) = track_under_test().await;
let err = moq_video::Error::Codec(anyhow::anyhow!("the codec lost the tail"));
finalize(producer, Err(err)).unwrap_err();
let Err(err) = track.recv_group().await else {
panic!("expected an abort, not a clean end");
};
assert!(
err.to_string().contains("the codec lost the tail"),
"the abort should carry the drain failure: {err}"
);
}
}