From 4135ef30cc66026ec33a3b83f9562672f25340a2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Th=C3=A9o=20Monnom?= Date: Wed, 28 Dec 2022 02:19:40 +0100 Subject: [PATCH] DataReceived event --- crates/livekit-core/src/room/mod.rs | 28 +++++++++++++++++++ .../livekit-core/src/room/participant/mod.rs | 2 +- crates/livekit-core/src/room/room_session.rs | 15 ++++++++++ crates/livekit-core/src/rtc_engine/mod.rs | 20 ++++++++++++- .../src/rtc_engine/rtc_session.rs | 12 ++++++-- 5 files changed, 72 insertions(+), 5 deletions(-) diff --git a/crates/livekit-core/src/room/mod.rs b/crates/livekit-core/src/room/mod.rs index aec30ce..8a6b17a 100644 --- a/crates/livekit-core/src/room/mod.rs +++ b/crates/livekit-core/src/room/mod.rs @@ -1,7 +1,10 @@ use self::room_session::{ConnectionState, RoomSession, SessionHandle}; +use crate::proto::data_packet; use crate::room::id::TrackSid; use crate::room::participant::remote_participant::RemoteParticipant; +use crate::room::participant::ParticipantHandle; use crate::room::publication::RemoteTrackPublication; +use crate::room::publication::TrackPublication; use crate::room::track::remote_track::RemoteTrackHandle; use crate::rtc_engine::EngineError; use std::fmt::Debug; @@ -48,11 +51,36 @@ pub enum RoomEvent { publication: RemoteTrackPublication, participant: Arc, }, + TrackUnpublished { + publication: RemoteTrackPublication, + participant: Arc, + }, + TrackUnsubscribed { + track: RemoteTrackHandle, + publication: RemoteTrackPublication, + participant: Arc, + }, TrackSubscriptionFailed { error: TrackError, sid: TrackSid, participant: Arc, }, + TrackMuted { + publication: TrackPublication, + participant: ParticipantHandle, + }, + TrackUnmuted { + publication: TrackPublication, + participant: ParticipantHandle, + }, + ActiveSpeakersChanged { + speakers: Vec, + }, + DataReceived { + payload: Vec, + kind: data_packet::Kind, + participant: Arc, + }, ConnectionStateChanged(ConnectionState), Connected, Disconnected, diff --git a/crates/livekit-core/src/room/participant/mod.rs b/crates/livekit-core/src/room/participant/mod.rs index 5b473ea..911f5d4 100644 --- a/crates/livekit-core/src/room/participant/mod.rs +++ b/crates/livekit-core/src/room/participant/mod.rs @@ -63,7 +63,7 @@ pub trait ParticipantTrait { fn metadata(&self) -> String; } -#[derive(Clone)] +#[derive(Debug, Clone)] pub enum ParticipantHandle { Local(Arc), Remote(Arc), diff --git a/crates/livekit-core/src/room/room_session.rs b/crates/livekit-core/src/room/room_session.rs index 9ecf65e..fecd1da 100644 --- a/crates/livekit-core/src/room/room_session.rs +++ b/crates/livekit-core/src/room/room_session.rs @@ -283,6 +283,21 @@ impl SessionInner { EngineEvent::Restarting => self.handle_restarting(), EngineEvent::Restarted => self.handle_restarted(), EngineEvent::Disconnected => self.handle_disconnected(), + EngineEvent::Data { + payload, + kind, + participant_sid, + } => { + if let Some(participant) = self.get_participant(&participant_sid.into()) { + let _ = self + .internal_tx + .send(SessionEvent::Room(RoomEvent::DataReceived { + payload, + kind, + participant, + })); + } + } } Ok(()) diff --git a/crates/livekit-core/src/rtc_engine/mod.rs b/crates/livekit-core/src/rtc_engine/mod.rs index dbe6b63..ab29a65 100644 --- a/crates/livekit-core/src/rtc_engine/mod.rs +++ b/crates/livekit-core/src/rtc_engine/mod.rs @@ -73,6 +73,11 @@ pub enum EngineEvent { stream: MediaStream, receiver: RtpReceiver, }, + Data { + participant_sid: String, + payload: Vec, + kind: data_packet::Kind, + }, Resuming, Resumed, Restarting, @@ -241,7 +246,20 @@ impl EngineInner { self.close().await; } } - SessionEvent::Data { data: _ } => {} + SessionEvent::Data { + participant_sid, + payload, + kind, + } => { + let _ = self + .engine_emitter + .send(EngineEvent::Data { + participant_sid, + payload, + kind, + }) + .await; + } SessionEvent::MediaTrack { track, stream, diff --git a/crates/livekit-core/src/rtc_engine/rtc_session.rs b/crates/livekit-core/src/rtc_engine/rtc_session.rs index 7e6d847..b375761 100644 --- a/crates/livekit-core/src/rtc_engine/rtc_session.rs +++ b/crates/livekit-core/src/rtc_engine/rtc_session.rs @@ -44,7 +44,9 @@ pub type SessionEvents = mpsc::UnboundedReceiver; #[derive(Debug)] pub enum SessionEvent { Data { - data: Vec, + participant_sid: String, + payload: Vec, + kind: proto::data_packet::Kind, }, MediaTrack { track: MediaStreamTrackHandle, @@ -501,8 +503,12 @@ impl SessionInner { let data = DataPacket::decode(&*data)?; match data.value.unwrap() { - Value::User(_user) => { - // TODO(theomonnom) Send event + Value::User(user) => { + let _ = self.emitter.send(SessionEvent::Data { + participant_sid: user.participant_sid, + payload: user.payload, + kind: data_packet::Kind::from_i32(data.kind).unwrap(), + }); } Value::Speaker(_) => { // TODO(theomonnonm)