From b234359c4a7e91b05f8cf71d32d29a36517264d4 Mon Sep 17 00:00:00 2001 From: Nell Date: Wed, 23 Sep 2026 13:58:26 +0200 Subject: [PATCH] add webrtc --- frontend/src/pages/rtc-test.vue | 360 ++++++++++++++++++++ frontend/src/router/index.ts | 5 + frontend/vite.config.mts | 4 + src/rtc/client.rs | 567 ++++++++++++++++++++++++++++++-- src/rtc/messages.rs | 15 +- src/rtc/mod.rs | 126 ++++++- src/rtc/test.html | 112 +++++++ src/rtc/ws_entrypoint.rs | 22 +- 8 files changed, 1163 insertions(+), 48 deletions(-) create mode 100644 frontend/src/pages/rtc-test.vue create mode 100644 src/rtc/test.html diff --git a/frontend/src/pages/rtc-test.vue b/frontend/src/pages/rtc-test.vue new file mode 100644 index 0000000..d342cc1 --- /dev/null +++ b/frontend/src/pages/rtc-test.vue @@ -0,0 +1,360 @@ + + + + + \ No newline at end of file diff --git a/frontend/src/router/index.ts b/frontend/src/router/index.ts index 0b73a87..83bd2e1 100644 --- a/frontend/src/router/index.ts +++ b/frontend/src/router/index.ts @@ -60,6 +60,11 @@ const router = createRouter({ }, ], }, + { + path: 'rtc-test', + name: 'rtc-test', + component: () => import('@/pages/rtc-test.vue'), + }, { path: 'server/:serverId(default|[0-9a-fA-F-]{36})', name: 'server-dashboard', diff --git a/frontend/vite.config.mts b/frontend/vite.config.mts index 0ea712a..62e754e 100644 --- a/frontend/vite.config.mts +++ b/frontend/vite.config.mts @@ -51,6 +51,10 @@ export default defineConfig({ target: 'http://localhost:8080', changeOrigin: true, }, + '/ws': { + target: 'ws://localhost:8080', + ws: true, + }, }, }, }) diff --git a/src/rtc/client.rs b/src/rtc/client.rs index 245fcbf..9d2a50e 100644 --- a/src/rtc/client.rs +++ b/src/rtc/client.rs @@ -1,9 +1,18 @@ -use super::messages::VoiceClientMessage; +use super::messages::{VoiceClientMessage, VoiceServerMessage}; +use super::{RoomEvent, VoiceRoomManager}; use crate::models::{channel, user}; use axum::extract::ws::Message; -use rustrtc::peer_connection::PeerConnection; +use rustrtc::media::{self, MediaStreamTrack}; +use rustrtc::peer_connection::{PeerConnection, RtpCodecParameters}; +use rustrtc::sdp::{SdpType, SessionDescription}; +use rustrtc::transports::ice::IceCandidate; +use rustrtc::{MediaKind, TransceiverDirection}; +use std::collections::HashMap; use std::sync::Arc; use tokio::sync::mpsc::UnboundedSender; +use tokio::sync::mpsc; +use parking_lot::Mutex; +use tracing::{info, warn}; /// Client connecté à un canal RTC. #[derive(Clone)] @@ -12,6 +21,12 @@ pub struct RTCClient { pub channel: Arc, pub peer_connection: Arc, pub websocket_sender: UnboundedSender, + pub room_manager: Arc, + pub client_id: uuid::Uuid, + room_task: Arc>>>, + negotiation_done: mpsc::UnboundedSender<()>, + negotiation_rx: Arc>>>, + pending_room: Arc, u8, String)>>>, } impl RTCClient { @@ -20,41 +35,551 @@ impl RTCClient { channel: channel::Model, peer_connection: Arc, websocket_sender: UnboundedSender, + room_manager: Arc, ) -> Self { + let (negotiation_done, negotiation_rx) = mpsc::unbounded_channel(); Self { user: Arc::new(user), channel: Arc::new(channel), peer_connection, websocket_sender, + room_manager, + client_id: uuid::Uuid::new_v4(), + room_task: Arc::new(Mutex::new(None)), + negotiation_done, + negotiation_rx: Arc::new(Mutex::new(Some(negotiation_rx))), + pending_room: Arc::new(Mutex::new(None)), } } - pub async fn ws_on_message(&self, raw_message: Message) { - let Message::Text(raw_message) = raw_message else { - return; - }; - - let parsed = match serde_json::from_str::(&raw_message) { - Ok(message) => message, - Err(_) => return, - }; - - match parsed { - VoiceClientMessage::SDPOffer { channel_id, sdp } => {} - VoiceClientMessage::IceCandidate { - channel_id, - candidate, - } => {} - VoiceClientMessage::Leave { channel_id } => {} + /// Retourne `false` quand le client demande à quitter le canal. + pub async fn ws_on_message(&self, raw_message: Message) -> bool { + match raw_message { + Message::Close(_) => false, + Message::Text(text) => match serde_json::from_str::(&text) { + Ok(VoiceClientMessage::Leave) => false, + Ok(VoiceClientMessage::SdpAnswer { sdp }) => { + match SessionDescription::parse(SdpType::Answer, &sdp) { + Ok(answer) => match self.peer_connection.set_remote_description(answer).await { + Ok(()) => { let _ = self.negotiation_done.send(()); } + Err(error) => warn!(%error, "Réponse de renégociation refusée"), + }, + Err(error) => warn!(%error, "Réponse SDP invalide"), + } + true + } + Ok(VoiceClientMessage::SDPOffer { sdp }) => { + info!(channel_id = %self.channel.id, "Offre SDP reçue"); + let response = match self.handle_sdp_offer(&sdp).await { + Ok(sdp) => { + info!(channel_id = %self.channel.id, "Réponse SDP créée"); + VoiceServerMessage::Answer { sdp } + } + Err(error) => { + warn!(channel_id = %self.channel.id, %error, "Négociation SDP échouée"); + VoiceServerMessage::Error { message: error.to_string() } + } + }; + self.send_response(response); + self.start_room(); + true + } + Ok(VoiceClientMessage::IceCandidate { candidate }) => { + if candidate.is_empty() { + info!(channel_id = %self.channel.id, "Fin de collecte ICE distante"); + return true; + } + if let Err(error) = self.handle_ice_candidate(&candidate) { + if candidate.split_whitespace().nth(4).is_some_and(|address| address.ends_with(".local")) { + warn!(channel_id = %self.channel.id, %error, candidate = %candidate, "Candidat ICE mDNS du navigateur refusé par rustrtc"); + } else { + warn!(channel_id = %self.channel.id, %error, candidate = %candidate, "Candidat ICE distant refusé"); + } + self.send_response(VoiceServerMessage::Error { + message: error.to_string(), + }); + } else { + info!(channel_id = %self.channel.id, "Candidat ICE distant accepté"); + } + true + } + _ => true, + }, + _ => true, } } - pub async fn ws_send_message(&self, raw_message: Message) { - let _ = self.websocket_sender.send(raw_message); + async fn handle_sdp_offer(&self, sdp: &str) -> anyhow::Result { + let opus_payload_type = sdp.split("m=").find(|section| section.starts_with("audio ")) + .and_then(|audio| { + let formats = audio.lines().next()?.split_whitespace().skip(3).collect::>(); + audio.lines().find_map(|line| { + let (payload_type, codec) = line.strip_prefix("a=rtpmap:")?.split_once(' ')?; + (codec.eq_ignore_ascii_case("opus/48000/2") && formats.contains(&payload_type)) + .then(|| payload_type.parse::().ok()).flatten() + }) + }); + let offer = SessionDescription::parse(SdpType::Offer, sdp)?; + self.peer_connection.set_remote_description(offer).await?; + if let Some(transceiver) = self.peer_connection.get_transceivers().into_iter() + .find(|transceiver| transceiver.kind() == MediaKind::Audio) + { + if let Some(receiver) = transceiver.receiver() { + let incoming = receiver.track(); + let payload_type = opus_payload_type.ok_or_else(|| anyhow::anyhow!("Opus absent de l'offre audio"))?; + // Réserver la m-line du microphone : add_track réutiliserait sinon + // cette m-line pour le premier autre participant. + let (_, reserved, _) = media::sample_track(media::MediaKind::Audio, 48000); + self.peer_connection.add_track(reserved, RtpCodecParameters { + payload_type, name: "opus".to_string(), clock_rate: 48000, channels: 2, + })?; + info!(channel_id = %self.channel.id, "Audio du salon configuré"); + let channel_id = self.channel.id; + let client_id = self.client_id; + let room_events = self.room_manager.join(channel_id, client_id); + let room_manager = self.room_manager.clone(); + tokio::spawn(async move { Self::forward_audio(incoming, room_manager, channel_id, client_id).await }); + let answer = self.create_initial_answer(opus_payload_type).await?; + *self.pending_room.lock() = Some((room_events, payload_type, answer.clone())); + return Ok(answer); + } + } + self.create_initial_answer(opus_payload_type).await + } + + fn start_room(&self) { + if let Some((events, payload_type, initial_sdp)) = self.pending_room.lock().take() { + let client = self.clone(); + *self.room_task.lock() = Some(tokio::spawn(async move { client.run_room(events, payload_type, initial_sdp).await })); + } + } + + async fn create_initial_answer(&self, opus_payload_type: Option) -> anyhow::Result { + let answer = self.peer_connection.create_answer().await?; + let sdp = Self::opus_sdp(&answer.to_sdp_string(), opus_payload_type.unwrap_or(111)); + self.peer_connection.set_local_description(SessionDescription::parse(SdpType::Answer, &sdp)?)?; + Ok(sdp) + } + + fn opus_sdp(sdp: &str, payload_type: u8) -> String { + if payload_type == 111 { return sdp.to_string(); } + let mut in_audio = false; + sdp.split_inclusive('\n').map(|line| { + if line.starts_with("m=") { in_audio = line.starts_with("m=audio "); } + if !in_audio { return line.to_string(); } + let ending = if line.ends_with("\r\n") { "\r\n" } else { "\n" }; + let content = line.trim_end_matches(['\r', '\n']); + let adjusted = if content.starts_with("m=audio ") { + content.split_whitespace().map(|part| if part == "111" { payload_type.to_string() } else { part.to_string() }) + .collect::>().join(" ") + } else if let Some(rest) = content.strip_prefix("a=rtpmap:111 ") { + format!("a=rtpmap:{payload_type} {rest}") + } else if let Some(rest) = content.strip_prefix("a=fmtp:111 ") { + format!("a=fmtp:{payload_type} {rest}") + } else { content.to_string() }; + format!("{adjusted}{ending}") + }).collect() + } + + fn stable_extmaps(sdp: &str, initial_sdp: &str) -> String { + let mut ids = HashMap::::new(); + let mut used = std::collections::HashSet::new(); + for line in initial_sdp.lines() { + if let Some((id, uri)) = line.strip_prefix("a=extmap:") + .and_then(|value| value.split_once(' ')) + .and_then(|(id, rest)| Some((id.split('/').next()?.parse::().ok()?, rest.split_whitespace().next()?))) { + ids.entry(uri.to_string()).or_insert(id); + used.insert(id); + } + } + sdp.split_inclusive('\n').map(|line| { + let Some((_, rest)) = line.strip_prefix("a=extmap:").and_then(|value| value.split_once(' ')) else { + return line.to_string(); + }; + let Some(uri) = rest.split_whitespace().next() else { return line.to_string() }; + let id = *ids.entry(uri.to_string()).or_insert_with(|| { + let free = (1..=14).chain(4096..=4351).find(|id| !used.contains(id)).expect("identifiants RTP épuisés"); + used.insert(free); + free + }); + let ending = if line.ends_with("\r\n") { "\r\n" } else { "\n" }; + format!("a=extmap:{id} {}{ending}", rest.trim_end_matches(['\r', '\n'])) + }).collect() + } + + async fn run_room(&self, mut events: mpsc::UnboundedReceiver, payload_type: u8, initial_sdp: String) { + let mut negotiation_rx = self.negotiation_rx.lock().take().expect("une seule tâche de salon"); + let mut tracks: HashMap, tokio::task::JoinHandle<()>)> = HashMap::new(); + while let Some(event) = events.recv().await { + match event { + RoomEvent::Joined(id, mut receiver) => { + let (source, outgoing, _) = media::sample_track(media::MediaKind::Audio, 48000); + let sender = match self.peer_connection.add_track_with_stream_id(outgoing, id.to_string(), RtpCodecParameters { + payload_type, name: "opus".to_string(), clock_rate: 48000, channels: 2, + }) { + Ok(sender) => sender, + Err(error) => { warn!(%error, "Ajout de la piste du salon impossible"); continue; } + }; + let Some(transceiver) = self.peer_connection.get_transceivers().into_iter().find(|t| t.sender().is_some_and(|s| s.ssrc() == sender.ssrc())) else { continue }; + transceiver.set_direction(TransceiverDirection::SendRecv); + let task = tokio::spawn(async move { + loop { + match receiver.recv().await { + Ok(mut sample) => { + if let media::MediaSample::Audio(frame) = &mut sample { + frame.payload_type = None; + frame.sequence_number = None; + frame.raw_packet = None; + } + if source.send(sample).is_err() { break; } + } + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue, + Err(_) => break, + } + } + }); + tracks.insert(id, (transceiver, task)); + } + RoomEvent::Left(id) => { + if let Some((transceiver, task)) = tracks.remove(&id) { + task.abort(); + // Conserver l'émetteur pour éviter que rustrtc réutilise cette + // m-line pour un nouveau participant après son départ. + transceiver.set_direction(TransceiverDirection::Inactive); + } + self.send_response(VoiceServerMessage::SourceLeft { id }); + } + } + match self.peer_connection.create_offer().await { + Ok(offer) => { + let sdp = Self::stable_extmaps(&Self::opus_sdp(&offer.to_sdp_string(), payload_type), &initial_sdp); + if let Err(error) = SessionDescription::parse(SdpType::Offer, &sdp) + .map_err(anyhow::Error::from) + .and_then(|offer| self.peer_connection.set_local_description(offer).map_err(anyhow::Error::from)) { + warn!(%error, "Offre du salon impossible"); break; + } + self.send_response(VoiceServerMessage::SdpOffer { sdp }); + if tokio::time::timeout(std::time::Duration::from_secs(10), negotiation_rx.recv()).await.ok().flatten().is_none() { + warn!("Renégociation du salon expirée"); break; + } + } + Err(error) => { warn!(%error, "Renégociation du salon impossible"); break; } + } + } + for (_, (_, task)) in tracks { task.abort(); } + } + + async fn forward_audio( + incoming: Arc, + room_manager: Arc, + channel_id: uuid::Uuid, + client_id: uuid::Uuid, + ) { + while let Ok(sample) = incoming.recv().await { + room_manager.forward(channel_id, client_id, sample); + } + } + + fn handle_ice_candidate(&self, candidate: &str) -> anyhow::Result<()> { + self.peer_connection + .add_ice_candidate(IceCandidate::from_sdp(candidate.strip_prefix("candidate:").unwrap_or(candidate))?)?; + Ok(()) + } + + /// Transmet au WebSocket les candidats ICE locaux produits par la connexion WebRTC. + /// À démarrer avant l'offre SDP pour ne pas manquer les premiers candidats ; + /// la tâche retournée doit être arrêtée à la déconnexion. + pub fn forward_ice_candidates(&self) -> tokio::task::JoinHandle<()> { + // Chaque connexion WebSocket a son propre abonnement aux candidats locaux. + let mut candidates = self.peer_connection.subscribe_ice_candidates(); + // La tâche doit conserver le client pendant qu'elle attend les candidats. + let client = self.clone(); + tokio::spawn(async move { + // recv attend un candidat ; une erreur termine l'écoute. + while let Ok(candidate) = candidates.recv().await { + // Le candidat est sérialisé pour être envoyé au navigateur via send_response. + info!(channel_id = %client.channel.id, "Candidat ICE local envoyé"); + client.send_response(VoiceServerMessage::IceCandidate { + candidate: format!("candidate:{}", candidate.to_sdp()), + }); + } + }) + } + + fn send_response(&self, response: VoiceServerMessage) { + if let Ok(json) = serde_json::to_string(&response) { + let _ = self.websocket_sender.send(Message::Text(json.into())); + } } /// Ferme proprement la connexion WebRTC. pub fn close(&self) { + info!(channel_id = %self.channel.id, "Fermeture de la connexion RTC"); + self.room_manager.leave(self.channel.id, self.client_id); + if let Some(task) = self.room_task.lock().take() { task.abort(); } self.peer_connection.close(); } } + +#[cfg(test)] +mod tests { + use super::RTCClient; + use crate::models::{channel, user}; + use axum::extract::ws::Message; + use chrono::Utc; + use rustrtc::peer_connection::TransceiverDirection; + use rustrtc::sdp::{SdpType, SessionDescription}; + use rustrtc::{MediaKind, PeerConnection, RtcConfigurationBuilder}; + use std::sync::Arc; + use tokio::sync::mpsc; + use uuid::Uuid; + + fn client() -> (RTCClient, mpsc::UnboundedReceiver) { + let now = Utc::now(); + let (sender, receiver) = mpsc::unbounded_channel(); + ( + RTCClient::new( + user::Model { + id: Uuid::new_v4(), + username: "test".into(), + password: String::new(), + pub_key: None, + created_at: now, + updated_at: now, + is_superuser: false, + }, + channel::Model { + id: Uuid::new_v4(), + server_id: None, + category_id: None, + channel_type: channel::ChannelType::Voice, + name: None, + created_at: now, + updated_at: now, + }, + Arc::new(PeerConnection::new( + RtcConfigurationBuilder::new().build(), + )), + sender, + Arc::new(super::super::VoiceRoomManager::new()), + ), + receiver, + ) + } + + #[tokio::test] + async fn leave_and_close_end_connection() { + let (client, _) = client(); + let leave = r#"{"action":"leave"}"#; + assert!(!client.ws_on_message(Message::Text(leave.into())).await); + assert!(!client.ws_on_message(Message::Close(None)).await); + } + + #[tokio::test] + async fn other_messages_keep_connection_open() { + let (client, mut receiver) = client(); + assert!(client.ws_on_message(Message::Text("invalid".into())).await); + assert!(client.ws_on_message(Message::Ping(vec![].into())).await); + let offer = r#"{"action":"sdp-offer","sdp":"test"}"#; + assert!(serde_json::from_str::(&offer).is_ok()); + assert!(client.ws_on_message(Message::Text(offer.into())).await); + let response = receiver.try_recv().unwrap(); + let Message::Text(response) = response else { + panic!("expected text") + }; + assert_eq!( + serde_json::from_str::(&response).unwrap()["action"], + "error" + ); + } + + #[tokio::test] + async fn offer_returns_answer() { + let (client, mut receiver) = client(); + let remote = PeerConnection::new(RtcConfigurationBuilder::new().build()); + remote.create_data_channel("test", None).unwrap(); + let offer = remote.create_offer().await.unwrap(); + let message = serde_json::json!({ + "action": "sdp-offer", + "sdp": offer.to_sdp_string(), + }); + assert!( + client + .ws_on_message(Message::Text(message.to_string().into())) + .await + ); + let Message::Text(response) = receiver.try_recv().unwrap() else { + panic!("expected text") + }; + let response: serde_json::Value = serde_json::from_str(&response).unwrap(); + assert_eq!(response["action"], "answer"); + assert!(response.get("channel_id").is_none()); + assert!(response["sdp"].as_str().unwrap().starts_with("v=0")); + client.close(); + remote.close(); + } + + #[tokio::test] + async fn invalid_ice_candidate_returns_error() { + let (client, mut receiver) = client(); + let message = r#"{"action":"ice-candidate","candidate":"invalid"}"#; + assert!(client.ws_on_message(Message::Text(message.into())).await); + let Message::Text(response) = receiver.try_recv().unwrap() else { + panic!("expected text") + }; + let response: serde_json::Value = serde_json::from_str(&response).unwrap(); + assert_eq!(response["action"], "error"); + } + + #[tokio::test] + async fn empty_ice_candidate_is_ignored() { + let (client, mut receiver) = client(); + let message = r#"{"action":"ice-candidate","candidate":""}"#; + assert!(client.ws_on_message(Message::Text(message.into())).await); + assert!(receiver.try_recv().is_err()); + client.close(); + } + + #[tokio::test] + async fn audio_offer_prepares_room_without_self_echo() { + let (client, mut receiver) = client(); + let remote = PeerConnection::new(RtcConfigurationBuilder::new().build()); + remote.add_transceiver(MediaKind::Audio, TransceiverDirection::SendRecv); + let offer = remote.create_offer().await.unwrap(); + let message = serde_json::json!({"action": "sdp-offer", "sdp": offer.to_sdp_string()}); + assert!(client.ws_on_message(Message::Text(message.to_string().into())).await); + let Message::Text(response) = receiver.try_recv().unwrap() else { + panic!("expected text") + }; + let response: serde_json::Value = serde_json::from_str(&response).unwrap(); + assert_eq!(response["action"], "answer", "{response}"); + let audio = response["sdp"].as_str().unwrap().split("m=") + .find(|section| section.starts_with("audio ")).unwrap(); + assert!(audio.lines().any(|line| line == "a=sendrecv"), "{audio}"); + assert!(client.peer_connection.get_transceivers()[0].sender().is_some()); + client.close(); + remote.close(); + } + + #[tokio::test] + async fn audio_answer_uses_offered_opus_payload_type() { + let (client, mut receiver) = client(); + let remote = PeerConnection::new(RtcConfigurationBuilder::new().build()); + remote.add_transceiver(MediaKind::Audio, TransceiverDirection::SendRecv); + let offer = remote.create_offer().await.unwrap().to_sdp_string().replace("111", "109"); + let message = serde_json::json!({"action": "sdp-offer", "sdp": offer}); + client.ws_on_message(Message::Text(message.to_string().into())).await; + let Message::Text(response) = receiver.try_recv().unwrap() else { panic!("expected text") }; + let response: serde_json::Value = serde_json::from_str(&response).unwrap(); + assert_eq!(response["action"], "answer", "{response}"); + let audio = response["sdp"].as_str().unwrap().split("m=") + .find(|section| section.starts_with("audio ")).unwrap(); + assert!(audio.lines().any(|line| line == "a=rtpmap:109 opus/48000/2"), "{audio}"); + client.close(); + remote.close(); + } + + #[tokio::test] + async fn joining_listener_gets_a_distinct_audio_offer() { + let (client, mut messages) = client(); + let remote = PeerConnection::new(RtcConfigurationBuilder::new().build()); + remote.add_transceiver(MediaKind::Audio, TransceiverDirection::SendRecv); + let offer = remote.create_offer().await.unwrap().to_sdp_string().replace("111", "109"); + remote.set_local_description(rustrtc::sdp::SessionDescription::parse(rustrtc::sdp::SdpType::Offer, &offer).unwrap()).unwrap(); + let message = serde_json::json!({"action": "sdp-offer", "sdp": offer}); + client.ws_on_message(Message::Text(message.to_string().into())).await; + let Message::Text(answer) = messages.recv().await.unwrap() else { panic!("réponse attendue") }; + let answer: serde_json::Value = serde_json::from_str(&answer).unwrap(); + remote.set_remote_description(rustrtc::sdp::SessionDescription::parse( + rustrtc::sdp::SdpType::Answer, answer["sdp"].as_str().unwrap(), + ).unwrap()).await.unwrap(); + let other = Uuid::new_v4(); + let _other_events = client.room_manager.join(client.channel.id, other); + let Message::Text(offer) = tokio::time::timeout(std::time::Duration::from_secs(2), messages.recv()).await.unwrap().unwrap() else { panic!("offre attendue") }; + let offer: serde_json::Value = serde_json::from_str(&offer).unwrap(); + assert_eq!(offer["action"], "sdp-offer", "{offer}"); + assert_eq!(offer["sdp"].as_str().unwrap().matches("m=audio ").count(), 2); + let outgoing = offer["sdp"].as_str().unwrap().split("m=").filter(|section| section.starts_with("audio ")).nth(1).unwrap(); + assert!(outgoing.lines().next().unwrap().ends_with(" 109"), "{outgoing}"); + let initial_mid = answer["sdp"].as_str().unwrap().lines().find(|line| line.contains("urn:ietf:params:rtp-hdrext:sdes:mid")); + let next_mid = offer["sdp"].as_str().unwrap().lines().find(|line| line.contains("urn:ietf:params:rtp-hdrext:sdes:mid")); + assert_eq!(initial_mid, next_mid); + client.close(); + remote.close(); + } + + #[tokio::test] + async fn returning_participant_gets_a_new_audio_section() { + let (client, mut messages) = client(); + let remote = PeerConnection::new(RtcConfigurationBuilder::new().build()); + remote.add_transceiver(MediaKind::Audio, TransceiverDirection::SendRecv); + let offer = remote.create_offer().await.unwrap(); + remote.set_local_description(offer.clone()).unwrap(); + let message = serde_json::json!({"action": "sdp-offer", "sdp": offer.to_sdp_string()}); + client.ws_on_message(Message::Text(message.to_string().into())).await; + let Message::Text(answer) = messages.recv().await.unwrap() else { panic!("réponse attendue") }; + let answer: serde_json::Value = serde_json::from_str(&answer).unwrap(); + remote.set_remote_description(SessionDescription::parse(SdpType::Answer, answer["sdp"].as_str().unwrap()).unwrap()).await.unwrap(); + + let first = Uuid::new_v4(); + let second = Uuid::new_v4(); + let _first_events = client.room_manager.join(client.channel.id, first); + for (step, expected_sections) in [2, 2, 3].into_iter().enumerate() { + if step == 1 { + let Message::Text(left) = messages.recv().await.unwrap() else { panic!("départ attendu") }; + assert_eq!(serde_json::from_str::(&left).unwrap()["action"], "source-left"); + } + let Message::Text(message) = tokio::time::timeout(std::time::Duration::from_secs(2), messages.recv()).await.unwrap().unwrap() else { panic!("message attendu") }; + let message: serde_json::Value = serde_json::from_str(&message).unwrap(); + assert_eq!(message["action"], "sdp-offer", "{message}"); + assert_eq!(message["sdp"].as_str().unwrap().matches("m=audio ").count(), expected_sections); + let offer = SessionDescription::parse(SdpType::Offer, message["sdp"].as_str().unwrap()).unwrap(); + remote.set_remote_description(offer).await.unwrap(); + let answer = remote.create_answer().await.unwrap(); + remote.set_local_description(answer.clone()).unwrap(); + let response = serde_json::json!({"action": "sdp-answer", "sdp": answer.to_sdp_string()}); + client.ws_on_message(Message::Text(response.to_string().into())).await; + if step == 0 { + client.room_manager.leave(client.channel.id, first); + } else if step == 1 { + let _second_events = client.room_manager.join(client.channel.id, second); + } + } + client.close(); + remote.close(); + } + + #[test] + fn renegotiation_keeps_extension_ids_across_audio_sections() { + let initial = "v=0\r\nm=audio 9 UDP/TLS/RTP/SAVPF 109\r\na=extmap:3 urn:ietf:params:rtp-hdrext:sdes:mid\r\n"; + let offer = "v=0\r\nm=audio 9 UDP/TLS/RTP/SAVPF 109\r\na=extmap:3 http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time\r\na=extmap:4 urn:ietf:params:rtp-hdrext:sdes:mid\r\nm=audio 9 UDP/TLS/RTP/SAVPF 109\r\na=extmap:3 http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time\r\na=extmap:4 urn:ietf:params:rtp-hdrext:sdes:mid\r\n"; + let fixed = RTCClient::stable_extmaps(offer, initial); + assert_eq!(fixed.matches("a=extmap:3 urn:ietf:params:rtp-hdrext:sdes:mid").count(), 2); + assert!(!fixed.contains("a=extmap:3 http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time")); + assert_eq!(fixed.matches("a=extmap:1 http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time").count(), 2); + } + + #[tokio::test] + async fn valid_ice_candidate_is_accepted() { + let (client, mut receiver) = client(); + let message = serde_json::json!({ + "action": "ice-candidate", + "candidate": "candidate:1 1 udp 2130706431 127.0.0.1 12345 typ host", + }); + assert!( + client + .ws_on_message(Message::Text(message.to_string().into())) + .await + ); + assert!(receiver.try_recv().is_err()); + client.close(); + } + + #[test] + fn local_candidate_uses_browser_format() { + let candidate = rustrtc::transports::ice::IceCandidate::from_sdp( + "1 1 udp 2130706431 127.0.0.1 12345 typ host", + ).unwrap(); + assert!(format!("candidate:{}", candidate.to_sdp()).starts_with("candidate:1 1 udp")); + } +} diff --git a/src/rtc/messages.rs b/src/rtc/messages.rs index c3512c1..95f42b2 100644 --- a/src/rtc/messages.rs +++ b/src/rtc/messages.rs @@ -1,18 +1,21 @@ use serde::{Deserialize, Serialize}; -use uuid::Uuid; #[derive(Debug, Deserialize)] #[serde(tag = "action", rename_all = "kebab-case")] pub enum VoiceClientMessage { - SDPOffer { channel_id: Uuid, sdp: String }, - IceCandidate { channel_id: Uuid, candidate: String }, - Leave { channel_id: Uuid }, + #[serde(rename = "sdp-offer")] + SDPOffer { sdp: String }, + SdpAnswer { sdp: String }, + IceCandidate { candidate: String }, + Leave, } #[derive(Debug, Serialize)] #[serde(tag = "action", rename_all = "kebab-case")] pub enum VoiceServerMessage { - Answer { channel_id: Uuid, sdp: String }, - IceCandidate { channel_id: Uuid, candidate: String }, + Answer { sdp: String }, + SdpOffer { sdp: String }, + SourceLeft { id: uuid::Uuid }, + IceCandidate { candidate: String }, Error { message: String }, } diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 9489151..fa7272e 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -4,12 +4,12 @@ mod metrics; pub mod ws_entrypoint; use crate::config::NetworkConfig; -use crate::models::channel; use crate::repositories::Repositories; -use crate::rtc::client::RTCClient; use crate::services::Services; use event_bus::EventBus; use rustrtc::{PeerConnection, RtcConfiguration, RtcConfigurationBuilder}; +use rustrtc::media::MediaSample; +use parking_lot::Mutex; use std::collections::HashMap; use std::fmt; use std::sync::Arc; @@ -24,10 +24,18 @@ use std::sync::Arc; // 9. ICE sélectionne un chemin réseau // 10. La connexion WebRTC devient active -pub struct VoiceRoom {} +pub struct VoiceRoom { + participants: HashMap>, + sources: HashMap>, +} + +pub enum RoomEvent { + Joined(uuid::Uuid, tokio::sync::broadcast::Receiver), + Left(uuid::Uuid), +} pub struct VoiceRoomManager { - rooms: HashMap, + rooms: Mutex>, } #[derive(Clone)] @@ -37,8 +45,7 @@ pub struct RTCManager { pub services: Arc, pub event_bus: Arc, - // - rooms: Arc, + pub(crate) rooms: Arc, } impl fmt::Debug for RTCManager { @@ -52,14 +59,50 @@ impl fmt::Debug for RTCManager { impl VoiceRoom { pub fn new() -> Self { - Self {} + Self { participants: HashMap::new(), sources: HashMap::new() } } } impl VoiceRoomManager { pub fn new() -> Self { Self { - rooms: HashMap::new(), + rooms: Mutex::new(HashMap::new()), + } + } + + pub fn join(&self, channel_id: uuid::Uuid, client_id: uuid::Uuid) -> tokio::sync::mpsc::UnboundedReceiver { + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + let mut rooms = self.rooms.lock(); + let room = rooms.entry(channel_id).or_insert_with(VoiceRoom::new); + let (source, _) = tokio::sync::broadcast::channel(120); + for (id, existing) in &room.sources { + let _ = tx.send(RoomEvent::Joined(*id, existing.subscribe())); + } + for participant in room.participants.values() { + let _ = participant.send(RoomEvent::Joined(client_id, source.subscribe())); + } + room.sources.insert(client_id, source); + room.participants.insert(client_id, tx); + rx + } + + pub fn forward(&self, channel_id: uuid::Uuid, client_id: uuid::Uuid, sample: MediaSample) { + if let Some(source) = self.rooms.lock().get(&channel_id).and_then(|room| room.sources.get(&client_id).cloned()) { + let _ = source.send(sample); + } + } + + pub fn leave(&self, channel_id: uuid::Uuid, client_id: uuid::Uuid) { + let mut rooms = self.rooms.lock(); + if let Some(room) = rooms.get_mut(&channel_id) { + room.participants.remove(&client_id); + room.sources.remove(&client_id); + for participant in room.participants.values() { + let _ = participant.send(RoomEvent::Left(client_id)); + } + if room.participants.is_empty() { + rooms.remove(&channel_id); + } } } } @@ -88,17 +131,68 @@ impl RTCManager { } } - /// Session Description Protocol - Permet de se mettre d'accord sur les paramètres media - pub async fn handle_sdp_offer(&self, rtc_client: &RTCClient, offer_sdp: String) { - let pc = rtc_client.peer_connection.clone(); - } - - /// Interactive Connectivity Establishment - Permet de négocier un chemin réseau - pub async fn handle_ice_candidate(&self) {} - pub fn new_peer_connection(&self) -> Arc { let pc = Arc::new(PeerConnection::new(self.config.clone())); // todo : some logging ? pc } } + +#[cfg(test)] +mod room_tests { + use super::{RoomEvent, VoiceRoomManager}; + use rustrtc::media::{AudioFrame, MediaSample}; + use uuid::Uuid; + + #[tokio::test] + async fn audio_reaches_only_other_clients_in_same_room() { + let rooms = VoiceRoomManager::new(); + let channel = Uuid::new_v4(); + let sender = Uuid::new_v4(); + let listener = Uuid::new_v4(); + let mut own = rooms.join(channel, sender); + let mut other = rooms.join(channel, listener); + let mut separate = rooms.join(Uuid::new_v4(), Uuid::new_v4()); + assert!(matches!(own.try_recv(), Ok(RoomEvent::Joined(id, _)) if id == listener)); + let mut incoming = match other.try_recv().unwrap() { RoomEvent::Joined(id, receiver) if id == sender => receiver, _ => panic!("mauvaise source") }; + rooms.forward(channel, sender, MediaSample::Audio(AudioFrame::default())); + assert!(matches!(incoming.try_recv(), Ok(MediaSample::Audio(_)))); + assert!(own.try_recv().is_err()); + assert!(separate.try_recv().is_err()); + rooms.leave(channel, listener); + rooms.forward(channel, sender, MediaSample::Audio(AudioFrame::default())); + assert!(other.try_recv().is_err()); + rooms.leave(channel, sender); + assert!(!rooms.rooms.lock().contains_key(&channel)); + } + + #[tokio::test] + async fn simultaneous_sources_remain_separate() { + let rooms = VoiceRoomManager::new(); + let channel = Uuid::new_v4(); + let first = Uuid::new_v4(); + let second = Uuid::new_v4(); + let listener = Uuid::new_v4(); + let mut first_events = rooms.join(channel, first); + let mut second_events = rooms.join(channel, second); + let mut listener_events = rooms.join(channel, listener); + let mut sources = std::collections::HashMap::new(); + for _ in 0..2 { + if let RoomEvent::Joined(id, rx) = listener_events.try_recv().unwrap() { + sources.insert(id, rx); + } + } + let mut first_audio = sources.remove(&first).unwrap(); + let mut second_audio = sources.remove(&second).unwrap(); + rooms.forward(channel, first, MediaSample::Audio(AudioFrame::default())); + rooms.forward(channel, second, MediaSample::Audio(AudioFrame::default())); + assert!(matches!(first_audio.try_recv(), Ok(MediaSample::Audio(_)))); + assert!(matches!(second_audio.try_recv(), Ok(MediaSample::Audio(_)))); + assert!(first_audio.try_recv().is_err()); + assert!(second_audio.try_recv().is_err()); + assert!(matches!(first_events.try_recv(), Ok(RoomEvent::Joined(id, _)) if id == second)); + assert!(matches!(first_events.try_recv(), Ok(RoomEvent::Joined(id, _)) if id == listener)); + assert!(matches!(second_events.try_recv(), Ok(RoomEvent::Joined(id, _)) if id == first)); + assert!(matches!(second_events.try_recv(), Ok(RoomEvent::Joined(id, _)) if id == listener)); + } +} diff --git a/src/rtc/test.html b/src/rtc/test.html new file mode 100644 index 0000000..8ce4d65 --- /dev/null +++ b/src/rtc/test.html @@ -0,0 +1,112 @@ + + + + + Test RTC + + +

