use std::ffi::c_void;
use std::time::Duration;
use tokio::sync::oneshot;
use crate::ffi::OnStatus;
use crate::{Error, Id, NonZeroSlab, State, ffi};
#[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 {
consumer_tasks: NonZeroSlab<Option<VideoTaskEntry>>,
frames: NonZeroSlab<VideoFrame>,
}
struct VideoFrame {
timestamp_us: u64,
width: u32,
height: u32,
data: bytes::Bytes,
}
struct VideoTaskEntry {
close: Option<oneshot::Sender<()>>,
callback: OnStatus,
}
impl Video {
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 frame = VideoFrame {
timestamp_us: frame.timestamp.as_micros() as u64,
width: frame.size.width,
height: frame.size.height,
data: frame.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(())
}
}
#[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)
})
}