Kanidm group-gated rooms and minimal mesh calling
Rooms can now require a Kanidm group (via the `groups` OIDC claim, mapped by `oauth2 update-claim-map` server-side) - dev/ops require `developers`, enforced at every message path (send, history, SSE). Adds a minimal WebRTC mesh call feature scoped to the lobby room, signaled over a separate `call.room.*` NATS subject kept out of the chat archive: public STUN only, no TURN, no SFU - small groups on friendly networks, by design.
This commit is contained in:
+328
@@ -0,0 +1,328 @@
|
||||
//! 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<Option<MediaStream>>,
|
||||
}
|
||||
|
||||
/// 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<Option<MediaStream>>,
|
||||
peers: StoredValue<HashMap<String, Peer>, LocalStorage>,
|
||||
pub in_call: RwSignal<bool>,
|
||||
}
|
||||
|
||||
fn new_peer_connection() -> Result<RtcPeerConnection, JsValue> {
|
||||
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<MediaStream, JsValue> {
|
||||
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::<MediaStream>()
|
||||
}
|
||||
|
||||
fn attach_local_tracks(pc: &RtcPeerConnection, stream: &MediaStream) {
|
||||
for track in stream.get_tracks().iter() {
|
||||
if let Ok(track) = track.dyn_into::<web_sys::MediaStreamTrack>() {
|
||||
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<String>`, 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<Option<MediaStream>> {
|
||||
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<Option<MediaStream>>)> {
|
||||
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::<web_sys::MediaStreamTrack>() {
|
||||
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<String>, 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<RtcPeerConnection> {
|
||||
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::<MediaStream>);
|
||||
{
|
||||
let remote_stream = remote_stream;
|
||||
let ontrack = Closure::<dyn FnMut(web_sys::RtcTrackEvent)>::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::<dyn FnMut(web_sys::RtcPeerConnectionIceEvent)>::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;
|
||||
});
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user