Test RTC (un client)

+

Ouvrir ce fichier dans un navigateur. Un canal existant et un JWT de connexion sont nécessaires. Ne partage pas le JWT ni les logs contenant l'URL du WebSocket.

+
+
+
+ + +

+    
+
+
\ No newline at end of file
diff --git a/src/rtc/ws_entrypoint.rs b/src/rtc/ws_entrypoint.rs
index 8edfadd..699d031 100644
--- a/src/rtc/ws_entrypoint.rs
+++ b/src/rtc/ws_entrypoint.rs
@@ -30,7 +30,8 @@ pub async fn ws_entrypoint_handler(
     let (tx, mut rx) = mpsc::unbounded_channel::();
 
     let peer_connection = state.rtc.new_peer_connection();
-    let mut rtc_client = RTCClient::new(user, channel, peer_connection, tx);
+    let rtc_client = RTCClient::new(user, channel, peer_connection, tx, state.rtc.rooms.clone());
+    let ice_task = rtc_client.forward_ice_candidates();
 
     // Task pour envoyer les message au frontend depuis RTCClient
     let mut send_task = tokio::spawn(async move {
@@ -45,12 +46,23 @@ pub async fn ws_entrypoint_handler(
     let client_clone = rtc_client.clone();
     let mut recv_task = tokio::spawn(async move {
         while let Some(Ok(message)) = receiver.next().await {
-            client_clone.ws_on_message(message).await;
+            if !client_clone.ws_on_message(message).await {
+                break;
+            }
         }
     });
 
     tokio::select! {
-        _ = (&mut send_task) => recv_task.abort(),
-        _ = (&mut recv_task) => send_task.abort(),
-    };
+        _ = &mut send_task => {
+            recv_task.abort();
+            let _ = recv_task.await;
+        }
+        _ = &mut recv_task => {
+            send_task.abort();
+            let _ = send_task.await;
+        }
+    }
+    ice_task.abort();
+    let _ = ice_task.await;
+    rtc_client.close();
 }