diff --git a/src/core/mod.rs b/src/core/mod.rs index 18e100b..36afa6b 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -6,6 +6,7 @@ use crate::http::server::HttpServer; use crate::metrics::{AppMetrics, reporter}; use crate::repositories::Repositories; use crate::routes::gateway::{GatewayManager, RealtimeRouter}; +use crate::rtc::RTCManager; use crate::services::Services; use crate::voice::VoiceService; use event_bus::EventBus; @@ -82,6 +83,13 @@ impl App { )) .start(event_bus.clone()); + let rtc = Arc::new(RTCManager::new( + &config.network, + repositories.clone(), + services.clone(), + event_bus.clone(), + )); + let state = AppState { db, config: Arc::new(config), @@ -93,6 +101,7 @@ impl App { event_bus, services, voice, + rtc, }; Ok(Self { state }) diff --git a/src/core/state.rs b/src/core/state.rs index c8e2a75..9b60379 100644 --- a/src/core/state.rs +++ b/src/core/state.rs @@ -3,6 +3,7 @@ use crate::metrics::AppMetrics; use crate::models::server; use crate::repositories::Repositories; use crate::routes::gateway::GatewayManager; +use crate::rtc::RTCManager; use crate::services::Services; use crate::voice::VoiceService; use event_bus::EventBus; @@ -20,6 +21,7 @@ pub struct AppState { pub gateway: Arc, pub event_bus: Arc, pub services: Arc, + pub rtc: Arc, pub voice: Arc, } diff --git a/src/rtc/client.rs b/src/rtc/client.rs index d085275..fd75d42 100644 --- a/src/rtc/client.rs +++ b/src/rtc/client.rs @@ -7,14 +7,14 @@ use tokio::sync::mpsc::UnboundedSender; /// Client connecté à un canal RTC. #[derive(Clone)] -pub struct RtcClient { +pub struct RTCClient { pub user: Arc, pub channel: Arc, pub peer_connection: Arc, pub websocket_sender: UnboundedSender, } -impl RtcClient { +impl RTCClient { pub fn new( user: user::Model, channel: channel::Model, @@ -29,7 +29,7 @@ impl RtcClient { } } - pub fn on_message(&self, raw_message: Message) { + pub async fn on_message(&self, raw_message: Message) { let Message::Text(raw_message) = raw_message else { return; }; diff --git a/src/rtc/manager.rs b/src/rtc/manager.rs deleted file mode 100644 index 604bf25..0000000 --- a/src/rtc/manager.rs +++ /dev/null @@ -1,36 +0,0 @@ -use crate::config::NetworkConfig; -use rustrtc::{RtcConfiguration, RtcConfigurationBuilder}; - -// 1. Client crée une RTCPeerConnection -// 2. Client crée une SDP offer -// 3. Client envoie l'offer via WebSocket -// 4. Serveur crée sa PeerConnection -// 5. Serveur applique l'offer -// 6. Serveur crée une SDP answer -// 7. Serveur renvoie l'answer via WebSocket -// 8. Client et serveur échangent les candidats ICE -// 9. ICE sélectionne un chemin réseau -// 10. La connexion WebRTC devient active - -#[derive(Debug, Clone)] -struct RTCManager { - pub config: RtcConfiguration, -} - -impl RTCManager { - pub fn new(network: &NetworkConfig) -> Self { - let builder = RtcConfigurationBuilder::new() - .ice_udp_mux(true) - .ice_udp_mux_port(network.udp_port) - .bind_ip(network.host.to_string()); - let config = builder.build(); - - Self { config } - } - - /// Permet de se mettre d'accord sur les paramètres media - pub async fn handle_sdp_offer(&self) {} - - /// Interactive Connectivity Establishment - Permet de négocier un chemin réseau - pub async fn handle_ice_candidate(&self) {} -} diff --git a/src/rtc/mod.rs b/src/rtc/mod.rs index 2590deb..7ad7dae 100644 --- a/src/rtc/mod.rs +++ b/src/rtc/mod.rs @@ -1,25 +1,64 @@ mod client; -mod manager; mod messages; mod metrics; pub mod ws_entrypoint; use crate::config::NetworkConfig; -use metrics::VoiceMetrics; -use rustrtc::{RtcConfiguration, RtcConfigurationBuilder}; +use crate::repositories::Repositories; +use crate::rtc::client::RTCClient; +use crate::services::Services; +use event_bus::EventBus; +use rustrtc::{PeerConnection, RtcConfiguration, RtcConfigurationBuilder}; use std::sync::Arc; +// 1. Client crée une RTCPeerConnection +// 2. Client crée une SDP offer +// 3. Client envoie l'offer via WebSocket +// 4. Serveur crée sa PeerConnection +// 5. Serveur applique l'offer +// 6. Serveur crée une SDP answer +// 7. Serveur renvoie l'answer via WebSocket +// 8. Client et serveur échangent les candidats ICE +// 9. ICE sélectionne un chemin réseau +// 10. La connexion WebRTC devient active +#[derive(Debug, Clone)] pub struct RTCManager { pub config: RtcConfiguration, + pub repositories: Arc, + pub services: Arc, + pub event_bus: Arc, } impl RTCManager { - pub fn new(network: &NetworkConfig, metrics: Arc) -> Self { + pub fn new( + network: &NetworkConfig, + repositories: Arc, + services: Arc, + event_bus: Arc, + ) -> Self { let builder = RtcConfigurationBuilder::new() .ice_udp_mux(true) .ice_udp_mux_port(network.udp_port) .bind_ip(network.host.to_string()); - let rtc_config = builder.build(); - Self { config: rtc_config } + let config = builder.build(); + + Self { + config, + repositories, + services, + event_bus, + } + } + + /// 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) {} + + /// 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 } } diff --git a/src/rtc/ws_entrypoint.rs b/src/rtc/ws_entrypoint.rs index 0bf4c90..ad72bb6 100644 --- a/src/rtc/ws_entrypoint.rs +++ b/src/rtc/ws_entrypoint.rs @@ -1,12 +1,10 @@ // This is the first point when WebRTC ask for a connection -use super::client::RtcClient; +use super::client::RTCClient; use crate::core::AppState; use crate::models::{channel, user}; use axum::extract::ws::{Message, WebSocket}; use futures_util::{SinkExt, StreamExt}; -use rustrtc::PeerConnection; -use std::sync::Arc; use tokio::sync::mpsc; // 1. Le frontend ouvre /rtc/{channel_id} @@ -31,10 +29,11 @@ pub async fn ws_entrypoint_handler( let (mut sender, mut receiver) = socket.split(); let (tx, mut rx) = mpsc::unbounded_channel::(); - let peer_connection = Arc::new(PeerConnection::new(state.rtc.config.clone())); - let mut rtc_client = RtcClient::new(user, channel, peer_connection, tx); + let peer_connection = state.rtc.new_peer_connection(); + let mut rtc_client = RTCClient::new(user, channel, peer_connection, tx); - let send_task = tokio::spawn(async move { + // Task pour envoyer les message au frontend depuis RTCClient + let mut send_task = tokio::spawn(async move { while let Some(message) = rx.recv().await { if sender.send(message).await.is_err() { break; @@ -42,11 +41,16 @@ pub async fn ws_entrypoint_handler( } }); + // Task pour recevoir les messages du frontend afin de les transmettre à RTCClient let client_clone = rtc_client.clone(); - let state_clone = state.clone(); let mut recv_task = tokio::spawn(async move { while let Some(Ok(message)) = receiver.next().await { - client_clone.on_message(message, &state_clone).await; + client_clone.on_message(message).await; } }); + + tokio::select! { + _ = (&mut send_task) => recv_task.abort(), + _ = (&mut recv_task) => send_task.abort(), + }; }