//! Minimal mesh WebRTC calling, browser-only. Public STUN, no TURN - calls //! across hostile NATs (symmetric NAT, restrictive corporate networks) //! simply won't connect. That's a known, accepted limitation, not a bug to //! fix later: real NAT traversal needs a TURN relay, which is real //! infrastructure this pass deliberately isn't standing up. Mesh topology //! (every pair of peers connects directly) is fine at the ~4-person scale //! this is scoped for; it does not scale further than that. #![cfg(feature = "hydrate")] use std::collections::HashMap; use js_sys::{Array, Reflect}; use leptos::prelude::*; use wasm_bindgen::{prelude::*, JsCast}; use wasm_bindgen_futures::JsFuture; use web_sys::{ MediaStream, MediaStreamConstraints, RtcConfiguration, RtcIceCandidateInit, RtcIceServer, RtcPeerConnection, RtcSdpType, RtcSessionDescriptionInit, }; use crate::call::{send_signal, CallSignalKind}; const STUN_URL: &str = "stun:stun.l.google.com:19302"; /// One remote participant: their peer connection plus the remote stream /// their video tile renders once `ontrack` fires. struct Peer { conn: RtcPeerConnection, stream: RwSignal>, } /// Call state for one room. Lives for as long as the user is in the call; /// dropped (and everything torn down) on "leave". #[derive(Clone)] pub struct CallState { room: String, me: String, local_stream: RwSignal>, peers: StoredValue, LocalStorage>, pub in_call: RwSignal, } fn new_peer_connection() -> Result { let config = RtcConfiguration::new(); let ice_server = RtcIceServer::new(); ice_server.set_urls(&JsValue::from_str(STUN_URL)); let servers = Array::new(); servers.push(&ice_server); config.set_ice_servers(&servers); RtcPeerConnection::new_with_configuration(&config) } async fn get_local_stream() -> Result { let window = web_sys::window().ok_or("no window")?; let media_devices = window.navigator().media_devices()?; let constraints = MediaStreamConstraints::new(); constraints.set_video(&JsValue::TRUE); constraints.set_audio(&JsValue::TRUE); let promise = media_devices.get_user_media_with_constraints(&constraints)?; let stream = JsFuture::from(promise).await?; stream.dyn_into::() } fn attach_local_tracks(pc: &RtcPeerConnection, stream: &MediaStream) { for track in stream.get_tracks().iter() { if let Ok(track) = track.dyn_into::() { pc.add_track_0(&track, stream); } } } /// Reads the `sdp` field off whatever `create_offer`/`create_answer` /// resolved to, and builds a fresh `RtcSessionDescriptionInit` from it - /// simpler and more reliable than trying to cast the resolved JsValue /// directly, since its concrete type varies by browser. Returns the sdp /// string alongside the desc (not `desc.get_sdp()` afterwards - that /// returns `Option`, and we already have it as a plain `String` /// right here). fn session_description_from_resolved( resolved: &JsValue, sdp_type: RtcSdpType, ) -> Result<(RtcSessionDescriptionInit, String), JsValue> { let sdp = Reflect::get(resolved, &JsValue::from_str("sdp"))? .as_string() .ok_or("resolved session description had no sdp field")?; let desc = RtcSessionDescriptionInit::new(sdp_type); desc.set_sdp(&sdp); Ok((desc, sdp)) } impl CallState { pub fn new(room: String, me: String) -> Self { Self { room, me, local_stream: RwSignal::new(None), peers: StoredValue::new_local(HashMap::new()), in_call: RwSignal::new(false), } } pub fn local_stream(&self) -> ReadSignal> { self.local_stream.read_only() } /// Streams for currently-known peers, keyed by their peer id (username). /// Recomputed each call - fine at mesh scale (~4 peers). pub fn peer_streams(&self) -> Vec<(String, RwSignal>)> { self.peers .with_value(|p| p.iter().map(|(id, peer)| (id.clone(), peer.stream)).collect()) } /// getUserMedia, then broadcast Join so existing participants know to /// offer us a connection. pub async fn join(&self) { match get_local_stream().await { Ok(stream) => self.local_stream.set(Some(stream)), Err(e) => { leptos::logging::error!("getUserMedia failed: {e:?}"); return; } } self.in_call.set(true); let _ = send_signal(self.room.clone(), None, CallSignalKind::Join).await; } /// Tears down every peer connection, stops all local tracks (releases /// the camera/mic), and tells the room we're gone. pub async fn leave(&self) { self.peers.update_value(|peers| { for (_, peer) in peers.drain() { peer.conn.close(); } }); if let Some(stream) = self.local_stream.get_untracked() { for track in stream.get_tracks().iter() { if let Ok(track) = track.dyn_into::() { track.stop(); } } } self.local_stream.set(None); self.in_call.set(false); let _ = send_signal(self.room.clone(), None, CallSignalKind::Leave).await; } /// One incoming signal from `/call-sse/{room}`. Ignores our own /// broadcasts and anything not addressed to us (directed messages are /// broadcast NATS-wide and filtered client-side - see call.rs). pub fn handle_signal(&self, from: String, to: Option, kind: CallSignalKind) { if from == self.me || !self.in_call.get_untracked() { return; } if let Some(to) = &to { if *to != self.me { return; } } match kind { CallSignalKind::Join => { // A new peer announced themselves - we initiate the offer. self.start_offer(from); } CallSignalKind::Leave => { self.peers.update_value(|peers| { if let Some(peer) = peers.remove(&from) { peer.conn.close(); } }); } CallSignalKind::Offer(sdp) => self.handle_offer(from, sdp), CallSignalKind::Answer(sdp) => self.handle_answer(from, sdp), CallSignalKind::IceCandidate(candidate_json) => { self.handle_ice_candidate(from, candidate_json) } } } fn ensure_peer(&self, peer_id: &str) -> Option { if let Some(pc) = self .peers .with_value(|peers| peers.get(peer_id).map(|p| p.conn.clone())) { return Some(pc); } let pc = new_peer_connection().ok()?; let Some(local) = self.local_stream.get_untracked() else { return None; }; attach_local_tracks(&pc, &local); let remote_stream = RwSignal::new(None::); { let remote_stream = remote_stream; let ontrack = Closure::::new(move |ev: web_sys::RtcTrackEvent| { remote_stream.set(Some(ev.streams().get(0).dyn_into().unwrap())); }); pc.set_ontrack(Some(ontrack.as_ref().unchecked_ref())); ontrack.forget(); } { let room = self.room.clone(); let peer_id = peer_id.to_string(); let onicecandidate = Closure::::new(move |ev: web_sys::RtcPeerConnectionIceEvent| { let Some(candidate) = ev.candidate() else { return; }; let Ok(candidate_json) = js_sys::JSON::stringify(&candidate.to_json()) .map(|s| s.as_string().unwrap_or_default()) else { return; }; let room = room.clone(); let peer_id = peer_id.clone(); leptos::task::spawn_local(async move { let _ = send_signal( room, Some(peer_id), CallSignalKind::IceCandidate(candidate_json), ) .await; }); }); pc.set_onicecandidate(Some(onicecandidate.as_ref().unchecked_ref())); onicecandidate.forget(); } self.peers.update_value(|peers| { peers.insert( peer_id.to_string(), Peer { conn: pc.clone(), stream: remote_stream, }, ); }); Some(pc) } fn start_offer(&self, peer_id: String) { let Some(pc) = self.ensure_peer(&peer_id) else { return; }; let room = self.room.clone(); leptos::task::spawn_local(async move { let Ok(resolved) = JsFuture::from(pc.create_offer()).await else { return; }; let Ok((desc, sdp)) = session_description_from_resolved(&resolved, RtcSdpType::Offer) else { return; }; if JsFuture::from(pc.set_local_description(&desc)).await.is_err() { return; } let _ = send_signal(room, Some(peer_id), CallSignalKind::Offer(sdp)).await; }); } fn handle_offer(&self, from: String, sdp: String) { let Some(pc) = self.ensure_peer(&from) else { return; }; let room = self.room.clone(); leptos::task::spawn_local(async move { let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Offer); remote_desc.set_sdp(&sdp); if JsFuture::from(pc.set_remote_description(&remote_desc)) .await .is_err() { return; } let Ok(resolved) = JsFuture::from(pc.create_answer()).await else { return; }; let Ok((answer_desc, answer_sdp)) = session_description_from_resolved(&resolved, RtcSdpType::Answer) else { return; }; if JsFuture::from(pc.set_local_description(&answer_desc)) .await .is_err() { return; } let _ = send_signal(room, Some(from), CallSignalKind::Answer(answer_sdp)).await; }); } fn handle_answer(&self, from: String, sdp: String) { let Some(pc) = self .peers .with_value(|peers| peers.get(&from).map(|p| p.conn.clone())) else { return; }; leptos::task::spawn_local(async move { let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Answer); remote_desc.set_sdp(&sdp); let _ = JsFuture::from(pc.set_remote_description(&remote_desc)).await; }); } fn handle_ice_candidate(&self, from: String, candidate_json: String) { let Some(pc) = self .peers .with_value(|peers| peers.get(&from).map(|p| p.conn.clone())) else { return; }; let Ok(parsed) = js_sys::JSON::parse(&candidate_json) else { return; }; let init: RtcIceCandidateInit = parsed.unchecked_into(); leptos::task::spawn_local(async move { let _ = JsFuture::from( pc.add_ice_candidate_with_opt_rtc_ice_candidate_init(Some(&init)), ) .await; }); } }