Compare commits

..
2 Commits
Author SHA1 Message Date
Nell d1ce0655fb add event_bus_typed 2026-09-23 14:30:19 +02:00
Nell b234359c4a add webrtc 2026-09-23 13:58:42 +02:00
8 changed files with 1164 additions and 48 deletions
+361
View File
@@ -0,0 +1,361 @@
<script lang="ts" setup>
import {onUnmounted, ref} from 'vue'
const channelId = ref('')
const logs = ref<string[]>([])
const audioStats = ref<string[]>([])
const connected = ref(false)
const audioFile = ref<File | null>(null)
const sourceLevel = ref(0)
const roomLevel = ref(0)
let socket: WebSocket | null = null
let peer: RTCPeerConnection | null = null
let microphone: MediaStream | null = null
let microphoneContext: AudioContext | null = null
let microphoneAnalyser: AnalyserNode | null = null
let roomContext: AudioContext | null = null
const roomAnalysers = new Map<MediaStreamTrack, AnalyserNode>()
let fileAudio: HTMLAudioElement | null = null
let fileUrl: string | null = null
let fileStream: MediaStream | null = null
let pendingCandidates: string[] = []
let audioStatsTimer: ReturnType<typeof setInterval> | null = null
let levelTimer: ReturnType<typeof setInterval> | null = null
const remoteAudio = ref<HTMLDivElement | null>(null)
const players = new Map<MediaStreamTrack, HTMLAudioElement>()
let roomStream: MediaStream | null = null
const sourceTracks = new Map<string, MediaStreamTrack>()
function rms(analyser: AnalyserNode | null, context: AudioContext | null) {
if (!analyser || context?.state !== 'running') return 0
const samples = new Float32Array(analyser.fftSize)
analyser.getFloatTimeDomainData(samples)
return Math.sqrt(samples.reduce((sum, sample) => sum + sample * sample, 0) / samples.length)
}
function updateLevels() {
sourceLevel.value = Math.min(1, Math.sqrt(rms(microphoneAnalyser, microphoneContext)))
roomLevel.value = Math.min(1, Math.sqrt([...roomAnalysers.values()].reduce((sum, analyser) => sum + rms(analyser, roomContext) ** 2, 0)))
}
function removeRoomTrack(track: MediaStreamTrack) {
roomStream?.removeTrack(track)
const player = players.get(track)
if (player) {
player.pause()
player.srcObject = null
player.remove()
players.delete(track)
}
roomAnalysers.get(track)?.disconnect()
roomAnalysers.delete(track)
roomLevel.value = 0
}
function log(message: string) {
logs.value.push(`${new Date().toLocaleTimeString()} ${message}`)
}
async function logSelectedIcePair(connection: RTCPeerConnection) {
const stats = await connection.getStats()
const pair = [...stats.values()].find(report => report.type === 'candidate-pair' && report.nominated && report.state === 'succeeded')
if (!pair) {
log('Paire ICE sélectionnée introuvable dans les statistiques.')
return
}
const local = stats.get(pair.localCandidateId)
const remote = stats.get(pair.remoteCandidateId)
log(`Paire ICE sélectionnée : local ${local?.address ?? local?.ip ?? '?'}:${local?.port ?? '?'} (${local?.candidateType ?? '?'}) → distant ${remote?.address ?? remote?.ip ?? '?'}:${remote?.port ?? '?'} (${remote?.candidateType ?? '?'})`)
}
async function logAudioStats(connection: RTCPeerConnection) {
const stats = await connection.getStats()
if (peer !== connection) return
const sent = [...stats.values()].find(report => report.type === 'outbound-rtp' && report.kind === 'audio')
const activeTracks = new Set(roomStream?.getAudioTracks().map(track => track.id) ?? [])
const inbound = [...stats.values()].filter(report => report.type === 'inbound-rtp' && report.kind === 'audio' && activeTracks.has(report.trackIdentifier))
const received = inbound.find(report => (report.packetsReceived ?? 0) > 0) ?? inbound[0]
const lines = [`Audio : envoyé ${sent?.packetsSent ?? '?'} paquets / ${sent?.bytesSent ?? '?'} octets ; reçu ${inbound.length ? inbound.reduce((sum, report) => sum + (report.packetsReceived ?? 0), 0) : '?'} paquets / ${inbound.length ? inbound.reduce((sum, report) => sum + (report.bytesReceived ?? 0), 0) : '?'} octets (${inbound.length} piste(s))`]
const level = (value: number | undefined) => typeof value === 'number' ? value.toFixed(3) : '?'
const micLevel = microphoneAnalyser && microphoneContext?.state === 'running' ? rms(microphoneAnalyser, microphoneContext) : undefined
lines.push(`Niveaux audio (0 à 1) : source RMS ${level(micLevel)} ; retour ${level(received?.audioLevel)} ; pertes ${received?.packetsLost ?? '?'} ; gigue ${received?.jitter ?? '?'} s ; échantillons masqués ${received?.concealedSamples ?? '?'}/${received?.totalSamplesReceived ?? '?'}`)
audioStats.value = lines
}
function disconnect() {
if (audioStatsTimer) clearInterval(audioStatsTimer)
audioStatsTimer = null
if (levelTimer) clearInterval(levelTimer)
levelTimer = null
sourceLevel.value = 0
roomLevel.value = 0
if (socket?.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({action: 'leave'}))
socket.close()
} else {
socket?.close()
}
socket = null
peer?.close()
peer = null
microphone?.getTracks().forEach(track => track.stop())
microphone = null
fileAudio?.pause()
fileAudio = null
fileStream?.getTracks().forEach(track => track.stop())
fileStream = null
if (fileUrl) URL.revokeObjectURL(fileUrl)
fileUrl = null
microphoneAnalyser = null
if (microphoneContext) void microphoneContext.close()
microphoneContext = null
roomAnalysers.forEach(analyser => analyser.disconnect())
roomAnalysers.clear()
if (roomContext) void roomContext.close()
roomContext = null
for (const track of players.keys()) removeRoomTrack(track)
roomStream = null
sourceTracks.clear()
pendingCandidates = []
connected.value = false
}
async function connect() {
const channel = channelId.value.trim()
if (!/^[0-9a-fA-F-]{36}$/.test(channel)) {
log('Renseigner un identifiant de canal valide.')
return
}
try {
audioStats.value = []
let stream: MediaStream
if (audioFile.value) {
const context = new AudioContext()
microphoneContext = context
fileUrl = URL.createObjectURL(audioFile.value)
const player = new Audio(fileUrl)
fileAudio = player
player.loop = true
const source = context.createMediaElementSource(player)
const destination = context.createMediaStreamDestination()
const analyser = context.createAnalyser()
analyser.fftSize = 2048
source.connect(destination)
source.connect(analyser)
microphoneAnalyser = analyser
await context.resume()
stream = destination.stream
fileStream = stream
await player.play()
log(`Fichier audio envoyé : ${audioFile.value.name} (en boucle, sans lecture locale)`)
} else {
stream = await navigator.mediaDevices.getUserMedia({audio: true})
microphone = stream
try {
const context = new AudioContext()
microphoneContext = context
const analyser = context.createAnalyser()
analyser.fftSize = 2048
context.createMediaStreamSource(stream).connect(analyser)
microphoneAnalyser = analyser
await context.resume()
} catch (error) {
log(`Mesure du niveau micro indisponible : ${String(error)}`)
}
}
const url = new URL(`/ws/rtc/${encodeURIComponent(channel)}`, window.location.href)
url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:'
const connection = new RTCPeerConnection()
const ws = new WebSocket(url)
peer = connection
socket = ws
connected.value = true
stream.getAudioTracks().forEach(track => connection.addTrack(track, stream))
roomStream = new MediaStream()
connection.ontrack = ({streams, track}) => {
if (peer !== connection || players.has(track)) return
const remoteStream = roomStream!
remoteStream.addTrack(track)
if (streams[0]) sourceTracks.set(streams[0].id, track)
track.onended = () => removeRoomTrack(track)
if (remoteAudio.value) {
const player = new Audio()
player.autoplay = true
player.controls = true
player.srcObject = new MediaStream([track])
players.set(track, player)
remoteAudio.value.append(player)
void player.play().catch(error => log(`Lecture audio impossible : ${String(error)}`))
}
if (track.kind === 'audio') {
if (!roomContext) {
roomContext = new AudioContext()
void roomContext.resume().catch(error => log(`Mesure du retour indisponible : ${String(error)}`))
}
const analyser = roomContext.createAnalyser()
analyser.fftSize = 2048
roomContext.createMediaStreamSource(new MediaStream([track])).connect(analyser)
roomAnalysers.set(track, analyser)
}
log('Piste audio distante reçue')
}
connection.onicecandidate = ({candidate}) => {
if (candidate && ws.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify({action: 'ice-candidate', candidate: candidate.candidate}))
log('Candidat ICE local envoyé')
}
}
connection.oniceconnectionstatechange = () => {
log(`ICE : ${connection.iceConnectionState}`)
if (connection.iceConnectionState === 'connected' || connection.iceConnectionState === 'completed') {
void logSelectedIcePair(connection).catch(error => log(`Statistiques ICE indisponibles : ${String(error)}`))
}
}
connection.onconnectionstatechange = () => {
log(`Connexion : ${connection.connectionState}`)
if (connection.connectionState === 'connected' && !audioStatsTimer) {
levelTimer = setInterval(updateLevels, 100)
void logAudioStats(connection).catch(error => log(`Statistiques audio indisponibles : ${String(error)}`))
audioStatsTimer = setInterval(() => {
void logAudioStats(connection).catch(error => log(`Statistiques audio indisponibles : ${String(error)}`))
}, 2000)
} else if (connection.connectionState !== 'connected' && audioStatsTimer) {
clearInterval(audioStatsTimer)
audioStatsTimer = null
if (levelTimer) clearInterval(levelTimer)
levelTimer = null
sourceLevel.value = 0
roomLevel.value = 0
}
}
ws.onopen = async () => {
log('WebSocket ouvert')
try {
const offer = await connection.createOffer()
await connection.setLocalDescription(offer)
if (ws.readyState !== WebSocket.OPEN) return
ws.send(JSON.stringify({action: 'sdp-offer', sdp: connection.localDescription?.sdp}))
log('Offre SDP envoyée')
} catch (error) {
log(`Offre impossible : ${String(error)}`)
disconnect()
}
}
let messageQueue = Promise.resolve()
ws.onmessage = ({data}) => {
messageQueue = messageQueue.then(async () => {
if (socket !== ws) return
try {
const message = JSON.parse(data)
if (message.action === 'answer') {
await connection.setRemoteDescription({type: 'answer', sdp: message.sdp})
log('Réponse SDP appliquée')
for (const candidate of pendingCandidates) {
await connection.addIceCandidate({candidate, sdpMLineIndex: 0})
}
pendingCandidates = []
} else if (message.action === 'sdp-offer') {
await connection.setRemoteDescription({type: 'offer', sdp: message.sdp})
const answer = await connection.createAnswer()
await connection.setLocalDescription(answer)
ws.send(JSON.stringify({action: 'sdp-answer', sdp: connection.localDescription?.sdp}))
log('Pistes du salon renégociées')
} else if (message.action === 'source-left') {
const track = sourceTracks.get(message.id)
if (track) {
removeRoomTrack(track)
sourceTracks.delete(message.id)
}
log(`Participant parti : ${message.id}`)
} else if (message.action === 'ice-candidate') {
if (connection.remoteDescription) {
await connection.addIceCandidate({candidate: message.candidate, sdpMLineIndex: 0})
} else {
pendingCandidates.push(message.candidate)
}
log(`Candidat ICE distant reçu : ${message.candidate}`)
} else if (message.action === 'error') {
log(`Erreur serveur : ${message.message}`)
} else {
log(`Message inconnu : ${message.action}`)
}
} catch (error) {
log(`Message impossible à traiter : ${String(error)}`)
}
})
}
ws.onerror = () => log('Erreur WebSocket (vérifier le serveur, le canal et la session).')
ws.onclose = ({code}) => {
log(`WebSocket fermé (${code})`)
if (socket === ws) disconnect()
}
} catch (error) {
log(`Connexion impossible : ${String(error)}`)
disconnect()
}
}
onUnmounted(disconnect)
</script>
<template>
<v-main class="rtc-test pa-6">
<div class="rtc-test-content">
<h1 class="text-h5 mb-4">Test RTC</h1>
<!-- <p class="mb-4">Ouvre cette page dans deux onglets avec le même canal. Chaque onglet entend l'autre, sans retour de son propre micro ou fichier audio.</p>-->
<v-text-field v-model="channelId" :disabled="connected" label="Identifiant du canal"/>
<input :disabled="connected" accept="audio/*" class="d-block mb-4" type="file"
@change="audioFile = ($event.target as HTMLInputElement).files?.[0] ?? null"/>
<v-btn v-if="!connected" color="primary" @click="connect">Connecter</v-btn>
<v-btn v-else color="primary" @click="disconnect">Quitter</v-btn>
<div ref="remoteAudio" class="mt-4"/>
<div class="mt-4 rtc-test-panel">
<h2 class="text-h6">Niveaux audio</h2>
<label for="source-level">Ma source (micro ou fichier)</label>
<progress id="source-level" :value="sourceLevel" class="rtc-test-level" max="1"/>
<label for="room-level">Son reçu du salon (autres participants)</label>
<progress id="room-level" :value="roomLevel" class="rtc-test-level" max="1"/>
</div>
<div class="mt-4 rtc-test-panel">
<h2 class="text-h6">Statistiques WebRTC</h2>
<pre class="rtc-test-text">{{ audioStats.join('\n') || 'En attente de connexion…' }}</pre>
</div>
<div class="mt-4 rtc-test-panel">
<h2 class="text-h6">Journal</h2>
<pre class="rtc-test-text rtc-test-logs">{{ logs.join('\n') }}</pre>
</div>
</div>
</v-main>
</template>
<style scoped>
.rtc-test-content {
max-width: 900px;
margin: 0 auto;
}
.rtc-test-panel {
padding: 16px;
border: 1px solid rgba(128, 128, 128, 0.5);
border-radius: 8px;
}
.rtc-test-level {
display: block;
width: 100%;
height: 20px;
margin-bottom: 12px;
}
.rtc-test-text {
white-space: pre-wrap;
overflow-wrap: anywhere;
margin: 0;
}
.rtc-test-logs {
max-height: 360px;
overflow-y: auto;
}
</style>
+5
View File
@@ -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})', path: 'server/:serverId(default|[0-9a-fA-F-]{36})',
name: 'server-dashboard', name: 'server-dashboard',
+4
View File
@@ -51,6 +51,10 @@ export default defineConfig({
target: 'http://localhost:8080', target: 'http://localhost:8080',
changeOrigin: true, changeOrigin: true,
}, },
'/ws': {
target: 'ws://localhost:8080',
ws: true,
},
}, },
}, },
}) })
+546 -21
View File
@@ -1,9 +1,18 @@
use super::messages::VoiceClientMessage; use super::messages::{VoiceClientMessage, VoiceServerMessage};
use super::{RoomEvent, VoiceRoomManager};
use crate::models::{channel, user}; use crate::models::{channel, user};
use axum::extract::ws::Message; 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 std::sync::Arc;
use tokio::sync::mpsc::UnboundedSender; use tokio::sync::mpsc::UnboundedSender;
use tokio::sync::mpsc;
use parking_lot::Mutex;
use tracing::{info, warn};
/// Client connecté à un canal RTC. /// Client connecté à un canal RTC.
#[derive(Clone)] #[derive(Clone)]
@@ -12,6 +21,12 @@ pub struct RTCClient {
pub channel: Arc<channel::Model>, pub channel: Arc<channel::Model>,
pub peer_connection: Arc<PeerConnection>, pub peer_connection: Arc<PeerConnection>,
pub websocket_sender: UnboundedSender<Message>, pub websocket_sender: UnboundedSender<Message>,
pub room_manager: Arc<VoiceRoomManager>,
pub client_id: uuid::Uuid,
room_task: Arc<Mutex<Option<tokio::task::JoinHandle<()>>>>,
negotiation_done: mpsc::UnboundedSender<()>,
negotiation_rx: Arc<Mutex<Option<mpsc::UnboundedReceiver<()>>>>,
pending_room: Arc<Mutex<Option<(mpsc::UnboundedReceiver<RoomEvent>, u8, String)>>>,
} }
impl RTCClient { impl RTCClient {
@@ -20,41 +35,551 @@ impl RTCClient {
channel: channel::Model, channel: channel::Model,
peer_connection: Arc<PeerConnection>, peer_connection: Arc<PeerConnection>,
websocket_sender: UnboundedSender<Message>, websocket_sender: UnboundedSender<Message>,
room_manager: Arc<VoiceRoomManager>,
) -> Self { ) -> Self {
let (negotiation_done, negotiation_rx) = mpsc::unbounded_channel();
Self { Self {
user: Arc::new(user), user: Arc::new(user),
channel: Arc::new(channel), channel: Arc::new(channel),
peer_connection, peer_connection,
websocket_sender, 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) { /// Retourne `false` quand le client demande à quitter le canal.
let Message::Text(raw_message) = raw_message else { pub async fn ws_on_message(&self, raw_message: Message) -> bool {
return; match raw_message {
}; Message::Close(_) => false,
Message::Text(text) => match serde_json::from_str::<VoiceClientMessage>(&text) {
let parsed = match serde_json::from_str::<VoiceClientMessage>(&raw_message) { Ok(VoiceClientMessage::Leave) => false,
Ok(message) => message, Ok(VoiceClientMessage::SdpAnswer { sdp }) => {
Err(_) => return, match SessionDescription::parse(SdpType::Answer, &sdp) {
}; Ok(answer) => match self.peer_connection.set_remote_description(answer).await {
Ok(()) => { let _ = self.negotiation_done.send(()); }
match parsed { Err(error) => warn!(%error, "Réponse de renégociation refusée"),
VoiceClientMessage::SDPOffer { channel_id, sdp } => {} },
VoiceClientMessage::IceCandidate { Err(error) => warn!(%error, "Réponse SDP invalide"),
channel_id, }
candidate, true
} => {} }
VoiceClientMessage::Leave { channel_id } => {} 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) { async fn handle_sdp_offer(&self, sdp: &str) -> anyhow::Result<String> {
let _ = self.websocket_sender.send(raw_message); 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::<Vec<_>>();
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::<u8>().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<u8>) -> anyhow::Result<String> {
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::<Vec<_>>().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::<String, u16>::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::<u16>().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<RoomEvent>, 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<uuid::Uuid, (Arc<rustrtc::peer_connection::RtpTransceiver>, 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<media::SampleStreamTrack>,
room_manager: Arc<VoiceRoomManager>,
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. /// Ferme proprement la connexion WebRTC.
pub fn close(&self) { 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(); 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<Message>) {
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::<super::VoiceClientMessage>(&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::<serde_json::Value>(&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::<serde_json::Value>(&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"));
}
}
+9 -6
View File
@@ -1,18 +1,21 @@
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use uuid::Uuid;
#[derive(Debug, Deserialize)] #[derive(Debug, Deserialize)]
#[serde(tag = "action", rename_all = "kebab-case")] #[serde(tag = "action", rename_all = "kebab-case")]
pub enum VoiceClientMessage { pub enum VoiceClientMessage {
SDPOffer { channel_id: Uuid, sdp: String }, #[serde(rename = "sdp-offer")]
IceCandidate { channel_id: Uuid, candidate: String }, SDPOffer { sdp: String },
Leave { channel_id: Uuid }, SdpAnswer { sdp: String },
IceCandidate { candidate: String },
Leave,
} }
#[derive(Debug, Serialize)] #[derive(Debug, Serialize)]
#[serde(tag = "action", rename_all = "kebab-case")] #[serde(tag = "action", rename_all = "kebab-case")]
pub enum VoiceServerMessage { pub enum VoiceServerMessage {
Answer { channel_id: Uuid, sdp: String }, Answer { sdp: String },
IceCandidate { channel_id: Uuid, candidate: String }, SdpOffer { sdp: String },
SourceLeft { id: uuid::Uuid },
IceCandidate { candidate: String },
Error { message: String }, Error { message: String },
} }
+110 -16
View File
@@ -4,12 +4,12 @@ mod metrics;
pub mod ws_entrypoint; pub mod ws_entrypoint;
use crate::config::NetworkConfig; use crate::config::NetworkConfig;
use crate::models::channel;
use crate::repositories::Repositories; use crate::repositories::Repositories;
use crate::rtc::client::RTCClient;
use crate::services::Services; use crate::services::Services;
use event_bus::EventBus; use event_bus::EventBus;
use rustrtc::{PeerConnection, RtcConfiguration, RtcConfigurationBuilder}; use rustrtc::{PeerConnection, RtcConfiguration, RtcConfigurationBuilder};
use rustrtc::media::MediaSample;
use parking_lot::Mutex;
use std::collections::HashMap; use std::collections::HashMap;
use std::fmt; use std::fmt;
use std::sync::Arc; use std::sync::Arc;
@@ -24,10 +24,18 @@ use std::sync::Arc;
// 9. ICE sélectionne un chemin réseau // 9. ICE sélectionne un chemin réseau
// 10. La connexion WebRTC devient active // 10. La connexion WebRTC devient active
pub struct VoiceRoom {} pub struct VoiceRoom {
participants: HashMap<uuid::Uuid, tokio::sync::mpsc::UnboundedSender<RoomEvent>>,
sources: HashMap<uuid::Uuid, tokio::sync::broadcast::Sender<MediaSample>>,
}
pub enum RoomEvent {
Joined(uuid::Uuid, tokio::sync::broadcast::Receiver<MediaSample>),
Left(uuid::Uuid),
}
pub struct VoiceRoomManager { pub struct VoiceRoomManager {
rooms: HashMap<channel::Model, VoiceRoom>, rooms: Mutex<HashMap<uuid::Uuid, VoiceRoom>>,
} }
#[derive(Clone)] #[derive(Clone)]
@@ -37,8 +45,7 @@ pub struct RTCManager {
pub services: Arc<Services>, pub services: Arc<Services>,
pub event_bus: Arc<EventBus>, pub event_bus: Arc<EventBus>,
// pub(crate) rooms: Arc<VoiceRoomManager>,
rooms: Arc<VoiceRoomManager>,
} }
impl fmt::Debug for RTCManager { impl fmt::Debug for RTCManager {
@@ -52,14 +59,50 @@ impl fmt::Debug for RTCManager {
impl VoiceRoom { impl VoiceRoom {
pub fn new() -> Self { pub fn new() -> Self {
Self {} Self { participants: HashMap::new(), sources: HashMap::new() }
} }
} }
impl VoiceRoomManager { impl VoiceRoomManager {
pub fn new() -> Self { pub fn new() -> Self {
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<RoomEvent> {
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<PeerConnection> { pub fn new_peer_connection(&self) -> Arc<PeerConnection> {
let pc = Arc::new(PeerConnection::new(self.config.clone())); let pc = Arc::new(PeerConnection::new(self.config.clone()));
// todo : some logging ? // todo : some logging ?
pc 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));
}
}
+112
View File
@@ -0,0 +1,112 @@
<!doctype html>
<html lang="fr">
<head>
<meta charset="utf-8">
<title>Test RTC</title>
</head>
<body>
<h1>Test RTC (un client)</h1>
<p>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.</p>
<label>Serveur <input id="server" value="http://localhost:8080"></label><br>
<label>Channel ID <input id="channel" size="40"></label><br>
<label>JWT <input id="token" type="password" size="60"></label><br>
<button id="connect">Connecter</button>
<button id="disconnect" disabled>Quitter</button>
<pre id="log"></pre>
<script>
const log = (message) => {
document.querySelector('#log').textContent += `${new Date().toLocaleTimeString()} ${message}\n`;
};
const connectButton = document.querySelector('#connect');
const disconnectButton = document.querySelector('#disconnect');
let socket;
let peer;
let pendingCandidates = [];
function disconnect() {
if (socket?.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({ action: 'leave' }));
socket.close();
}
peer?.close();
peer = null;
socket = null;
pendingCandidates = [];
connectButton.disabled = false;
disconnectButton.disabled = true;
}
connectButton.onclick = () => {
const channel = document.querySelector('#channel').value.trim();
const token = document.querySelector('#token').value.trim();
if (!channel || !token) {
log('Renseigner le canal et le JWT.');
return;
}
let url;
try {
url = new URL(document.querySelector('#server').value);
url.protocol = url.protocol === 'https:' ? 'wss:' : 'ws:';
url.pathname = `/ws/rtc/${encodeURIComponent(channel)}`;
url.search = new URLSearchParams({ token }).toString();
peer = new RTCPeerConnection();
peer.createDataChannel('test'); // Produit une offre SDP sans demander l'accès au microphone.
socket = new WebSocket(url);
} catch (error) {
log(`Configuration invalide : ${error.message}`);
disconnect();
return;
}
connectButton.disabled = true;
disconnectButton.disabled = false;
peer.onicecandidate = ({ candidate }) => {
if (candidate && socket?.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({ action: 'ice-candidate', candidate: candidate.candidate }));
log('Candidat ICE local envoyé');
}
};
peer.oniceconnectionstatechange = () => log(`ICE : ${peer.iceConnectionState}`);
peer.onconnectionstatechange = () => log(`Connexion : ${peer.connectionState}`);
socket.onopen = async () => {
log('WebSocket ouvert');
try {
const offer = await peer.createOffer();
await peer.setLocalDescription(offer);
socket.send(JSON.stringify({ action: 'sdp-offer', sdp: peer.localDescription.sdp }));
log('Offre SDP envoyée');
} catch (error) {
log(`Offre impossible : ${error.message}`);
disconnect();
}
};
socket.onmessage = async ({ data }) => {
try {
const message = JSON.parse(data);
if (message.action === 'answer') {
await peer.setRemoteDescription({ type: 'answer', sdp: message.sdp });
log('Réponse SDP appliquée');
for (const candidate of pendingCandidates) await peer.addIceCandidate({ candidate, sdpMLineIndex: 0 });
pendingCandidates = [];
} else if (message.action === 'ice-candidate') {
if (peer.remoteDescription) await peer.addIceCandidate({ candidate: message.candidate, sdpMLineIndex: 0 });
else pendingCandidates.push(message.candidate);
log('Candidat ICE distant reçu');
} else if (message.action === 'error') {
log(`Erreur serveur : ${message.message}`);
} else {
log(`Message inconnu : ${message.action}`);
}
} catch (error) {
log(`Message impossible à traiter : ${error.message}`);
}
};
socket.onerror = () => log('Erreur WebSocket (vérifier le serveur, le canal et le JWT).');
socket.onclose = ({ code }) => {
log(`WebSocket fermé (${code})`);
disconnect();
};
};
disconnectButton.onclick = disconnect;
</script>
</body>
</html>
+17 -5
View File
@@ -30,7 +30,8 @@ pub async fn ws_entrypoint_handler(
let (tx, mut rx) = mpsc::unbounded_channel::<Message>(); let (tx, mut rx) = mpsc::unbounded_channel::<Message>();
let peer_connection = state.rtc.new_peer_connection(); 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 // Task pour envoyer les message au frontend depuis RTCClient
let mut send_task = tokio::spawn(async move { 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 client_clone = rtc_client.clone();
let mut recv_task = tokio::spawn(async move { let mut recv_task = tokio::spawn(async move {
while let Some(Ok(message)) = receiver.next().await { 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! { tokio::select! {
_ = (&mut send_task) => recv_task.abort(), _ = &mut send_task => {
_ = (&mut recv_task) => send_task.abort(), 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();
} }