use wasm_bindgen::prelude::*;
use wasm_bindgen_futures::JsFuture;
use web_sys::{
RtcPeerConnection, RtcDataChannel, RtcConfiguration, RtcIceServer,
RtcDataChannelInit, RtcPeerConnectionIceEvent, RtcSessionDescriptionInit,
RtcSdpType, MessageEvent, Event,
};
use js_sys::{Object, Reflect, Array, Promise};
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use crate::error::Result;
use super::BrowserPeer;
pub struct WebRtcTransport {
entity_id: String,
peer_connections: Arc<Mutex<HashMap<String, WebRtcConnection>>>,
ice_servers: Vec<super::IceServer>,
signaling_server: Option<String>,
data_channel_config: RtcDataChannelInit,
}
pub struct WebRtcConnection {
peer_id: String,
peer_connection: RtcPeerConnection,
data_channel: Option<RtcDataChannel>,
connection_state: WebRtcConnectionState,
ice_candidates: Vec<RtcIceCandidate>,
}
#[derive(Debug, Clone)]
pub struct RtcIceCandidate {
pub candidate: String,
pub sdp_mid: Option<String>,
pub sdp_m_line_index: Option<u16>,
}
#[derive(Debug, Clone, PartialEq)]
pub enum WebRtcConnectionState {
New,
Connecting,
Connected,
Disconnected,
Failed,
Closed,
}
impl WebRtcTransport {
pub async fn new(entity_id: String, ice_servers: Vec<super::IceServer>) -> Result<Self> {
let mut data_channel_config = RtcDataChannelInit::new();
data_channel_config.ordered(false); data_channel_config.max_retransmits(3);
Ok(Self {
entity_id,
peer_connections: Arc::new(Mutex::new(HashMap::new())),
ice_servers,
signaling_server: None,
data_channel_config,
})
}
pub fn set_signaling_server(&mut self, server: String) {
self.signaling_server = Some(server);
}
pub async fn connect_to_peer(&self, peer_id: String) -> Result<()> {
web_sys::console::log_1(&format!("Connecting to peer via WebRTC: {}", peer_id).into());
let rtc_config = self.create_rtc_configuration()?;
let peer_connection = RtcPeerConnection::new_with_configuration(&rtc_config)
.map_err(|e| anyhow::anyhow!("Failed to create peer connection: {:?}", e))?;
let data_channel = peer_connection.create_data_channel_with_data_channel_dict(
"synapse",
&self.data_channel_config,
);
self.setup_peer_connection_handlers(&peer_connection, &peer_id).await?;
self.setup_data_channel_handlers(&data_channel, &peer_id).await?;
let connection = WebRtcConnection {
peer_id: peer_id.clone(),
peer_connection,
data_channel: Some(data_channel),
connection_state: WebRtcConnectionState::New,
ice_candidates: Vec::new(),
};
{
let mut connections = self.peer_connections.lock().unwrap();
connections.insert(peer_id, connection);
}
Ok(())
}
pub async fn send_message(&self, peer_id: &str, message: &str) -> Result<String> {
let connections = self.peer_connections.lock().unwrap();
if let Some(connection) = connections.get(peer_id) {
if let Some(ref data_channel) = connection.data_channel {
if data_channel.ready_state() == web_sys::RtcDataChannelState::Open {
data_channel.send_with_str(message)
.map_err(|e| anyhow::anyhow!("Failed to send via WebRTC: {:?}", e))?;
web_sys::console::log_1(&format!("Sent message via WebRTC to {}", peer_id).into());
return Ok(format!("webrtc://{}@{}", peer_id, self.entity_id));
}
}
}
Err(anyhow::anyhow!("No open WebRTC connection to {}", peer_id))
}
pub async fn discover_peers(&self) -> Result<Vec<BrowserPeer>> {
web_sys::console::log_1(&"Discovering WebRTC peers...".into());
Ok(Vec::new())
}
pub async fn handle_offer(&self, peer_id: String, offer_sdp: String) -> Result<String> {
web_sys::console::log_1(&format!("Handling WebRTC offer from {}", peer_id).into());
let rtc_config = self.create_rtc_configuration()?;
let peer_connection = RtcPeerConnection::new_with_configuration(&rtc_config)
.map_err(|e| anyhow::anyhow!("Failed to create peer connection: {:?}", e))?;
self.setup_peer_connection_handlers(&peer_connection, &peer_id).await?;
let mut offer_desc = RtcSessionDescriptionInit::new(RtcSdpType::Offer);
offer_desc.sdp(&offer_sdp);
let set_remote_promise = peer_connection.set_remote_description(&offer_desc);
JsFuture::from(set_remote_promise).await
.map_err(|e| anyhow::anyhow!("Failed to set remote description: {:?}", e))?;
let answer_promise = peer_connection.create_answer();
let answer = JsFuture::from(answer_promise).await
.map_err(|e| anyhow::anyhow!("Failed to create answer: {:?}", e))?;
let set_local_promise = peer_connection.set_local_description(&answer.into());
JsFuture::from(set_local_promise).await
.map_err(|e| anyhow::anyhow!("Failed to set local description: {:?}", e))?;
let answer_sdp = peer_connection.local_description()
.ok_or_else(|| anyhow::anyhow!("No local description available"))?
.sdp();
let connection = WebRtcConnection {
peer_id: peer_id.clone(),
peer_connection,
data_channel: None, connection_state: WebRtcConnectionState::Connecting,
ice_candidates: Vec::new(),
};
{
let mut connections = self.peer_connections.lock().unwrap();
connections.insert(peer_id, connection);
}
Ok(answer_sdp)
}
pub async fn handle_ice_candidate(&self, peer_id: &str, candidate: RtcIceCandidate) -> Result<()> {
let connections = self.peer_connections.lock().unwrap();
if let Some(connection) = connections.get(peer_id) {
let ice_candidate_init = web_sys::RtcIceCandidateInit::new(&candidate.candidate);
if let Some(ref sdp_mid) = candidate.sdp_mid {
ice_candidate_init.set_sdp_mid(Some(sdp_mid));
}
if let Some(sdp_m_line_index) = candidate.sdp_m_line_index {
ice_candidate_init.set_sdp_m_line_index(Some(sdp_m_line_index));
}
let ice_candidate = web_sys::RtcIceCandidate::new(&ice_candidate_init)
.map_err(|e| anyhow::anyhow!("Failed to create ICE candidate: {:?}", e))?;
let add_candidate_promise = connection.peer_connection.add_ice_candidate_with_opt_rtc_ice_candidate(Some(&ice_candidate));
JsFuture::from(add_candidate_promise).await
.map_err(|e| anyhow::anyhow!("Failed to add ICE candidate: {:?}", e))?;
web_sys::console::log_1(&format!("Added ICE candidate for {}", peer_id).into());
}
Ok(())
}
pub async fn disconnect_from_peer(&self, peer_id: &str) -> Result<()> {
let mut connections = self.peer_connections.lock().unwrap();
if let Some(connection) = connections.remove(peer_id) {
if let Some(ref data_channel) = connection.data_channel {
data_channel.close();
}
connection.peer_connection.close();
web_sys::console::log_1(&format!("Disconnected from WebRTC peer {}", peer_id).into());
}
Ok(())
}
async fn setup_peer_connection_handlers(&self, peer_connection: &RtcPeerConnection, peer_id: &str) -> Result<()> {
let pc_clone = peer_connection.clone();
let peer_id_clone = peer_id.to_string();
let on_ice_candidate = Closure::wrap(Box::new(move |event: RtcPeerConnectionIceEvent| {
let candidate = event.candidate();
if let Some(candidate) = candidate {
web_sys::console::log_1(&format!("ICE candidate: {:?}", candidate).into());
}
}) as Box<dyn FnMut(RtcPeerConnectionIceEvent)>);
peer_connection.set_onicecandidate(Some(on_ice_candidate.as_ref().unchecked_ref()));
on_ice_candidate.forget();
let pc_clone2 = pc_clone.clone();
let peer_id_clone2 = peer_id_clone.clone();
let on_connection_state_change = Closure::wrap(Box::new(move || {
let state = pc_clone2.connection_state();
web_sys::console::log_1(&format!("Connection state changed to: {:?}", state).into());
match state {
web_sys::RtcPeerConnectionState::Connected => {
web_sys::console::log_1(&format!("Connected to peer: {}", peer_id_clone2).into());
}
web_sys::RtcPeerConnectionState::Failed => {
web_sys::console::log_1(&format!("Connection failed with peer: {}", peer_id_clone2).into());
}
web_sys::RtcPeerConnectionState::Disconnected => {
web_sys::console::log_1(&format!("Disconnected from peer: {}", peer_id_clone2).into());
}
_ => {}
}
}) as Box<dyn FnMut()>);
peer_connection.set_onconnectionstatechange(Some(on_connection_state_change.as_ref().unchecked_ref()));
on_connection_state_change.forget();
let peer_id_clone3 = peer_id_clone.clone();
let self_entity_id = self.entity_id.clone();
let peer_connections = self.peer_connections.clone();
let on_data_channel = Closure::wrap(Box::new(move |event: web_sys::RtcDataChannelEvent| {
web_sys::console::log_1(&format!("Data channel received from peer: {}", peer_id_clone3).into());
let data_channel = event.channel();
let dc_clone = data_channel.clone();
let peer_id_clone4 = peer_id_clone3.clone();
let on_message = Closure::wrap(Box::new(move |event: MessageEvent| {
if let Ok(data) = event.data().dyn_into::<js_sys::ArrayBuffer>() {
web_sys::console::log_1(&format!("Received array buffer from {}: {} bytes",
peer_id_clone4, data.byte_length()).into());
} else if let Ok(text) = event.data().dyn_into::<js_sys::JsString>() {
web_sys::console::log_1(&format!("Received text from {}: {}",
peer_id_clone4, text.as_string().unwrap_or_default()).into());
}
}) as Box<dyn FnMut(MessageEvent)>);
dc_clone.set_onmessage(Some(on_message.as_ref().unchecked_ref()));
on_message.forget();
if let Ok(mut peer_connections) = peer_connections.lock() {
if let Some(connection) = peer_connections.get_mut(&peer_id_clone3) {
connection.data_channel = Some(data_channel);
connection.connection_state = WebRtcConnectionState::Connected;
}
}
}) as Box<dyn FnMut(web_sys::RtcDataChannelEvent)>);
peer_connection.set_ondatachannel(Some(on_data_channel.as_ref().unchecked_ref()));
on_data_channel.forget();
Ok(())
}
async fn setup_data_channel_handlers(&self, data_channel: &RtcDataChannel, peer_id: &str) -> Result<()> {
let dc_clone = data_channel.clone();
let peer_id_clone = peer_id.to_string();
let on_open = Closure::wrap(Box::new(move |_event: Event| {
web_sys::console::log_1(&format!("Data channel opened with peer: {}", peer_id_clone).into());
let message = format!("{{\"type\":\"handshake\",\"entity_id\":\"{}\"}}", peer_id_clone);
if let Err(e) = dc_clone.send_with_str(&message) {
web_sys::console::error_1(&format!("Error sending handshake: {:?}", e).into());
}
}) as Box<dyn FnMut(Event)>);
data_channel.set_onopen(Some(on_open.as_ref().unchecked_ref()));
on_open.forget();
let dc_clone = data_channel.clone();
let peer_id_clone = peer_id.to_string();
let on_error = Closure::wrap(Box::new(move |event: Event| {
web_sys::console::error_1(&format!("Data channel error with peer {}: {:?}", peer_id_clone, event).into());
}) as Box<dyn FnMut(Event)>);
data_channel.set_onerror(Some(on_error.as_ref().unchecked_ref()));
on_error.forget();
let peer_id_clone = peer_id.to_string();
let peer_connections = self.peer_connections.clone();
let on_close = Closure::wrap(Box::new(move |_event: Event| {
web_sys::console::log_1(&format!("Data channel closed with peer: {}", peer_id_clone).into());
if let Ok(mut peer_connections) = peer_connections.lock() {
if let Some(connection) = peer_connections.get_mut(&peer_id_clone) {
connection.connection_state = WebRtcConnectionState::Closed;
connection.data_channel = None;
}
}
}) as Box<dyn FnMut(Event)>);
data_channel.set_onclose(Some(on_close.as_ref().unchecked_ref()));
on_close.forget();
let peer_id_clone = peer_id.to_string();
let on_message = Closure::wrap(Box::new(move |event: MessageEvent| {
if let Ok(data) = event.data().dyn_into::<js_sys::ArrayBuffer>() {
web_sys::console::log_1(&format!("Received array buffer from {}: {} bytes",
peer_id_clone, data.byte_length()).into());
} else if let Ok(text) = event.data().dyn_into::<js_sys::JsString>() {
web_sys::console::log_1(&format!("Received text from {}: {}",
peer_id_clone, text.as_string().unwrap_or_default()).into());
}
}) as Box<dyn FnMut(MessageEvent)>);
data_channel.set_onmessage(Some(on_message.as_ref().unchecked_ref()));
on_message.forget();
Ok(())
}
fn create_rtc_configuration(&self) -> Result<RtcConfiguration> {
let config = RtcConfiguration::new();
let ice_servers = Array::new();
for server in &self.ice_servers {
let ice_server = RtcIceServer::new();
Reflect::set(&ice_server, &JsValue::from_str("urls"), &JsValue::from_str(&server.url))?;
if let Some(username) = &server.username {
Reflect::set(&ice_server, &JsValue::from_str("username"), &JsValue::from_str(username))?;
}
if let Some(credential) = &server.credential {
Reflect::set(&ice_server, &JsValue::from_str("credential"), &JsValue::from_str(credential))?;
}
ice_servers.push(&ice_server);
}
Reflect::set(&config, &JsValue::from_str("iceServers"), &ice_servers)?;
Ok(config)
}
}