diff --git a/crates/livekit-core/src/room.rs b/crates/livekit-core/src/room.rs index b6c7309..583044f 100644 --- a/crates/livekit-core/src/room.rs +++ b/crates/livekit-core/src/room.rs @@ -2,6 +2,7 @@ use std::sync::Arc; use thiserror::Error; use tokio::sync::Mutex; +use tracing::{event, Level, trace}; use crate::local_participant::LocalParticipant; use crate::rtc_engine; @@ -17,7 +18,7 @@ pub struct Room { sid: String, name: String, local_participant: LocalParticipant, - engine: Arc>, + internal: Arc } #[tracing::instrument(skip(url, token))] @@ -25,19 +26,21 @@ pub async fn connect(url: &str, token: &str) -> Result { let engine = rtc_engine::connect(url, token).await?; engine.on_data(Box::new(|packet| { + event!(Level::DEBUG, "received data"); Box::pin(async move {}) })).await; let join = engine.join_response().await; let engine = Arc::new(Mutex::new(engine)); - let lp = LocalParticipant::from(join.participant.unwrap(), engine.clone()); + let local_participant = LocalParticipant::from(join.participant.unwrap(), engine.clone()); + let internal = Arc::new(RoomInternal::new(engine)); let room_info = join.room.unwrap(); Ok(Room { sid: room_info.sid, name: room_info.name, - local_participant: lp, - engine, + local_participant, + internal, }) } @@ -54,3 +57,16 @@ impl Room { &self.name } } + + +struct RoomInternal { + engine: Arc>, +} + +impl RoomInternal { + pub fn new(engine: Arc>) -> Self { + Self { + engine + } + } +} \ No newline at end of file diff --git a/crates/livekit-core/src/rtc_engine.rs b/crates/livekit-core/src/rtc_engine.rs index dfc1054..c8cb166 100644 --- a/crates/livekit-core/src/rtc_engine.rs +++ b/crates/livekit-core/src/rtc_engine.rs @@ -55,6 +55,9 @@ pub struct Packet { pub kind: data_packet::Kind, } +pub type OnDataHandler = +Box Pin + Send + 'static>>) + Send + Sync>; + #[derive(Error, Debug)] pub enum EngineError { #[error("signal failure")] @@ -106,32 +109,6 @@ pub enum EngineMessage { }, } -pub type OnDataHandler = -Box Pin + Send + 'static>>) + Send + Sync>; - -struct EngineInternal { - publisher_pc: Arc>, - subscriber_pc: Arc>, - lossy_dc: Arc>, - reliable_dc: Arc>, - lossy_dc_sub: Arc>>, - reliable_dc_sub: Arc>>, - - msg_sender: mpsc::Sender, - join_response: Mutex, - pc_state: AtomicU8, - // PCState - has_published: AtomicBool, - - // Listeners - on_data_handler: Arc>>, -} - -impl Debug for EngineInternal { - fn fmt(&self, f: &mut Formatter) -> std::fmt::Result { - write!(f, "EngineInternal") - } -} #[derive(Debug)] pub struct RTCEngine { @@ -159,7 +136,7 @@ pub async fn connect(url: &str, token: &str) -> Result { if let Some(signal_response::Message::Join(join_response)) = signal_client.recv().await { event!(Level::DEBUG, "received JoinResponse: {:?}", join_response); let (sender, receiver) = mpsc::channel(8); - let internal = Arc::new(RTCEngine::configure( + let internal = Arc::new(EngineInternal::new( lk_runtime.clone(), sender, join_response.clone(), @@ -174,7 +151,7 @@ pub async fn connect(url: &str, token: &str) -> Result { let internal = internal.clone(); async move { - RTCEngine::handle_loop(receiver, signal_client, internal).await; + internal.run(receiver, signal_client).await; } }); @@ -196,16 +173,15 @@ impl RTCEngine { data: &DataPacket, kind: data_packet::Kind, ) -> Result<(), EngineError> { - self.ensure_publisher_connected(kind).await?; - - self.data_channel(kind) + self.internal.ensure_publisher_connected(kind).await?; + self.internal.data_channel(kind) .lock() .await .send(&data.encode_to_vec(), true) .map_err(Into::into) } - /// Return the last JoinResponse from the server + /// Return the last received JoinResponse pub async fn join_response(&self) -> JoinResponse { self.internal.join_response.lock().await.clone() } @@ -213,335 +189,35 @@ impl RTCEngine { pub async fn on_data(&self, f: OnDataHandler) { *self.internal.on_data_handler.lock().await = Some(f); } +} - fn data_channel(&self, kind: data_packet::Kind) -> Arc> { - if kind == data_packet::Kind::Reliable { - self.internal.reliable_dc.clone() - } else { - self.internal.lossy_dc.clone() - } - } +struct EngineInternal { + publisher_pc: Arc>, + subscriber_pc: Arc>, + lossy_dc: Arc>, + reliable_dc: Arc>, + lossy_dc_sub: Arc>>, + reliable_dc_sub: Arc>>, - #[tracing::instrument] - async fn ensure_publisher_connected( - &mut self, - kind: data_packet::Kind, - ) -> Result<(), EngineError> { - if !self.join_response().await.subscriber_primary { - return Ok(()); - } + msg_sender: mpsc::Sender, + join_response: Mutex, + pc_state: AtomicU8, + // PCState + has_published: AtomicBool, - let publisher = &self.internal.publisher_pc; - { - let mut publisher = publisher.lock().await; - if !publisher.is_connected() - && publisher.peer_connection().ice_connection_state() - != IceConnectionState::IceConnectionChecking - { - tokio::spawn({ - let rtc_internal = self.internal.clone(); - async move { - let _ = Self::negotiate_publisher(rtc_internal).await; - } - }); - } - } + // Listeners + on_data_handler: Arc>>, +} - let dc = self.data_channel(kind); - if dc.lock().await.state() == DataState::Open { - return Ok(()); - } - - let res = time::timeout(MAX_ICE_CONNECT_TIMEOUT, async move { - let mut interval = time::interval(Duration::from_millis(50)); - - loop { - if publisher.lock().await.is_connected() - && dc.lock().await.state() == DataState::Open - { - break; - } - - interval.tick().await; - } - }) - .await; - - if res.is_err() { - let err = - EngineError::Connection("could not establish publisher connection".to_string()); - event!(Level::ERROR, error = ?err); - Err(err) - } else { - Ok(()) - } - } - - #[tracing::instrument] - async fn handle_signal( - signal: signal_response::Message, - signal_client: Arc, - rtc_internal: Arc, - ) -> Result<(), EngineError> { - match signal { - signal_response::Message::Answer(answer) => { - event!(Level::TRACE, "received answer for publisher: {:?}", answer); - let sdp = SessionDescription::from(answer.r#type.parse().unwrap(), &answer.sdp)?; - rtc_internal - .publisher_pc - .lock() - .await - .set_remote_description(sdp) - .await?; - } - signal_response::Message::Offer(offer) => { - event!(Level::TRACE, "received offer for subscriber: {:?}", offer); - let sdp = SessionDescription::from(offer.r#type.parse().unwrap(), &offer.sdp)?; - let subscriber_pc = &rtc_internal.subscriber_pc; - - subscriber_pc - .lock() - .await - .set_remote_description(sdp) - .await?; - let answer = subscriber_pc - .lock() - .await - .peer_connection() - .create_answer(RTCOfferAnswerOptions::default()) - .await?; - subscriber_pc - .lock() - .await - .peer_connection() - .set_local_description(answer.clone()) - .await?; - - tokio::spawn(async move { - let _ = signal_client.send(signal_request::Message::Answer( - proto::SessionDescription { - r#type: "answer".to_string(), - sdp: answer.to_string(), - }, - )).await; - }); - } - signal_response::Message::Trickle(trickle) => { - let json: IceCandidateJSON = serde_json::from_str(&trickle.candidate_init)?; - let ice = IceCandidate::from(&json.sdpMid, json.sdpMLineIndex, &json.candidate)?; - - event!( - Level::TRACE, - "received ice_candidate ({:?}) - {:?}", - SignalTarget::from_i32(trickle.target).unwrap(), - ice - ); - - if trickle.target == SignalTarget::Publisher as i32 { - rtc_internal - .publisher_pc - .lock() - .await - .add_ice_candidate(ice) - .await?; - } else { - rtc_internal - .subscriber_pc - .lock() - .await - .add_ice_candidate(ice) - .await?; - } - } - _ => {} - } - - Ok(()) - } - - #[tracing::instrument] - async fn handle_message( - msg: EngineMessage, - signal_client: Arc, - rtc_internal: Arc, - ) -> Result<(), EngineError> { - match msg { - EngineMessage::IceCandidate { - ice_candidate, - publisher, - } => { - let target = if publisher { - SignalTarget::Publisher - } else { - SignalTarget::Subscriber - }; - - event!( - Level::TRACE, - "sending ice_candidate ({:?}) - {:?}", - target, - ice_candidate - ); - - let json = serde_json::to_string(&IceCandidateJSON { - sdpMid: ice_candidate.sdp_mid(), - sdpMLineIndex: ice_candidate.sdp_mline_index(), - candidate: ice_candidate.candidate(), - })?; - - // Send the ice_candidate to the server - tokio::spawn(async move { - let _ = signal_client.send(signal_request::Message::Trickle( - TrickleRequest { - candidate_init: json, - target: target as i32, - }, - )).await; - }); - } - EngineMessage::ConnectionChange { state, primary } => { - if primary && state == PeerConnectionState::Connected { - let old_state = rtc_internal.pc_state.load(Ordering::SeqCst); - rtc_internal - .pc_state - .store(PCState::Connected as u8, Ordering::SeqCst); - - if old_state == PCState::New as u8 { - // TODO(theomonnom) OnConnected - } - } else if state == PeerConnectionState::Failed { - rtc_internal - .pc_state - .store(PCState::Disconnected as u8, Ordering::SeqCst); - - // TODO(theomonnom) handle Disconnect - } - } - EngineMessage::PrimaryDataChannel { mut data_channel } => { - let reliable = data_channel.label() == RELIABLE_DC_LABEL; - Self::configure_dc(&mut data_channel, rtc_internal.msg_sender.clone()); - - event!( - Level::TRACE, - "received subscriber data_channel - {:?}", - data_channel - ); - - if reliable { - *rtc_internal.reliable_dc_sub.lock().await = Some(data_channel); - } else { - *rtc_internal.lossy_dc_sub.lock().await = Some(data_channel); - } - } - EngineMessage::PublisherOffer { offer } => { - event!( - Level::TRACE, - "sending publisher offer - {:?}", - offer - ); - - // Send the offer to the server - tokio::spawn(async move { - let _ = signal_client.send(signal_request::Message::Offer( - proto::SessionDescription { - r#type: "offer".to_string(), - sdp: offer.to_string(), - }, - )).await; - }); - } - EngineMessage::Data { - data, - binary, - } => { - if !binary { - return Err(EngineError::Internal( - "text messages aren't supported by LiveKit".to_string(), - )); - } - - let data = DataPacket::decode(&*data)?; - match data.value.unwrap() { - Value::User(user) => { - let mut handler = rtc_internal.on_data_handler.lock().await; - if let Some(f) = &mut *handler { - f(Packet { - data: user, - kind: data_packet::Kind::from_i32(data.kind).unwrap(), - }) - .await; - } - } - Value::Speaker(_) => { - // TODO(theomonnonm) - } - } - } - } - - Ok(()) - } - - #[tracing::instrument] - async fn handle_loop( - mut receiver: mpsc::Receiver, - signal_client: Arc, - rtc_internal: Arc, - ) { - loop { - tokio::select! { - signal = signal_client.recv() => { - match signal { - Some(signal) => { - if let Err(err) = Self::handle_signal(signal, signal_client.clone(), rtc_internal.clone()).await { - event!( - Level::ERROR, - "failed to handle signal: {:?}", - err, - ); - } - } - None => { - // TODO(theomonnom) Trigger reconnect - } - } - }, - Some(msg) = receiver.recv() => { - if let Err(err) = Self::handle_message(msg, signal_client.clone(), rtc_internal.clone()).await { - event!( - Level::ERROR, - "failed to handle engine message: {:?}", - err, - ); - } - }, - } - } - } - - #[tracing::instrument] - async fn negotiate_publisher(rtc_internal: Arc) -> Result<(), EngineError> { - rtc_internal.has_published.store(true, Ordering::SeqCst); - if let Err(err) = rtc_internal.publisher_pc.lock().await.negotiate().await { - event!( - Level::ERROR, - "failed to negotiate the publisher: {:?}", - err, - ); - Err(EngineError::Rtc(err)) - } else { - Ok(()) - } - } - - /// This function is called on connect & on reconnect +impl EngineInternal { + /// New internal is created on connect & on reconnect /// It creates the PeerConnections, the DataChannels and the libwebrtc listeners #[tracing::instrument] - fn configure( + fn new( lk_runtime: Arc, sender: mpsc::Sender, join: JoinResponse, - ) -> Result { + ) -> Result { let rtc_config = RTCConfiguration { ice_servers: { let mut servers = vec![]; @@ -658,7 +334,7 @@ impl RTCEngine { Self::configure_dc(&mut lossy_dc, sender.clone()); Self::configure_dc(&mut reliable_dc, sender.clone()); - Ok(EngineInternal { + Ok(Self { publisher_pc: Arc::new(Mutex::new(publisher_pc)), subscriber_pc: Arc::new(Mutex::new(subscriber_pc)), lossy_dc: Arc::new(Mutex::new(lossy_dc)), @@ -686,4 +362,323 @@ impl RTCEngine { }); })); } + + #[tracing::instrument] + async fn ensure_publisher_connected( + self: &Arc, + kind: data_packet::Kind, + ) -> Result<(), EngineError> { + if !self.join_response.lock().await.subscriber_primary { + return Ok(()); + } + + let publisher = &self.publisher_pc; + { + let mut publisher = publisher.lock().await; + if !publisher.is_connected() + && publisher.peer_connection().ice_connection_state() + != IceConnectionState::IceConnectionChecking + { + tokio::spawn({ + let internal = self.clone(); + async move { + let _ = internal.negotiate_publisher().await; + } + }); + } + } + + let dc = self.data_channel(kind); + if dc.lock().await.state() == DataState::Open { + return Ok(()); + } + + let res = time::timeout(MAX_ICE_CONNECT_TIMEOUT, async move { + let mut interval = time::interval(Duration::from_millis(50)); + + loop { + if publisher.lock().await.is_connected() + && dc.lock().await.state() == DataState::Open + { + break; + } + + interval.tick().await; + } + }) + .await; + + if res.is_err() { + let err = + EngineError::Connection("could not establish publisher connection".to_string()); + event!(Level::ERROR, error = ?err); + Err(err) + } else { + Ok(()) + } + } + + #[tracing::instrument] + async fn handle_signal( + self: &Arc, + signal: signal_response::Message, + signal_client: Arc, + ) -> Result<(), EngineError> { + match signal { + signal_response::Message::Answer(answer) => { + event!(Level::TRACE, "received answer for publisher: {:?}", answer); + let sdp = SessionDescription::from(answer.r#type.parse().unwrap(), &answer.sdp)?; + self.publisher_pc + .lock() + .await + .set_remote_description(sdp) + .await?; + } + signal_response::Message::Offer(offer) => { + event!(Level::TRACE, "received offer for subscriber: {:?}", offer); + let sdp = SessionDescription::from(offer.r#type.parse().unwrap(), &offer.sdp)?; + + self.subscriber_pc + .lock() + .await + .set_remote_description(sdp) + .await?; + let answer = self.subscriber_pc + .lock() + .await + .peer_connection() + .create_answer(RTCOfferAnswerOptions::default()) + .await?; + self.subscriber_pc + .lock() + .await + .peer_connection() + .set_local_description(answer.clone()) + .await?; + + tokio::spawn(async move { + let _ = signal_client.send(signal_request::Message::Answer( + proto::SessionDescription { + r#type: "answer".to_string(), + sdp: answer.to_string(), + }, + )).await; + }); + } + signal_response::Message::Trickle(trickle) => { + let json: IceCandidateJSON = serde_json::from_str(&trickle.candidate_init)?; + let ice = IceCandidate::from(&json.sdpMid, json.sdpMLineIndex, &json.candidate)?; + + event!( + Level::TRACE, + "received ice_candidate ({:?}) - {:?}", + SignalTarget::from_i32(trickle.target).unwrap(), + ice + ); + + if trickle.target == SignalTarget::Publisher as i32 { + self.publisher_pc + .lock() + .await + .add_ice_candidate(ice) + .await?; + } else { + self.subscriber_pc + .lock() + .await + .add_ice_candidate(ice) + .await?; + } + } + _ => {} + } + + Ok(()) + } + + #[tracing::instrument] + pub async fn run( + self: &Arc, + mut receiver: mpsc::Receiver, + signal_client: Arc, + ) { + loop { + tokio::select! { + signal = signal_client.recv() => { + match signal { + Some(signal) => { + if let Err(err) = self.handle_signal(signal, signal_client.clone()).await { + event!( + Level::ERROR, + "failed to handle signal: {:?}", + err, + ); + } + } + None => { + // TODO(theomonnom) Trigger reconnect + } + } + }, + Some(msg) = receiver.recv() => { + if let Err(err) = self.handle_message(msg, signal_client.clone()).await { + event!( + Level::ERROR, + "failed to handle engine message: {:?}", + err, + ); + } + }, + } + } + } + + #[tracing::instrument] + async fn handle_message( + self: &Arc, + msg: EngineMessage, + signal_client: Arc, + ) -> Result<(), EngineError> { + match msg { + EngineMessage::IceCandidate { + ice_candidate, + publisher, + } => { + let target = if publisher { + SignalTarget::Publisher + } else { + SignalTarget::Subscriber + }; + + event!( + Level::TRACE, + "sending ice_candidate ({:?}) - {:?}", + target, + ice_candidate + ); + + let json = serde_json::to_string(&IceCandidateJSON { + sdpMid: ice_candidate.sdp_mid(), + sdpMLineIndex: ice_candidate.sdp_mline_index(), + candidate: ice_candidate.candidate(), + })?; + + // Send the ice_candidate to the server + tokio::spawn(async move { + let _ = signal_client.send(signal_request::Message::Trickle( + TrickleRequest { + candidate_init: json, + target: target as i32, + }, + )).await; + }); + } + EngineMessage::ConnectionChange { state, primary } => { + if primary && state == PeerConnectionState::Connected { + let old_state = self.pc_state.load(Ordering::SeqCst); + self.pc_state + .store(PCState::Connected as u8, Ordering::SeqCst); + + if old_state == PCState::New as u8 { + // TODO(theomonnom) OnConnected + } + } else if state == PeerConnectionState::Failed { + self.pc_state.store(PCState::Disconnected as u8, Ordering::SeqCst); + + // TODO(theomonnom) handle Disconnect + } + } + EngineMessage::PrimaryDataChannel { mut data_channel } => { + let reliable = data_channel.label() == RELIABLE_DC_LABEL; + Self::configure_dc(&mut data_channel, self.msg_sender.clone()); + + event!( + Level::TRACE, + "received subscriber data_channel - {:?}", + data_channel + ); + + if reliable { + *self.reliable_dc_sub.lock().await = Some(data_channel); + } else { + *self.lossy_dc_sub.lock().await = Some(data_channel); + } + } + EngineMessage::PublisherOffer { offer } => { + event!( + Level::TRACE, + "sending publisher offer - {:?}", + offer + ); + + // Send the offer to the server + tokio::spawn(async move { + let _ = signal_client.send(signal_request::Message::Offer( + proto::SessionDescription { + r#type: "offer".to_string(), + sdp: offer.to_string(), + }, + )).await; + }); + } + EngineMessage::Data { + data, + binary, + } => { + if !binary { + return Err(EngineError::Internal( + "text messages aren't supported by LiveKit".to_string(), + )); + } + + let data = DataPacket::decode(&*data)?; + match data.value.unwrap() { + Value::User(user) => { + let mut handler = self.on_data_handler.lock().await; + if let Some(f) = &mut *handler { + f(Packet { + data: user, + kind: data_packet::Kind::from_i32(data.kind).unwrap(), + }) + .await; + } + } + Value::Speaker(_) => { + // TODO(theomonnonm) + } + } + } + } + + Ok(()) + } + + #[tracing::instrument] + async fn negotiate_publisher(self: &Arc) -> Result<(), EngineError> { + self.has_published.store(true, Ordering::SeqCst); + if let Err(err) = self.publisher_pc.lock().await.negotiate().await { + event!( + Level::ERROR, + "failed to negotiate the publisher: {:?}", + err, + ); + Err(EngineError::Rtc(err)) + } else { + Ok(()) + } + } + + fn data_channel(&self, kind: data_packet::Kind) -> Arc> { + if kind == data_packet::Kind::Reliable { + self.reliable_dc.clone() + } else { + self.lossy_dc.clone() + } + } +} + +impl Debug for EngineInternal { + fn fmt(&self, f: &mut Formatter) -> std::fmt::Result { + write!(f, "EngineInternal") + } } diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp index 985d8a3..41ecf56 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/peer_connection_factory.cpp @@ -30,6 +30,7 @@ PeerConnectionFactory::PeerConnectionFactory( dependencies.task_queue_factory = webrtc::CreateDefaultTaskQueueFactory(); dependencies.event_log_factory = std::make_unique( dependencies.task_queue_factory.get()); + dependencies.call_factory = webrtc::CreateCallFactory(); cricket::MediaEngineDependencies media_deps; media_deps.task_queue_factory = dependencies.task_queue_factory.get(); diff --git a/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp b/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp index bee2fc1..0afc9d9 100644 --- a/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp +++ b/crates/livekit-webrtc/libwebrtc-sys/src/webrtc.cpp @@ -8,7 +8,7 @@ namespace livekit { RTCRuntime::RTCRuntime() { - // rtc::LogMessage::LogToDebug(rtc::LS_INFO); + rtc::LogMessage::LogToDebug(rtc::LS_INFO); RTC_LOG(LS_INFO) << "RTCRuntime()"; RTC_CHECK(rtc::InitializeSSL()) << "Failed to InitializeSSL()"; diff --git a/examples/simple_room/src/main.rs b/examples/simple_room/src/main.rs index d22fecc..1c09fbe 100644 --- a/examples/simple_room/src/main.rs +++ b/examples/simple_room/src/main.rs @@ -1,10 +1,12 @@ +use std::time::Duration; +use tokio::time::sleep; use livekit::proto::data_packet; use livekit::room; const URL: &str = "ws://localhost:7880"; -const TOKEN : &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE2NjgxMzc0NDgsImlzcyI6IkFQSXpLYkFTaUNWYWtnSiIsIm5hbWUiOiJ3ZWIiLCJuYmYiOjE2NjQ1Mzc0NDgsInN1YiI6IndlYiIsInZpZGVvIjp7InJvb21DcmVhdGUiOnRydWUsInJvb21Kb2luIjp0cnVlfX0.6VMDdXJYrW3KWrEzxx4hzbmMQnjQIRILQ48Qrbx5j44"; +const TOKEN : &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjIzODQ4MDY0NzMsImlzcyI6IkFQSXpLYkFTaUNWYWtnSiIsIm5hbWUiOiJuYXRpdmUiLCJuYmYiOjE2NjQ4MDY0NzMsInN1YiI6Im5hdGl2ZSIsInZpZGVvIjp7InJvb21DcmVhdGUiOnRydWUsInJvb21Kb2luIjp0cnVlfX0.BgVdBnq3XFD3_BQHoe1azqjifYysubgFl6Qlzu9IQGI"; -// eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjE2NjgxMzc0NDgsImlzcyI6IkFQSXpLYkFTaUNWYWtnSiIsIm5hbWUiOiJ3ZWIiLCJuYmYiOjE2NjQ1Mzc0NDgsInN1YiI6IndlYiIsInZpZGVvIjp7InJvb21DcmVhdGUiOnRydWUsInJvb21Kb2luIjp0cnVlfX0.6VMDdXJYrW3KWrEzxx4hzbmMQnjQIRILQ48Qrbx5j44 +// eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJleHAiOjIzODQ4MDY3MzAsImlzcyI6IkFQSXpLYkFTaUNWYWtnSiIsIm5hbWUiOiJ3ZWIiLCJuYmYiOjE2NjQ4MDY3MzAsInN1YiI6IndlYiIsInZpZGVvIjp7InJvb21DcmVhdGUiOnRydWUsInJvb21Kb2luIjp0cnVlfX0.VbDoULjX1CVGZu2sPy3SvWYlVZUBXxQVPmdB9BnmlN4 #[tokio::main] async fn main() -> Result<(), room::RoomError> { @@ -12,8 +14,9 @@ async fn main() -> Result<(), room::RoomError> { let mut room = room::connect(URL, TOKEN).await?; room.local_participant() - .publish_data(b"this is a test", data_packet::Kind::Reliable) + .publish_data(b"some data", data_packet::Kind::Reliable) .await?; + sleep(Duration::from_secs(120)).await; Ok(()) }