finish proto of livekit-ffi (#36)
* ffi req * handle req_id * remaining room events * fix * Update ffi.proto * use sid instead of info * state -> stream_state * "s" * participant is optional on data received * add ReleaseHandleRequest * Update ffi.proto * add DataPacketKind * VideoSinkInfo * append info It looks nicer on the FFI languages * wip * keep it simple * Update ffi.proto * update server to latest proto * to_i420 * i420_to_abgr * to_argb * add publications to ParticipantInfo It'll be used when we join a Room * cleanup + DisposeRequest * fix compilation * send the whole buffer info with to_i420
This commit is contained in:
@@ -10,7 +10,7 @@ use std::sync::Arc;
|
||||
|
||||
impl From<FFIHandleId> for proto::FfiHandleId {
|
||||
fn from(id: FFIHandleId) -> Self {
|
||||
Self { id: id as u32 }
|
||||
Self { id: id as u64 }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -23,6 +23,7 @@ macro_rules! impl_participant_into {
|
||||
sid: p.sid().to_string(),
|
||||
identity: p.identity().to_string(),
|
||||
metadata: p.metadata(),
|
||||
publications: p.tracks().iter().map(|(_, p)| p.into()).collect(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -33,6 +34,18 @@ impl_participant_into!(&Arc<LocalParticipant>);
|
||||
impl_participant_into!(&Arc<RemoteParticipant>);
|
||||
impl_participant_into!(&Participant);
|
||||
|
||||
impl From<TrackSource> for proto::TrackSource {
|
||||
fn from(source: TrackSource) -> proto::TrackSource {
|
||||
match source {
|
||||
TrackSource::Unknown => proto::TrackSource::SourceUnknown,
|
||||
TrackSource::Camera => proto::TrackSource::SourceCamera,
|
||||
TrackSource::Microphone => proto::TrackSource::SourceMicrophone,
|
||||
TrackSource::Screenshare => proto::TrackSource::SourceScreenshare,
|
||||
TrackSource::ScreenshareAudio => proto::TrackSource::SourceScreenshareAudio,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
macro_rules! impl_publication_into {
|
||||
($p:ty) => {
|
||||
impl From<$p> for proto::TrackPublicationInfo {
|
||||
@@ -41,6 +54,14 @@ macro_rules! impl_publication_into {
|
||||
name: p.name(),
|
||||
sid: p.sid().to_string(),
|
||||
kind: proto::TrackKind::from(p.kind()).into(),
|
||||
source: proto::TrackSource::from(p.source()).into(),
|
||||
dimension: Some(proto::Dimension {
|
||||
width: p.dimension().0,
|
||||
height: p.dimension().1,
|
||||
}),
|
||||
mime_type: p.mime_type(),
|
||||
simulcasted: p.simulcasted(),
|
||||
muted: p.muted(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -57,7 +78,7 @@ macro_rules! impl_track_into {
|
||||
fn from(track: $t) -> Self {
|
||||
Self {
|
||||
name: track.name(),
|
||||
state: proto::StreamState::from(track.stream_state()).into(),
|
||||
stream_state: proto::StreamState::from(track.stream_state()).into(),
|
||||
sid: track.sid().to_string(),
|
||||
kind: proto::TrackKind::from(track.kind()).into(),
|
||||
muted: track.muted(),
|
||||
@@ -125,7 +146,7 @@ impl proto::RoomEvent {
|
||||
} => Some(proto::room_event::Message::TrackUnpublished(
|
||||
proto::TrackUnpublished {
|
||||
participant_sid: participant.sid().to_string(),
|
||||
publication: Some((&publication).into()),
|
||||
publication_sid: publication.sid().into(),
|
||||
},
|
||||
)),
|
||||
RoomEvent::TrackSubscribed {
|
||||
@@ -136,6 +157,9 @@ impl proto::RoomEvent {
|
||||
proto::TrackSubscribed {
|
||||
participant_sid: participant.sid().to_string(),
|
||||
track: Some((&track).into()),
|
||||
sink: Some(proto::VideoSinkInfo {
|
||||
track_sid: track.sid().to_string(),
|
||||
}),
|
||||
},
|
||||
)),
|
||||
RoomEvent::TrackUnsubscribed {
|
||||
@@ -145,7 +169,7 @@ impl proto::RoomEvent {
|
||||
} => Some(proto::room_event::Message::TrackUnsubscribed(
|
||||
proto::TrackUnsubscribed {
|
||||
participant_sid: participant.sid().to_string(),
|
||||
track: Some((&track).into()),
|
||||
track_sid: track.sid().to_string(),
|
||||
},
|
||||
)),
|
||||
_ => None,
|
||||
@@ -169,7 +193,7 @@ impl From<VideoRotation> for proto::VideoRotation {
|
||||
}
|
||||
}
|
||||
|
||||
impl From<VideoFrame> for proto::VideoFrame {
|
||||
impl From<VideoFrame> for proto::VideoFrameInfo {
|
||||
fn from(frame: VideoFrame) -> Self {
|
||||
Self {
|
||||
width: frame.width(),
|
||||
@@ -201,7 +225,7 @@ impl From<VideoFrameBufferType> for proto::VideoFrameBufferType {
|
||||
|
||||
macro_rules! impl_yuv_into {
|
||||
($b:ty) => {
|
||||
impl From<$b> for proto::PlanarYuvBuffer {
|
||||
impl From<$b> for proto::PlanarYuvBufferInfo {
|
||||
fn from(buffer: $b) -> Self {
|
||||
Self {
|
||||
chroma_width: buffer.chroma_width(),
|
||||
@@ -226,7 +250,7 @@ impl_yuv_into!(&I010Buffer);
|
||||
|
||||
macro_rules! impl_biyuv_into {
|
||||
($b:ty) => {
|
||||
impl From<$b> for proto::BiplanarYuvBuffer {
|
||||
impl From<$b> for proto::BiplanarYuvBufferInfo {
|
||||
fn from(buffer: $b) -> Self {
|
||||
Self {
|
||||
chroma_width: buffer.chroma_width(),
|
||||
@@ -243,7 +267,7 @@ macro_rules! impl_biyuv_into {
|
||||
|
||||
impl_biyuv_into!(&NV12Buffer);
|
||||
|
||||
impl proto::VideoFrameBuffer {
|
||||
impl proto::VideoFrameBufferInfo {
|
||||
pub fn from(handle_id: FFIHandleId, buffer: &VideoFrameBuffer) -> Self {
|
||||
Self {
|
||||
handle: Some(handle_id.into()),
|
||||
@@ -252,19 +276,54 @@ impl proto::VideoFrameBuffer {
|
||||
height: buffer.height(),
|
||||
buffer: Some(match &buffer {
|
||||
VideoFrameBuffer::Native(_) => {
|
||||
proto::video_frame_buffer::Buffer::Native(proto::NativeBuffer {})
|
||||
proto::video_frame_buffer_info::Buffer::Native(proto::NativeBufferInfo {})
|
||||
}
|
||||
VideoFrameBuffer::I420(i420) => {
|
||||
proto::video_frame_buffer_info::Buffer::Yuv(i420.into())
|
||||
}
|
||||
VideoFrameBuffer::I420(i420) => proto::video_frame_buffer::Buffer::Yuv(i420.into()),
|
||||
VideoFrameBuffer::I420A(i420a) => {
|
||||
proto::video_frame_buffer::Buffer::Yuv(i420a.into())
|
||||
proto::video_frame_buffer_info::Buffer::Yuv(i420a.into())
|
||||
}
|
||||
VideoFrameBuffer::I422(i422) => {
|
||||
proto::video_frame_buffer_info::Buffer::Yuv(i422.into())
|
||||
}
|
||||
VideoFrameBuffer::I444(i444) => {
|
||||
proto::video_frame_buffer_info::Buffer::Yuv(i444.into())
|
||||
}
|
||||
VideoFrameBuffer::I010(i010) => {
|
||||
proto::video_frame_buffer_info::Buffer::Yuv(i010.into())
|
||||
}
|
||||
VideoFrameBuffer::I422(i422) => proto::video_frame_buffer::Buffer::Yuv(i422.into()),
|
||||
VideoFrameBuffer::I444(i444) => proto::video_frame_buffer::Buffer::Yuv(i444.into()),
|
||||
VideoFrameBuffer::I010(i010) => proto::video_frame_buffer::Buffer::Yuv(i010.into()),
|
||||
VideoFrameBuffer::NV12(nv12) => {
|
||||
proto::video_frame_buffer::Buffer::BiYuv(nv12.into())
|
||||
proto::video_frame_buffer_info::Buffer::BiYuv(nv12.into())
|
||||
}
|
||||
}),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<proto::VideoFormatType> for VideoFormatType {
|
||||
fn from(format: proto::VideoFormatType) -> Self {
|
||||
match format {
|
||||
proto::VideoFormatType::FormatArgb => Self::ARGB,
|
||||
proto::VideoFormatType::FormatBgra => Self::BGRA,
|
||||
proto::VideoFormatType::FormatAbgr => Self::ABGR,
|
||||
proto::VideoFormatType::FormatRgba => Self::RGBA,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<&RoomSession> for proto::RoomInfo {
|
||||
fn from(session: &RoomSession) -> Self {
|
||||
Self {
|
||||
sid: session.sid().into(),
|
||||
name: session.name(),
|
||||
metadata: session.metadata(),
|
||||
local_participant: Some((&session.local_participant()).into()),
|
||||
participants: session
|
||||
.participants()
|
||||
.iter()
|
||||
.map(|(_, p)| p.into())
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+168
-118
@@ -1,30 +1,30 @@
|
||||
use crate::{
|
||||
proto, proto::ffi_request::Message as FFIRequest, proto::ffi_response::Message as FFIResponse,
|
||||
};
|
||||
use crate::proto;
|
||||
use lazy_static::lazy_static;
|
||||
use livekit::prelude::*;
|
||||
use livekit::webrtc::media_stream::OnFrameHandler;
|
||||
use parking_lot::{Mutex, RwLock};
|
||||
use prost::Message;
|
||||
use std::any::Any;
|
||||
use std::collections::HashMap;
|
||||
use std::panic;
|
||||
use std::slice;
|
||||
use std::sync::atomic::AtomicU32;
|
||||
use std::sync::atomic::AtomicU64;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use thiserror::Error;
|
||||
use tokio::sync::oneshot;
|
||||
use tokio::task::JoinHandle;
|
||||
|
||||
mod conversion;
|
||||
mod room;
|
||||
|
||||
#[derive(Error, Debug)]
|
||||
pub enum FFIError {
|
||||
#[error("the FFIServer isn't configured")]
|
||||
NotConfigured,
|
||||
#[error("failed to execute the ffi callback")]
|
||||
#[error("failed to execute the FFICallback")]
|
||||
CallbackFailed,
|
||||
}
|
||||
|
||||
pub type FFIHandleId = u32;
|
||||
pub type FFIHandleId = usize;
|
||||
pub type FFIHandle = Box<dyn Any + Send + Sync>;
|
||||
|
||||
type CallbackFn = unsafe extern "C" fn(*const u8, usize); // This "C" callback must be threadsafe
|
||||
@@ -42,10 +42,13 @@ pub struct FFIConfig {
|
||||
pub struct FFIServer {
|
||||
// Object owned by the foreign language
|
||||
// The foreign language is responsible for freeing this memory
|
||||
//
|
||||
// NOTE: For VideoBuffers, we always store the enum VideoFrameBuffer
|
||||
ffi_owned: RwLock<HashMap<FFIHandleId, FFIHandle>>,
|
||||
next_handle: AtomicU32, // FFIHandle
|
||||
next_handle_id: AtomicU64, // FFIHandleId
|
||||
next_async_id: AtomicU64,
|
||||
|
||||
rooms: RwLock<HashMap<RoomSid, Room>>,
|
||||
rooms: RwLock<HashMap<RoomSid, (JoinHandle<()>, oneshot::Sender<()>)>>,
|
||||
async_runtime: tokio::runtime::Runtime,
|
||||
initialized: AtomicBool,
|
||||
config: Mutex<Option<FFIConfig>>,
|
||||
@@ -55,7 +58,8 @@ impl Default for FFIServer {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
ffi_owned: RwLock::new(HashMap::new()),
|
||||
next_handle: Default::default(),
|
||||
next_handle_id: AtomicU64::new(1), // 0 is considered invalid
|
||||
next_async_id: AtomicU64::new(1),
|
||||
rooms: RwLock::new(HashMap::new()),
|
||||
async_runtime: tokio::runtime::Builder::new_multi_thread()
|
||||
.enable_all()
|
||||
@@ -68,8 +72,45 @@ impl Default for FFIServer {
|
||||
}
|
||||
|
||||
impl FFIServer {
|
||||
pub fn initialize(&self, init: &proto::InitializeRequest) {
|
||||
if self.initialized() {
|
||||
self.dispose();
|
||||
}
|
||||
|
||||
self.initialized.store(true, Ordering::SeqCst);
|
||||
*self.config.lock() = Some(FFIConfig {
|
||||
callback_fn: unsafe { std::mem::transmute(init.event_callback_ptr) },
|
||||
});
|
||||
}
|
||||
|
||||
pub fn dispose(&self) {
|
||||
self.initialized.store(false, Ordering::SeqCst);
|
||||
*self.config.lock() = None;
|
||||
self.async_runtime.block_on(self.close());
|
||||
}
|
||||
|
||||
pub async fn close(&self) {
|
||||
// Close all rooms
|
||||
for (k, (handle, shutdown_tx)) in self.rooms.write().drain() {
|
||||
let _ = shutdown_tx.send(());
|
||||
let _ = handle.await;
|
||||
}
|
||||
}
|
||||
|
||||
pub fn add_room(&self, sid: RoomSid, handle: (JoinHandle<()>, oneshot::Sender<()>)) {
|
||||
self.rooms.write().insert(sid, handle);
|
||||
}
|
||||
|
||||
pub fn initialized(&self) -> bool {
|
||||
self.initialized.load(Ordering::SeqCst)
|
||||
}
|
||||
|
||||
pub fn next_handle_id(&self) -> FFIHandleId {
|
||||
self.next_handle.fetch_add(1, Ordering::SeqCst) as FFIHandleId
|
||||
self.next_handle_id.fetch_add(1, Ordering::SeqCst) as FFIHandleId
|
||||
}
|
||||
|
||||
pub fn next_async_id(&self) -> u64 {
|
||||
self.next_async_id.fetch_add(1, Ordering::SeqCst)
|
||||
}
|
||||
|
||||
pub fn insert_handle(&self, handle_id: FFIHandleId, handle: FFIHandle) {
|
||||
@@ -80,142 +121,151 @@ impl FFIServer {
|
||||
self.ffi_owned.write().remove(&handle_id)
|
||||
}
|
||||
|
||||
pub fn send_response(&self, message: FFIResponse) -> Result<(), FFIError> {
|
||||
if !self.initialized.load(Ordering::SeqCst) {
|
||||
pub fn send_event(
|
||||
&self,
|
||||
message: proto::ffi_event::Message,
|
||||
async_id: Option<u64>,
|
||||
) -> Result<(), FFIError> {
|
||||
let config = self.config.lock();
|
||||
|
||||
if !self.initialized() {
|
||||
Err(FFIError::NotConfigured)?
|
||||
}
|
||||
|
||||
let message = proto::FfiResponse {
|
||||
let message = proto::FfiEvent {
|
||||
async_id,
|
||||
message: Some(message),
|
||||
}
|
||||
.encode_to_vec();
|
||||
|
||||
let callback_fn = self.config.lock().as_ref().unwrap().callback_fn;
|
||||
let config = config.as_ref().unwrap();
|
||||
if let Err(err) = panic::catch_unwind(|| unsafe {
|
||||
callback_fn(message.as_ptr(), message.len());
|
||||
(config.callback_fn)(message.as_ptr(), message.len());
|
||||
}) {
|
||||
eprintln!("panic when sending ffi response: {:?}", err);
|
||||
eprintln!("panic when sending ffi event: {:?}", err);
|
||||
Err(FFIError::CallbackFailed)?
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn on_request_received(&self, message: FFIRequest) -> Result<(), FFIError> {
|
||||
if let FFIRequest::Configure(ref init) = message {
|
||||
self.initialized.store(true, Ordering::SeqCst);
|
||||
*self.config.lock() = Some(FFIConfig {
|
||||
callback_fn: unsafe { std::mem::transmute(init.callback_ptr) },
|
||||
});
|
||||
}
|
||||
|
||||
if !self.initialized.load(Ordering::SeqCst) {
|
||||
Err(FFIError::NotConfigured)?
|
||||
}
|
||||
|
||||
pub fn handle_request(&self, message: proto::ffi_request::Message) -> proto::FfiResponse {
|
||||
match message {
|
||||
proto::ffi_request::Message::AsyncConnect(connect) => {
|
||||
self.async_runtime.spawn(room_task(connect));
|
||||
let async_id = self.next_async_id();
|
||||
self.async_runtime
|
||||
.spawn(room::create_room(&FFI_SERVER, async_id, connect));
|
||||
|
||||
return proto::FfiResponse {
|
||||
async_id: Some(async_id),
|
||||
..Default::default()
|
||||
};
|
||||
}
|
||||
_ => {}
|
||||
};
|
||||
proto::ffi_request::Message::ToI420(to_i420) => {
|
||||
let mut buffer_info = None;
|
||||
let buffer = self.release_handle(to_i420.buffer.unwrap().id as FFIHandleId);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
if let Some(buffer) = buffer {
|
||||
if let Ok(buffer) = buffer.downcast::<VideoFrameBuffer>() {
|
||||
let handle_id = self.next_handle_id();
|
||||
let i420 = VideoFrameBuffer::I420(buffer.to_i420());
|
||||
buffer_info = Some(proto::VideoFrameBufferInfo::from(handle_id, &i420));
|
||||
self.insert_handle(handle_id, Box::new(i420));
|
||||
}
|
||||
}
|
||||
|
||||
#[no_mangle]
|
||||
pub extern "C" fn livekit_ffi_request(data: *const u8, len: usize) {
|
||||
let data = unsafe { slice::from_raw_parts(data, len) };
|
||||
let request = proto::FfiRequest::decode(data).expect("Failed to decode the FFIRequest");
|
||||
let res = FFI_SERVER.on_request_received(request.message.unwrap());
|
||||
if let Err(err) = res {
|
||||
eprintln!("failed to handle ffi request: {:?}", err);
|
||||
}
|
||||
}
|
||||
|
||||
// Connect a listen to Room events
|
||||
async fn room_task(connect: proto::ConnectRequest) {
|
||||
let res = Room::connect(&connect.url, &connect.token).await;
|
||||
|
||||
if res.is_err() {
|
||||
let _ = FFI_SERVER.send_response(FFIResponse::AsyncConnect(proto::ConnectResponse {
|
||||
success: false,
|
||||
room: None,
|
||||
}));
|
||||
return;
|
||||
}
|
||||
|
||||
// Send connect response before listening to events
|
||||
let (room, mut events) = res.unwrap();
|
||||
let session = room.session();
|
||||
|
||||
let _ = FFI_SERVER.send_response(FFIResponse::AsyncConnect(proto::ConnectResponse {
|
||||
success: true,
|
||||
room: Some(proto::RoomInfo {
|
||||
sid: session.sid(),
|
||||
name: session.name(),
|
||||
local_participant: Some((&room.session().local_participant()).into()),
|
||||
participants: room
|
||||
.session()
|
||||
.participants()
|
||||
.iter()
|
||||
.map(|(_, p)| p.into())
|
||||
.collect(),
|
||||
}),
|
||||
}));
|
||||
|
||||
// Listen to events
|
||||
tokio::spawn(participant_task(Participant::Local(
|
||||
session.local_participant(),
|
||||
)));
|
||||
|
||||
while let Some(event) = events.recv().await {
|
||||
if let Some(event) = proto::RoomEvent::from(session.sid(), event.clone()) {
|
||||
let _ = FFI_SERVER.send_response(FFIResponse::RoomEvent(event));
|
||||
}
|
||||
|
||||
match event {
|
||||
RoomEvent::ParticipantConnected(p) => {
|
||||
tokio::spawn(participant_task(Participant::Remote(p)));
|
||||
return proto::FfiResponse {
|
||||
message: Some(proto::ffi_response::Message::ToI420(
|
||||
proto::ToI420Response {
|
||||
new_buffer: buffer_info,
|
||||
},
|
||||
)),
|
||||
..Default::default()
|
||||
};
|
||||
}
|
||||
RoomEvent::TrackSubscribed {
|
||||
track,
|
||||
publication,
|
||||
participant,
|
||||
} => {
|
||||
if let RemoteTrackHandle::Video(video_track) = track {
|
||||
let rtc_track = video_track.rtc_track();
|
||||
rtc_track.on_frame(on_video_frame(video_track.sid()));
|
||||
proto::ffi_request::Message::ToArgb(to_argb) => {
|
||||
let ffi_owned = self.ffi_owned.read();
|
||||
let buffer = ffi_owned.get(&(to_argb.buffer.unwrap().id as FFIHandleId));
|
||||
|
||||
if let Some(buffer) = buffer {
|
||||
if let Some(buffer) = buffer.downcast_ref::<VideoFrameBuffer>() {
|
||||
let dst_buf = unsafe {
|
||||
slice::from_raw_parts_mut(
|
||||
to_argb.dst_ptr as *mut u8,
|
||||
(to_argb.dst_stride * to_argb.dst_height) as usize,
|
||||
)
|
||||
};
|
||||
|
||||
if let Err(err) = buffer.to_argb(
|
||||
proto::VideoFormatType::from_i32(to_argb.dst_format)
|
||||
.unwrap()
|
||||
.into(),
|
||||
dst_buf,
|
||||
to_argb.dst_stride,
|
||||
to_argb.dst_width,
|
||||
to_argb.dst_height,
|
||||
) {
|
||||
eprintln!("failed to convert videoframe to argb: {:?}", err);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
|
||||
proto::FfiResponse::default()
|
||||
}
|
||||
}
|
||||
|
||||
// Listen to participant events
|
||||
async fn participant_task(participant: Participant) {
|
||||
let mut participant_events = participant.register_observer();
|
||||
while let Some(event) = participant_events.recv().await {
|
||||
// TODO convert event to proto
|
||||
/// This function is threadsafe, this is useful to run synchronous requests in another thread (e.g
|
||||
/// color conversion)
|
||||
#[no_mangle]
|
||||
pub extern "C" fn livekit_ffi_request(
|
||||
data: *const u8,
|
||||
len: usize,
|
||||
data_ptr: *mut *const u8,
|
||||
data_len: *mut usize,
|
||||
) -> FFIHandleId {
|
||||
let data = unsafe { slice::from_raw_parts(data, len) };
|
||||
let res = proto::FfiRequest::decode(data);
|
||||
if let Err(ref err) = res {
|
||||
eprintln!("failed to decode FfiRequest: {:?}", err);
|
||||
return 0;
|
||||
}
|
||||
|
||||
if res.as_ref().unwrap().message.is_none() {
|
||||
eprintln!("request message is empty");
|
||||
return 0;
|
||||
}
|
||||
|
||||
let message = res.unwrap().message.unwrap();
|
||||
if let proto::ffi_request::Message::Initialize(ref init) = message {
|
||||
FFI_SERVER.initialize(init);
|
||||
}
|
||||
|
||||
if let proto::ffi_request::Message::Dispose(_) = message {
|
||||
FFI_SERVER.dispose();
|
||||
}
|
||||
|
||||
if !FFI_SERVER.initialized() {
|
||||
eprintln!("the FFIServer isn't initialized");
|
||||
return 0;
|
||||
}
|
||||
|
||||
let res = FFI_SERVER.handle_request(message);
|
||||
let buf = res.encode_to_vec();
|
||||
|
||||
unsafe {
|
||||
*data_ptr = buf.as_ptr();
|
||||
*data_len = buf.len();
|
||||
}
|
||||
|
||||
let handle_id = FFI_SERVER.next_handle_id();
|
||||
FFI_SERVER.insert_handle(handle_id, Box::new(buf));
|
||||
handle_id
|
||||
}
|
||||
|
||||
fn on_video_frame(track_sid: TrackSid) -> OnFrameHandler {
|
||||
Box::new(move |frame, buffer| {
|
||||
let handle_id = FFI_SERVER.next_handle_id();
|
||||
let proto_buffer = proto::VideoFrameBuffer::from(handle_id, &buffer);
|
||||
FFI_SERVER.insert_handle(handle_id, Box::new(buffer));
|
||||
|
||||
let _ = FFI_SERVER.send_response(FFIResponse::TrackEvent(proto::TrackEvent {
|
||||
track_sid: track_sid.to_string(),
|
||||
message: Some(proto::track_event::Message::FrameReceived(
|
||||
proto::FrameReceived {
|
||||
frame: Some(frame.into()),
|
||||
frame_buffer: Some(proto_buffer),
|
||||
},
|
||||
)),
|
||||
}));
|
||||
})
|
||||
#[no_mangle]
|
||||
pub extern "C" fn livekit_ffi_drop_handle(handle_id: FFIHandleId) -> bool {
|
||||
FFI_SERVER.release_handle(handle_id).is_some() // Free the memory
|
||||
}
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
use crate::proto::{self};
|
||||
use crate::server::FFIServer;
|
||||
use livekit::prelude::*;
|
||||
use tokio::sync::{mpsc, oneshot};
|
||||
|
||||
pub async fn create_room(
|
||||
server: &'static FFIServer,
|
||||
async_id: u64,
|
||||
connect: proto::ConnectRequest,
|
||||
) {
|
||||
let res = Room::connect(&connect.url, &connect.token).await;
|
||||
if let Err(err) = &res {
|
||||
// Failed to connect to the room
|
||||
let _ = server.send_event(
|
||||
proto::ffi_event::Message::ConnectEvent(proto::ConnectEvent {
|
||||
success: false,
|
||||
room: None,
|
||||
}),
|
||||
Some(async_id),
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
let (room, events) = res.unwrap();
|
||||
let session = room.session();
|
||||
|
||||
// Successfully connected to the room
|
||||
let _ = server.send_event(
|
||||
proto::ffi_event::Message::ConnectEvent(proto::ConnectEvent {
|
||||
success: true,
|
||||
room: Some((&session).into()),
|
||||
}),
|
||||
Some(async_id),
|
||||
);
|
||||
|
||||
// Add the room to the server and listen to the incoming events
|
||||
let (close_tx, close_rx) = oneshot::channel();
|
||||
let room_handle = tokio::spawn(room_task(server, room, events, close_rx));
|
||||
server.add_room(session.sid(), (room_handle, close_tx));
|
||||
}
|
||||
|
||||
async fn room_task(
|
||||
server: &'static FFIServer,
|
||||
room: Room,
|
||||
mut events: mpsc::UnboundedReceiver<livekit::RoomEvent>,
|
||||
mut close_rx: oneshot::Receiver<()>,
|
||||
) {
|
||||
let session = room.session();
|
||||
|
||||
tokio::spawn(participant_task(Participant::Local(
|
||||
session.local_participant(),
|
||||
)));
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
Some(event) = events.recv() => {
|
||||
if let Some(event) = proto::RoomEvent::from(session.sid(), event.clone()) {
|
||||
let _ = server.send_event(proto::ffi_event::Message::RoomEvent(event), None);
|
||||
}
|
||||
|
||||
match event {
|
||||
RoomEvent::ParticipantConnected(p) => {
|
||||
tokio::spawn(participant_task(Participant::Remote(p)));
|
||||
}
|
||||
RoomEvent::TrackSubscribed {
|
||||
track,
|
||||
publication: _,
|
||||
participant: _,
|
||||
} => {
|
||||
if let RemoteTrackHandle::Video(video_track) = track {
|
||||
let rtc_track = video_track.rtc_track();
|
||||
rtc_track.on_frame(on_video_frame(server, video_track.sid()));
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
},
|
||||
_ = &mut close_rx => {
|
||||
break;
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
room.close().await;
|
||||
}
|
||||
|
||||
async fn participant_task(participant: Participant) {
|
||||
let mut participant_events = participant.register_observer();
|
||||
while let Some(event) = participant_events.recv().await {
|
||||
// TODO convert event to proto
|
||||
}
|
||||
}
|
||||
|
||||
fn on_video_frame(server: &'static FFIServer, track_sid: TrackSid) -> OnFrameHandler {
|
||||
// TODO(theomonnom): Should I use VideoSinkInfo here? (It'll help to have a more verbose
|
||||
// lifetime)
|
||||
|
||||
Box::new(move |frame, buffer| {
|
||||
// Frame received, create a new FFIHandle from the video buffer.
|
||||
let handle_id = server.next_handle_id();
|
||||
let proto_buffer = proto::VideoFrameBufferInfo::from(handle_id, &buffer);
|
||||
server.insert_handle(handle_id, Box::new(buffer));
|
||||
|
||||
// Send the received frame to the FFI language.
|
||||
let _ = server.send_event(
|
||||
proto::ffi_event::Message::TrackEvent(proto::TrackEvent {
|
||||
track_sid: track_sid.to_string(),
|
||||
message: Some(proto::track_event::Message::FrameReceived(
|
||||
proto::FrameReceived {
|
||||
frame: Some(frame.into()),
|
||||
frame_buffer: Some(proto_buffer),
|
||||
},
|
||||
)),
|
||||
}),
|
||||
None,
|
||||
);
|
||||
})
|
||||
}
|
||||
//
|
||||
Reference in New Issue
Block a user