use std::collections::{HashMap, VecDeque};
use std::ffi::c_void;
use std::ffi::CStr;
use std::os::raw::c_int;
use std::sync::{Arc, Mutex};
use crate::media::MediaKind;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FrameDirection {
Send,
Receive,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FrameAction {
Forward,
Drop,
}
pub struct EncodedFrame<'a> {
pub direction: FrameDirection,
pub kind: MediaKind,
pub is_key_frame: bool,
pub payload_type: u8,
pub ssrc: u32,
pub timestamp: u32,
pub mime_type: &'a str,
pub data: &'a [u8],
frame: *mut c_void,
}
impl EncodedFrame<'_> {
pub fn replace_data(&self, data: &[u8]) {
unsafe {
reactor_webrtc_sys::reactor_webrtc_encoded_frame_set_data(
self.frame,
data.as_ptr(),
data.len(),
);
}
}
}
type EncodedCb = Box<dyn for<'a> FnMut(&EncodedFrame<'a>) -> FrameAction + Send>;
struct TransformState {
cb: Mutex<EncodedCb>,
}
extern "C" fn encoded_tramp(
ud: *mut c_void,
frame: *const reactor_webrtc_sys::ReactorEncodedFrame,
) -> c_int {
let f = match unsafe { frame.as_ref() } {
Some(f) => f,
None => return 0,
};
let st = unsafe { &*(ud as *const TransformState) };
let data = if f.data.is_null() || f.data_len == 0 {
&[][..]
} else {
unsafe { std::slice::from_raw_parts(f.data, f.data_len) }
};
let mime = if f.mime_type.is_null() {
""
} else {
unsafe { CStr::from_ptr(f.mime_type) }
.to_str()
.unwrap_or("")
};
let ef = EncodedFrame {
direction: if f.direction == 1 {
FrameDirection::Receive
} else {
FrameDirection::Send
},
kind: if f.is_audio == 1 {
MediaKind::Audio
} else {
MediaKind::Video
},
is_key_frame: f.is_key_frame == 1,
payload_type: f.payload_type,
ssrc: f.ssrc,
timestamp: f.timestamp,
mime_type: mime,
data,
frame: f.frame,
};
let action = match st.cb.lock() {
Ok(mut cb) => cb(&ef),
Err(_) => FrameAction::Forward,
};
match action {
FrameAction::Forward => 0,
FrameAction::Drop => 1,
}
}
extern "C" fn free_state_tramp(ud: *mut c_void) {
drop(unsafe { Box::from_raw(ud as *mut TransformState) });
}
pub struct FrameTransform {
raw: *mut reactor_webrtc_sys::FrameTransformer,
}
unsafe impl Send for FrameTransform {}
unsafe impl Sync for FrameTransform {}
impl FrameTransform {
pub fn new(cb: impl for<'a> FnMut(&EncodedFrame<'a>) -> FrameAction + Send + 'static) -> Self {
let state = Box::into_raw(Box::new(TransformState {
cb: Mutex::new(Box::new(cb)),
}));
let raw = unsafe {
reactor_webrtc_sys::reactor_webrtc_frame_transformer_create(
encoded_tramp,
state as *mut c_void,
free_state_tramp,
)
};
if raw.is_null() {
drop(unsafe { Box::from_raw(state) });
}
Self { raw }
}
pub(crate) fn raw(&self) -> *mut reactor_webrtc_sys::FrameTransformer {
self.raw
}
}
impl Drop for FrameTransform {
fn drop(&mut self) {
if !self.raw.is_null() {
unsafe { reactor_webrtc_sys::reactor_webrtc_frame_transformer_destroy(self.raw) }
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum VideoCodec {
Vp8 = 1,
Vp9 = 2,
Av1 = 3,
H264 = 4,
H265 = 5,
}
impl VideoCodec {
fn from_u32(v: u32) -> Option<Self> {
match v {
1 => Some(Self::Vp8),
2 => Some(Self::Vp9),
3 => Some(Self::Av1),
4 => Some(Self::H264),
5 => Some(Self::H265),
_ => None,
}
}
}
pub struct RawVideoFrame<'a> {
pub codec: VideoCodec,
pub y: &'a [u8],
pub y_stride: u32,
pub u: &'a [u8],
pub u_stride: u32,
pub v: &'a [u8],
pub v_stride: u32,
pub width: u32,
pub height: u32,
pub rtp_timestamp: u32,
pub request_key_frame: bool,
}
pub struct EncodedVideoFrame {
pub data: Vec<u8>,
pub is_key_frame: bool,
pub width: u32,
pub height: u32,
pub rtp_timestamp: u32,
}
type EncodeCallbackBox = Box<dyn FnMut(&RawVideoFrame<'_>) -> Option<EncodedVideoFrame> + Send>;
struct CustomEncoderState {
cb: Mutex<EncodeCallbackBox>,
}
extern "C" fn encode_tramp(
ud: *mut c_void,
raw: *const reactor_webrtc_sys::ReactorRawVideoFrame,
out: *mut reactor_webrtc_sys::ReactorEncodedVideoOutput,
) -> c_int {
let Some(r) = (unsafe { raw.as_ref() }) else {
return 1;
};
let st = unsafe { &*(ud as *const CustomEncoderState) };
let y_len = (r.y_stride.max(0) as usize) * r.height as usize;
let uv_len = (r.u_stride.max(0) as usize) * (r.height as usize).div_ceil(2);
let frame = RawVideoFrame {
codec: VideoCodec::from_u32(r.codec).unwrap_or(VideoCodec::H264),
y: if r.y.is_null() {
&[]
} else {
unsafe { std::slice::from_raw_parts(r.y, y_len) }
},
y_stride: r.y_stride as u32,
u: if r.u.is_null() {
&[]
} else {
unsafe { std::slice::from_raw_parts(r.u, uv_len) }
},
u_stride: r.u_stride as u32,
v: if r.v.is_null() {
&[]
} else {
unsafe { std::slice::from_raw_parts(r.v, uv_len) }
},
v_stride: r.v_stride as u32,
width: r.width,
height: r.height,
rtp_timestamp: r.rtp_timestamp,
request_key_frame: r.request_key_frame != 0,
};
let result = match st.cb.lock() {
Ok(mut cb) => cb(&frame),
Err(_) => return 1,
};
match result {
None => 1,
Some(encoded) => fill_output(encoded, out),
}
}
extern "C" fn free_encoded_data(data: *const u8, len: usize) {
unsafe { drop(Vec::from_raw_parts(data as *mut u8, len, len)) };
}
extern "C" fn free_encoder_state_tramp(ud: *mut c_void) {
drop(unsafe { Box::from_raw(ud as *mut CustomEncoderState) });
}
#[derive(Clone)]
pub(crate) enum RegistrySlot {
Custom(Arc<Mutex<VecDeque<EncodedVideoFrame>>>),
Builtin,
}
pub(crate) struct EncoderRegistry {
pending: Mutex<VecDeque<RegistrySlot>>,
assigned: Mutex<HashMap<u64, RegistrySlot>>,
}
impl EncoderRegistry {
pub(crate) fn new() -> Arc<Self> {
Arc::new(Self {
pending: Mutex::new(VecDeque::new()),
assigned: Mutex::new(HashMap::new()),
})
}
pub(crate) fn add_encoded_slot(&self) -> Arc<Mutex<VecDeque<EncodedVideoFrame>>> {
let q = Arc::new(Mutex::new(VecDeque::new()));
self.pending
.lock()
.unwrap()
.push_back(RegistrySlot::Custom(q.clone()));
q
}
pub(crate) fn add_raw_slot(&self) {
self.pending
.lock()
.unwrap()
.push_back(RegistrySlot::Builtin);
}
pub(crate) fn use_builtin_for(&self, encoder_id: u64) -> bool {
let mut assigned = self.assigned.lock().unwrap();
if let Some(slot) = assigned.get(&encoder_id) {
return matches!(slot, RegistrySlot::Builtin);
}
let slot = self
.pending
.lock()
.unwrap()
.pop_front()
.unwrap_or(RegistrySlot::Builtin);
let is_builtin = matches!(slot, RegistrySlot::Builtin);
assigned.insert(encoder_id, slot);
is_builtin
}
fn pop_for(&self, encoder_id: u64) -> Option<EncodedVideoFrame> {
let mut assigned = self.assigned.lock().unwrap();
let slot = assigned.entry(encoder_id).or_insert_with(|| {
self.pending
.lock()
.unwrap()
.pop_front()
.unwrap_or(RegistrySlot::Builtin)
});
match slot {
RegistrySlot::Custom(q) => {
let q = q.clone();
drop(assigned);
let frame = q.lock().unwrap().pop_front();
frame
}
_ => None,
}
}
}
struct RegistryState {
registry: Arc<EncoderRegistry>,
}
fn fill_output(
encoded: EncodedVideoFrame,
out: *mut reactor_webrtc_sys::ReactorEncodedVideoOutput,
) -> c_int {
let mut v = encoded.data;
v.shrink_to_fit();
let ptr = v.as_ptr();
let len = v.len();
std::mem::forget(v);
unsafe {
let o = &mut *out;
o.data = ptr;
o.len = len;
o.is_key_frame = encoded.is_key_frame as c_int;
o.width = encoded.width;
o.height = encoded.height;
o.rtp_timestamp = encoded.rtp_timestamp;
o.free_data = Some(free_encoded_data);
}
0
}
extern "C" fn registry_encode_tramp(
ud: *mut c_void,
raw: *const reactor_webrtc_sys::ReactorRawVideoFrame,
out: *mut reactor_webrtc_sys::ReactorEncodedVideoOutput,
) -> c_int {
let Some(r) = (unsafe { raw.as_ref() }) else {
return 1;
};
let st = unsafe { &*(ud as *const RegistryState) };
match st.registry.pop_for(r.encoder_id) {
None => 1,
Some(encoded) => fill_output(encoded, out),
}
}
extern "C" fn registry_use_builtin_tramp(ud: *mut c_void, encoder_id: u64) -> c_int {
let st = unsafe { &*(ud as *const RegistryState) };
st.registry.use_builtin_for(encoder_id) as c_int
}
extern "C" fn free_registry_state_tramp(ud: *mut c_void) {
drop(unsafe { Box::from_raw(ud as *mut RegistryState) });
}
pub struct CustomVideoEncoder {
pub(crate) encode_fn: extern "C" fn(
*mut c_void,
*const reactor_webrtc_sys::ReactorRawVideoFrame,
*mut reactor_webrtc_sys::ReactorEncodedVideoOutput,
) -> c_int,
pub(crate) userdata: *mut c_void,
pub(crate) free_ud: Option<extern "C" fn(*mut c_void)>,
pub(crate) use_builtin: Option<extern "C" fn(*mut c_void, u64) -> c_int>,
}
unsafe impl Send for CustomVideoEncoder {}
unsafe impl Sync for CustomVideoEncoder {}
impl CustomVideoEncoder {
pub fn new(
cb: impl FnMut(&RawVideoFrame<'_>) -> Option<EncodedVideoFrame> + Send + 'static,
) -> Self {
let state = Box::into_raw(Box::new(CustomEncoderState {
cb: Mutex::new(Box::new(cb)),
}));
Self {
encode_fn: encode_tramp,
userdata: state as *mut c_void,
free_ud: Some(free_encoder_state_tramp),
use_builtin: None,
}
}
pub(crate) fn from_queue(queue: Arc<Mutex<VecDeque<EncodedVideoFrame>>>) -> Self {
Self::new(move |_raw| queue.lock().unwrap().pop_front())
}
pub(crate) fn from_registry(registry: Arc<EncoderRegistry>) -> Self {
let state = Box::into_raw(Box::new(RegistryState { registry }));
Self {
encode_fn: registry_encode_tramp,
userdata: state as *mut c_void,
free_ud: Some(free_registry_state_tramp),
use_builtin: Some(registry_use_builtin_tramp),
}
}
}
pub struct EncodedVideoTrack {
pub(crate) track: crate::media::Track,
pub(crate) queue: Arc<Mutex<VecDeque<EncodedVideoFrame>>>,
dummy: Vec<u8>,
width: u32,
height: u32,
}
unsafe impl Send for EncodedVideoTrack {}
unsafe impl Sync for EncodedVideoTrack {}
impl EncodedVideoTrack {
pub(crate) fn new(
track: crate::media::Track,
queue: Arc<Mutex<VecDeque<EncodedVideoFrame>>>,
width: u32,
height: u32,
) -> Self {
let dummy = vec![0u8; (width * height * 4) as usize];
Self {
track,
queue,
dummy,
width,
height,
}
}
pub fn track(&self) -> &crate::media::Track {
&self.track
}
pub fn push_encoded_frame(&self, frame: EncodedVideoFrame) {
self.queue.lock().unwrap().push_back(frame);
self.track
.push_video_frame(&self.dummy, self.width, self.height);
}
}