use std::ffi::{c_char, c_void};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
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 max_age_us: u64,
pub format: u32,
pub width: u32,
pub height: u32,
}
#[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)
}
fn video_ceiling(
config: &moq_video::encode::Config,
rendition: &hang::catalog::VideoConfig,
) -> moq_net::bandwidth::Rate {
config
.bitrate
.or_else(|| rendition.bitrate.map(moq_net::bandwidth::Rate::from_bps))
.unwrap_or_else(|| {
moq_net::bandwidth::Rate::from_bps(
(config.size().pixels() as f64 * config.framerate.as_f64() * 0.07) as u64,
)
})
}
async fn follow_reservation(
inner: Shared<VideoEncoder>,
mut consumer: moq_net::bandwidth::Consumer,
ceiling: Arc<AtomicU64>,
) {
use moq_mux::rate::{Control, Policy};
let mut max = moq_net::bandwidth::Rate::from_bps(ceiling.load(Ordering::SeqCst));
let mut control = Control::new(Policy::new(max));
loop {
let estimate = match consumer.changed().await {
Ok(estimate) => estimate,
Err(_) => return,
};
let next = moq_net::bandwidth::Rate::from_bps(ceiling.load(Ordering::SeqCst));
if next != max {
max = next;
control = Control::new(Policy::new(max));
}
let Some(bitrate) = control.update(estimate, Instant::now()) else {
continue;
};
let mut guard = inner.lock();
let Some(producer) = guard.as_mut() else {
return;
};
match block_on(producer.encoder.set_bitrate(bitrate)) {
Ok(()) => tracing::debug!(bitrate = bitrate.as_bps(), "adjusted encoder bitrate"),
Err(moq_video::Error::BitrateUnsupported(name)) => {
tracing::warn!(encoder = name, "encoder cannot follow the bandwidth estimate");
return;
}
Err(err) => {
tracing::warn!(error = %err, bitrate = bitrate.as_bps(), "failed to adjust encoder bitrate");
}
}
}
}
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,
reservation: Option<Arc<moq_net::bandwidth::Reservation>>,
follow: Option<oneshot::Sender<()>>,
ceiling: Option<Arc<AtomicU64>>,
}
struct VideoFrame {
timestamp_us: u64,
width: u32,
height: u32,
data: bytes::Bytes,
}
#[derive(Clone, Copy)]
pub struct DecoderOutput {
format: moq_video_pixel_format,
size: Option<moq_video::Size>,
}
fn finalize(
mut 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, 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) -> Result<(), Error> {
Ok(block_on(self.encoder.cut())?)
}
fn publish_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
block_on(self.encoder.set_bitrate(moq_net::bandwidth::Rate::from_bps(bitrate)))?;
if let Some(ceiling) = &self.ceiling {
ceiling.store(bitrate, Ordering::SeqCst);
}
if let Some(reservation) = &self.reservation {
reservation.update(moq_net::bandwidth::Rate::from_bps(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(),
reservation: None,
follow: None,
ceiling: None,
}))
}
pub(crate) fn follow(
&self,
id: Id,
allocator: &moq_net::bandwidth::Allocator,
ceiling: moq_net::bandwidth::Rate,
) -> Result<(), Error> {
let shared = self.producer(id)?;
let max = ceiling;
let ceiling = Arc::new(AtomicU64::new(max.as_bps()));
let reservation = {
let mut guard = shared.lock();
let encoder = guard.as_mut().ok_or(Error::MediaNotFound)?;
let reservation = Arc::new(allocator.reserve(&encoder.producer.demand(), max));
encoder.reservation = Some(reservation.clone());
encoder.ceiling = Some(ceiling.clone());
reservation
};
let (close, closed) = oneshot::channel();
shared.lock().as_mut().ok_or(Error::MediaNotFound)?.follow = Some(close);
let follower = shared.clone();
tokio::spawn(async move {
tokio::select! {
biased;
_ = closed => {}
_ = follow_reservation(follower, reservation.consumer(), ceiling) => {}
}
});
Ok(())
}
pub(crate) fn reservation(&self, id: Id) -> Result<Option<Arc<moq_net::bandwidth::Reservation>>, Error> {
Ok(self
.producer(id)?
.lock()
.as_ref()
.ok_or(Error::MediaNotFound)?
.reservation
.clone())
}
pub(crate) fn demand(&self, id: Id) -> Result<moq_net::track::Demand, Error> {
Ok(self
.producer(id)?
.lock()
.as_ref()
.ok_or(Error::MediaNotFound)?
.producer
.demand())
}
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,
options: moq_video::decode::Options,
output: DecoderOutput,
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, options).await?;
Self::run(on_frame, consumer, channel.1, output).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<()>,
output: DecoderOutput,
) -> Result<(), Error> {
loop {
let frame = tokio::select! {
biased;
_ = &mut close => return Ok(()),
frame = consumer.read() => match frame? {
Some(frame) => frame,
None => return Ok(()),
},
};
let mut frame = frame;
if let Some(size) = output.size
&& frame.size() != size
{
frame = frame.resize(size, &moq_video::resize::Config::default())?;
}
let size = frame.size();
let data = match output.format {
moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 => {
bytes::Bytes::from(frame.surface.into_i420()?.into_data())
}
moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA => bytes::Bytes::from(
frame
.surface
.to_rgba(&moq_video::convert::Config::default())?
.into_data(),
),
};
let frame = VideoFrame {
timestamp_us: frame.timestamp.as_micros() as u64,
width: size.width,
height: size.height,
data,
};
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 decoder_size(width: u32, height: u32) -> Result<Option<moq_video::Size>, Error> {
if width == 0 && height == 0 {
return Ok(None);
}
if width == 0 || height == 0 || !width.is_multiple_of(2) || !height.is_multiple_of(2) {
return Err(Error::InvalidConfig(format!(
"decode size {width}x{height}: use 0x0 for the native size or even non-zero dimensions"
)));
}
let size = moq_video::Size::new(width, height);
if size
.pixels()
.checked_mul(4)
.is_none_or(|bytes| usize::try_from(bytes).is_err())
{
return Err(Error::InvalidConfig(format!(
"decode size {width}x{height}: dimensions too large to represent"
)));
}
Ok(Some(size))
}
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_encode_video(
broadcast: u32,
input: *const moq_video_encoder_input,
output: *const moq_video_encoder_output,
bandwidth: u32,
) -> 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 framerate = moq_video::Rate::new(raw_input.framerate, 1)
.map_err(|_| Error::Video(moq_video::Error::InvalidFramerate(raw_input.framerate).into()))?;
let mut config = moq_video::encode::Config::new(raw_input.width, raw_input.height, framerate);
config.codec = codec_from_u32(raw_output.codec)?;
config.kind = unsafe { encoder_kind(raw_output)? };
config.bitrate = (raw_output.bitrate != 0).then(|| moq_net::bandwidth::Rate::from_bps(raw_output.bitrate));
if raw_output.gop != 0 {
config.gop = moq_video::encode::Gop::Keyframe {
interval: raw_output.gop,
};
}
let rendition = block_on(config.probe())?;
let encoder = block_on(moq_video::encode::Sink::open(&config))?;
let bandwidth = ffi::parse_id_optional(bandwidth)?;
let mut state = State::lock();
let allocator = bandwidth.map(|id| state.bandwidth.allocator(id)).transpose()?;
let State { publish, video, .. } = &mut *state;
let (broadcast_producer, catalog) = publish.pair_mut(broadcast)?;
let id = video.publish(
broadcast_producer,
catalog.clone(),
format,
&config,
rendition.clone(),
encoder,
)?;
if let Some(allocator) = allocator.as_ref() {
video.follow(id, allocator, video_ceiling(&config, &rendition))?;
}
Ok(id)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_encode_video_reservation(producer: u32) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let mut state = State::lock();
match state.video.reservation(producer)? {
Some(reservation) => Ok(i32::from(state.bandwidth.hold(reservation)?)),
None => Ok(0),
}
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_encode_video_demand(
producer: u32,
on_demand: crate::moq_status_callback,
user_data: *mut c_void,
) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let on_demand = unsafe { OnStatus::new(user_data, on_demand)? };
let mut state = State::lock();
let demand = state.video.demand(producer)?;
state.publish.demand(demand, on_demand)
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_encode_video_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_encode_video_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()
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_encode_video_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_encode_video_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_decode_video(
catalog: u32,
index: u32,
output: *const moq_video_decoder_output,
on_frame: crate::moq_status_callback,
user_data: *mut c_void,
) -> i32 {
ffi::enter(move || {
let raw = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
let format = pixel_format_from_u32(raw.format)?;
let size = decoder_size(raw.width, raw.height)?;
let catalog = ffi::parse_id(catalog)?;
let mut options = moq_video::decode::Options::new();
options.start = moq_video::decode::Start::Latest;
options.max_age = Duration::from_micros(raw.max_age_us);
options.decoder.output = moq_video::Output::Cpu;
options.decoder.scale_hint = size;
let output = DecoderOutput { format, size };
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, options, output, on_frame)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_decode_video_cancel(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_decode_video_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_decode_video_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 config = moq_mux::catalog::Config::default()
.with_catalog(moq_mux::catalog::hang::Catalog::<moq_mux::catalog::hang::Extra>::default());
let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config).unwrap();
let consumer = broadcast.consume();
let rendition = moq_video::encode::Config::new(320, 240, moq_video::Rate::new(30, 1).unwrap())
.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}"
);
}
}