use crate::error::{Error, Result};
use crate::network::Message;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::{Duration, Instant};
#[cfg(feature = "webrtc")]
use async_tungstenite::tungstenite::Message as WsMessage;
#[cfg(feature = "webrtc")]
use futures_util::{SinkExt, StreamExt};
#[cfg(feature = "webrtc")]
use smol::channel::{Receiver, Sender};
#[cfg(feature = "webrtc")]
use smol::lock::RwLock;
#[cfg(feature = "webrtc")]
use webrtc::api::APIBuilder;
#[cfg(feature = "webrtc")]
use webrtc::data_channel::RTCDataChannel;
#[cfg(feature = "webrtc")]
use webrtc::ice_transport::ice_server::RTCIceServer;
#[cfg(feature = "webrtc")]
use webrtc::peer_connection::configuration::RTCConfiguration;
#[cfg(feature = "webrtc")]
use webrtc::peer_connection::RTCPeerConnection;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum SignalingMessage {
Offer {
from: String,
to: String,
sdp: String,
},
Answer {
from: String,
to: String,
sdp: String,
},
IceCandidate {
from: String,
to: String,
candidate: String,
sdp_mid: Option<String>,
sdp_mline_index: Option<u16>,
},
Join { peer_id: String },
Leave { peer_id: String },
}
#[derive(Debug, Clone)]
pub struct WebRtcConfig {
pub stun_server: String,
pub turn_server: Option<String>,
pub turn_username: Option<String>,
pub turn_credential: Option<String>,
pub signaling_port: u16,
pub ice_timeout: Duration,
pub channel_label: String,
}
impl Default for WebRtcConfig {
fn default() -> Self {
Self {
stun_server: "stun:stun.l.google.com:19302".to_string(),
turn_server: None,
turn_username: None,
turn_credential: None,
signaling_port: 19080,
ice_timeout: Duration::from_secs(30),
channel_label: "aingle".to_string(),
}
}
}
impl WebRtcConfig {
pub fn with_stun(stun_server: &str) -> Self {
Self {
stun_server: stun_server.to_string(),
..Default::default()
}
}
pub fn with_turn(mut self, server: &str, username: &str, credential: &str) -> Self {
self.turn_server = Some(server.to_string());
self.turn_username = Some(username.to_string());
self.turn_credential = Some(credential.to_string());
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ConnectionState {
New,
Connecting,
Connected,
Disconnected,
Failed,
Closed,
}
#[derive(Debug, Clone, Default)]
pub struct WebRtcStats {
pub messages_sent: u64,
pub messages_received: u64,
pub bytes_sent: u64,
pub bytes_received: u64,
pub rtt_ms: u32,
pub ice_candidates: u32,
pub using_relay: bool,
}
#[derive(Debug)]
pub struct PeerConnection {
pub peer_id: String,
pub state: ConnectionState,
pub remote_addr: Option<SocketAddr>,
pub stats: WebRtcStats,
pub connected_at: Option<Instant>,
}
impl PeerConnection {
pub fn new(peer_id: &str) -> Self {
Self {
peer_id: peer_id.to_string(),
state: ConnectionState::New,
remote_addr: None,
stats: WebRtcStats::default(),
connected_at: None,
}
}
pub fn is_connected(&self) -> bool {
self.state == ConnectionState::Connected
}
pub fn connection_duration(&self) -> Option<Duration> {
self.connected_at.map(|t| t.elapsed())
}
}
pub struct WebRtcServer {
config: WebRtcConfig,
peers: HashMap<String, PeerConnection>,
running: bool,
local_peer_id: String,
#[cfg(feature = "webrtc")]
api: Option<webrtc::api::API>,
#[cfg(feature = "webrtc")]
rtc_peers: HashMap<String, Arc<RTCPeerConnection>>,
#[cfg(feature = "webrtc")]
data_channels: HashMap<String, Arc<RTCDataChannel>>,
#[cfg(feature = "webrtc")]
message_rx: Option<Receiver<(String, Vec<u8>)>>,
#[cfg(feature = "webrtc")]
message_tx: Option<Sender<(String, Vec<u8>)>>,
#[cfg(feature = "webrtc")]
signaling_queue: Vec<SignalingMessage>,
}
impl WebRtcServer {
pub fn new(config: WebRtcConfig) -> Self {
let local_peer_id = Self::generate_peer_id();
Self {
config,
peers: HashMap::new(),
running: false,
local_peer_id,
#[cfg(feature = "webrtc")]
api: None,
#[cfg(feature = "webrtc")]
rtc_peers: HashMap::new(),
#[cfg(feature = "webrtc")]
data_channels: HashMap::new(),
#[cfg(feature = "webrtc")]
message_rx: None,
#[cfg(feature = "webrtc")]
message_tx: None,
#[cfg(feature = "webrtc")]
signaling_queue: Vec::new(),
}
}
fn generate_peer_id() -> String {
use rand::Rng;
let mut rng = rand::rng();
let bytes: [u8; 16] = rng.random();
bytes.iter().map(|b| format!("{:02x}", b)).collect()
}
pub fn local_peer_id(&self) -> &str {
&self.local_peer_id
}
pub async fn start(&mut self) -> Result<()> {
if self.running {
return Ok(());
}
log::info!(
"Starting WebRTC server on signaling port {}",
self.config.signaling_port
);
#[cfg(feature = "webrtc")]
{
let api = APIBuilder::new().build();
self.api = Some(api);
let (tx, rx) = smol::channel::unbounded();
self.message_tx = Some(tx);
self.message_rx = Some(rx);
log::info!("WebRTC API initialized");
}
#[cfg(not(feature = "webrtc"))]
{
log::warn!("WebRTC feature not enabled, using simulated mode");
}
self.running = true;
Ok(())
}
pub async fn stop(&mut self) -> Result<()> {
if !self.running {
return Ok(());
}
log::info!("Stopping WebRTC server");
#[cfg(feature = "webrtc")]
{
for (peer_id, pc) in self.rtc_peers.drain() {
if let Err(e) = pc.close().await {
log::warn!("Error closing peer connection {}: {}", peer_id, e);
}
}
self.data_channels.clear();
self.api = None;
}
for (peer_id, mut peer) in self.peers.drain() {
peer.state = ConnectionState::Closed;
log::debug!("Closed connection to peer: {}", peer_id);
}
self.running = false;
Ok(())
}
pub fn is_running(&self) -> bool {
self.running
}
pub async fn connect(&mut self, peer_id: &str) -> Result<()> {
if self.peers.contains_key(peer_id) {
return Err(Error::network(format!(
"Already connected to peer: {}",
peer_id
)));
}
let mut peer = PeerConnection::new(peer_id);
peer.state = ConnectionState::Connecting;
log::info!("Initiating WebRTC connection to peer: {}", peer_id);
#[cfg(feature = "webrtc")]
{
let api = self
.api
.as_ref()
.ok_or_else(|| Error::network("WebRTC API not initialized".to_string()))?;
let mut ice_servers = vec![RTCIceServer {
urls: vec![self.config.stun_server.clone()],
..Default::default()
}];
if let Some(turn) = &self.config.turn_server {
ice_servers.push(RTCIceServer {
urls: vec![turn.clone()],
username: self.config.turn_username.clone().unwrap_or_default(),
credential: self.config.turn_credential.clone().unwrap_or_default(),
..Default::default()
});
}
let rtc_config = RTCConfiguration {
ice_servers,
..Default::default()
};
let pc = api
.new_peer_connection(rtc_config)
.await
.map_err(|e| Error::network(format!("Failed to create peer connection: {}", e)))?;
let dc = pc
.create_data_channel(&self.config.channel_label, None)
.await
.map_err(|e| Error::network(format!("Failed to create data channel: {}", e)))?;
log::debug!(
"Created data channel '{}' for peer {}",
self.config.channel_label,
peer_id
);
let pc = Arc::new(pc);
self.rtc_peers.insert(peer_id.to_string(), pc);
self.data_channels.insert(peer_id.to_string(), dc);
log::info!("WebRTC peer connection created for: {}", peer_id);
}
self.peers.insert(peer_id.to_string(), peer);
Ok(())
}
pub async fn disconnect(&mut self, peer_id: &str) -> Result<()> {
#[cfg(feature = "webrtc")]
{
if let Some(pc) = self.rtc_peers.remove(peer_id) {
if let Err(e) = pc.close().await {
log::warn!("Error closing peer connection {}: {}", peer_id, e);
}
}
self.data_channels.remove(peer_id);
}
if let Some(mut peer) = self.peers.remove(peer_id) {
peer.state = ConnectionState::Closed;
log::info!("Disconnected from peer: {}", peer_id);
Ok(())
} else {
Err(Error::network(format!("Peer not found: {}", peer_id)))
}
}
pub async fn send(&mut self, peer_id: &str, message: &Message) -> Result<()> {
let peer = self
.peers
.get_mut(peer_id)
.ok_or_else(|| Error::network(format!("Peer not found: {}", peer_id)))?;
if peer.state != ConnectionState::Connected {
return Err(Error::network(format!("Peer not connected: {}", peer_id)));
}
let payload =
serde_json::to_vec(message).map_err(|e| Error::Serialization(e.to_string()))?;
#[cfg(feature = "webrtc")]
{
let dc = self
.data_channels
.get(peer_id)
.ok_or_else(|| Error::network(format!("Data channel not found: {}", peer_id)))?;
dc.send(&bytes::Bytes::from(payload.clone()))
.await
.map_err(|e| Error::network(format!("Failed to send: {}", e)))?;
}
peer.stats.messages_sent += 1;
peer.stats.bytes_sent += payload.len() as u64;
log::debug!("Sent message to peer {} ({} bytes)", peer_id, payload.len());
Ok(())
}
pub async fn recv(&mut self) -> Result<Option<(String, Message)>> {
#[cfg(feature = "webrtc")]
{
if let Some(ref rx) = self.message_rx {
match rx.try_recv() {
Ok((peer_id, data)) => {
let message: Message = serde_json::from_slice(&data)
.map_err(|e| Error::Serialization(e.to_string()))?;
if let Some(peer) = self.peers.get_mut(&peer_id) {
peer.stats.messages_received += 1;
peer.stats.bytes_received += data.len() as u64;
}
log::debug!(
"Received message from peer {} ({} bytes)",
peer_id,
data.len()
);
return Ok(Some((peer_id, message)));
}
Err(smol::channel::TryRecvError::Empty) => {
return Ok(None);
}
Err(smol::channel::TryRecvError::Closed) => {
return Err(Error::network("Message channel closed".to_string()));
}
}
}
}
Ok(None)
}
pub fn peer_count(&self) -> usize {
self.peers
.values()
.filter(|p| p.state == ConnectionState::Connected)
.count()
}
pub fn peer_ids(&self) -> Vec<String> {
self.peers.keys().cloned().collect()
}
pub fn get_peer(&self, peer_id: &str) -> Option<&PeerConnection> {
self.peers.get(peer_id)
}
pub fn stats(&self) -> WebRtcStats {
let mut stats = WebRtcStats::default();
for peer in self.peers.values() {
stats.messages_sent += peer.stats.messages_sent;
stats.messages_received += peer.stats.messages_received;
stats.bytes_sent += peer.stats.bytes_sent;
stats.bytes_received += peer.stats.bytes_received;
}
stats
}
#[cfg(feature = "webrtc")]
pub fn queue_signaling(&mut self, message: SignalingMessage) {
self.signaling_queue.push(message);
}
#[cfg(feature = "webrtc")]
pub fn drain_signaling_queue(&mut self) -> Vec<SignalingMessage> {
std::mem::take(&mut self.signaling_queue)
}
#[cfg(feature = "webrtc")]
pub async fn handle_signaling(&mut self, message: SignalingMessage) -> Result<()> {
match message {
SignalingMessage::Offer { from, to, sdp } => {
if to == self.local_peer_id {
log::debug!("Received offer from {}: {} bytes", from, sdp.len());
}
}
SignalingMessage::Answer { from, to, sdp } => {
if to == self.local_peer_id {
log::debug!("Received answer from {}: {} bytes", from, sdp.len());
}
}
SignalingMessage::IceCandidate {
from,
to,
candidate,
sdp_mid,
sdp_mline_index,
} => {
if to == self.local_peer_id {
log::debug!(
"Received ICE candidate from {}: mid={:?}, index={:?}",
from,
sdp_mid,
sdp_mline_index
);
let _ = candidate; }
}
SignalingMessage::Join { peer_id } => {
log::info!("Peer joined: {}", peer_id);
}
SignalingMessage::Leave { peer_id } => {
log::info!("Peer left: {}", peer_id);
let _ = self.disconnect(&peer_id).await;
}
}
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct SignalingConfig {
pub bind_addr: String,
pub max_connections: usize,
pub heartbeat_interval: Duration,
pub connection_timeout: Duration,
}
impl Default for SignalingConfig {
fn default() -> Self {
Self {
bind_addr: "0.0.0.0:19080".to_string(),
max_connections: 100,
heartbeat_interval: Duration::from_secs(30),
connection_timeout: Duration::from_secs(60),
}
}
}
#[cfg(feature = "webrtc")]
#[derive(Debug)]
#[allow(dead_code)] struct ConnectedPeer {
peer_id: String,
tx: Sender<SignalingMessage>,
last_activity: Instant,
}
#[cfg(feature = "webrtc")]
pub struct SignalingServer {
config: SignalingConfig,
clients: Arc<RwLock<HashMap<String, ConnectedPeer>>>,
running: Arc<std::sync::atomic::AtomicBool>,
server_task: Option<smol::Task<()>>,
}
#[cfg(feature = "webrtc")]
impl SignalingServer {
pub fn new(config: SignalingConfig) -> Self {
Self {
config,
clients: Arc::new(RwLock::new(HashMap::new())),
running: Arc::new(std::sync::atomic::AtomicBool::new(false)),
server_task: None,
}
}
pub async fn start(&mut self) -> Result<()> {
if self.running.load(std::sync::atomic::Ordering::SeqCst) {
return Ok(());
}
let addr = &self.config.bind_addr;
log::info!("Starting signaling server on {}", addr);
let listener = smol::net::TcpListener::bind(addr)
.await
.map_err(|e| Error::network(format!("Failed to bind signaling server: {}", e)))?;
self.running
.store(true, std::sync::atomic::Ordering::SeqCst);
let clients = self.clients.clone();
let running = self.running.clone();
let max_connections = self.config.max_connections;
let task = smol::spawn(async move {
log::info!("Signaling server listening for connections");
while running.load(std::sync::atomic::Ordering::SeqCst) {
match listener.accept().await {
Ok((stream, addr)) => {
let client_count = clients.read().await.len();
if client_count >= max_connections {
log::warn!(
"Max connections reached ({}), rejecting {}",
max_connections,
addr
);
continue;
}
log::debug!("New signaling connection from {}", addr);
let clients = clients.clone();
smol::spawn(async move {
if let Err(e) = Self::handle_connection(stream, addr, clients).await {
log::warn!("Connection error from {}: {}", addr, e);
}
})
.detach();
}
Err(e) => {
log::error!("Accept error: {}", e);
}
}
}
log::info!("Signaling server stopped");
});
self.server_task = Some(task);
Ok(())
}
async fn handle_connection(
stream: smol::net::TcpStream,
addr: SocketAddr,
clients: Arc<RwLock<HashMap<String, ConnectedPeer>>>,
) -> Result<()> {
let ws_stream = async_tungstenite::accept_async(stream)
.await
.map_err(|e| Error::network(format!("WebSocket upgrade failed: {}", e)))?;
let (mut ws_sink, mut ws_stream) = ws_stream.split();
let peer_id = loop {
match ws_stream.next().await {
Some(Ok(WsMessage::Text(text))) => {
match serde_json::from_str::<SignalingMessage>(&text) {
Ok(SignalingMessage::Join { peer_id }) => break peer_id,
Ok(_) => {
log::warn!("Expected Join message from {}", addr);
}
Err(e) => {
log::warn!("Invalid message from {}: {}", addr, e);
}
}
}
Some(Ok(WsMessage::Close(_))) | None => {
return Ok(());
}
_ => continue,
}
};
log::info!("Peer '{}' joined from {}", peer_id, addr);
let (tx, rx) = smol::channel::bounded(32);
{
let mut clients_guard = clients.write().await;
clients_guard.insert(
peer_id.clone(),
ConnectedPeer {
peer_id: peer_id.clone(),
tx: tx.clone(),
last_activity: Instant::now(),
},
);
let join_msg = SignalingMessage::Join {
peer_id: peer_id.clone(),
};
for (id, client) in clients_guard.iter() {
if id != &peer_id {
let _ = client.tx.try_send(join_msg.clone());
}
}
}
let forward_task = smol::spawn(async move {
while let Ok(msg) = rx.recv().await {
let json = match serde_json::to_string(&msg) {
Ok(j) => j,
Err(e) => {
log::warn!("Failed to serialize message: {}", e);
continue;
}
};
if ws_sink.send(WsMessage::Text(json)).await.is_err() {
break;
}
}
});
while let Some(msg_result) = ws_stream.next().await {
match msg_result {
Ok(WsMessage::Text(text)) => {
match serde_json::from_str::<SignalingMessage>(&text) {
Ok(msg) => {
Self::route_message(&clients, &peer_id, msg).await;
}
Err(e) => {
log::warn!("Invalid message from {}: {}", peer_id, e);
}
}
}
Ok(WsMessage::Ping(data)) => {
if let Some(client) = clients.write().await.get_mut(&peer_id) {
client.last_activity = Instant::now();
}
let _ = data; }
Ok(WsMessage::Close(_)) => {
log::info!("Peer '{}' disconnected", peer_id);
break;
}
Err(e) => {
log::warn!("WebSocket error from {}: {}", peer_id, e);
break;
}
_ => {}
}
}
forward_task.cancel().await;
{
let mut clients_guard = clients.write().await;
clients_guard.remove(&peer_id);
let leave_msg = SignalingMessage::Leave {
peer_id: peer_id.clone(),
};
for client in clients_guard.values() {
let _ = client.tx.try_send(leave_msg.clone());
}
}
log::info!("Peer '{}' cleanup complete", peer_id);
Ok(())
}
async fn route_message(
clients: &Arc<RwLock<HashMap<String, ConnectedPeer>>>,
from: &str,
message: SignalingMessage,
) {
let target_peer_id = match &message {
SignalingMessage::Offer { to, .. } => Some(to.clone()),
SignalingMessage::Answer { to, .. } => Some(to.clone()),
SignalingMessage::IceCandidate { to, .. } => Some(to.clone()),
SignalingMessage::Join { .. } | SignalingMessage::Leave { .. } => None,
};
if let Some(target) = target_peer_id {
let clients_guard = clients.read().await;
if let Some(client) = clients_guard.get(&target) {
if let Err(e) = client.tx.try_send(message) {
log::warn!("Failed to route message from {} to {}: {}", from, target, e);
}
} else {
log::warn!(
"Target peer '{}' not found for message from '{}'",
target,
from
);
}
}
}
pub async fn stop(&mut self) -> Result<()> {
if !self.running.load(std::sync::atomic::Ordering::SeqCst) {
return Ok(());
}
log::info!("Stopping signaling server");
self.running
.store(false, std::sync::atomic::Ordering::SeqCst);
if let Some(task) = self.server_task.take() {
task.cancel().await;
}
self.clients.write().await.clear();
Ok(())
}
pub fn is_running(&self) -> bool {
self.running.load(std::sync::atomic::Ordering::SeqCst)
}
pub async fn client_count(&self) -> usize {
self.clients.read().await.len()
}
pub async fn peer_ids(&self) -> Vec<String> {
self.clients.read().await.keys().cloned().collect()
}
}
#[cfg(feature = "webrtc")]
pub struct SignalingClient {
server_url: String,
peer_id: String,
tx: Option<Sender<SignalingMessage>>,
rx: Option<Receiver<SignalingMessage>>,
connected: bool,
}
#[cfg(feature = "webrtc")]
impl SignalingClient {
pub fn new(server_url: &str, peer_id: &str) -> Self {
Self {
server_url: server_url.to_string(),
peer_id: peer_id.to_string(),
tx: None,
rx: None,
connected: false,
}
}
pub async fn connect(&mut self) -> Result<()> {
if self.connected {
return Ok(());
}
log::info!("Connecting to signaling server: {}", self.server_url);
let (ws_stream, _) = async_tungstenite::async_std::connect_async(&self.server_url)
.await
.map_err(|e| Error::network(format!("Failed to connect to signaling server: {}", e)))?;
let (mut ws_sink, mut ws_stream) = ws_stream.split();
let join_msg = SignalingMessage::Join {
peer_id: self.peer_id.clone(),
};
let json =
serde_json::to_string(&join_msg).map_err(|e| Error::Serialization(e.to_string()))?;
ws_sink
.send(WsMessage::Text(json))
.await
.map_err(|e| Error::network(format!("Failed to send Join: {}", e)))?;
let (out_tx, out_rx) = smol::channel::bounded::<SignalingMessage>(32);
let (in_tx, in_rx) = smol::channel::bounded::<SignalingMessage>(32);
self.tx = Some(out_tx);
self.rx = Some(in_rx);
smol::spawn(async move {
while let Ok(msg) = out_rx.recv().await {
let json = match serde_json::to_string(&msg) {
Ok(j) => j,
Err(_) => continue,
};
if ws_sink.send(WsMessage::Text(json)).await.is_err() {
break;
}
}
})
.detach();
smol::spawn(async move {
while let Some(Ok(msg)) = ws_stream.next().await {
if let WsMessage::Text(text) = msg {
if let Ok(signaling_msg) = serde_json::from_str::<SignalingMessage>(&text) {
if in_tx.send(signaling_msg).await.is_err() {
break;
}
}
}
}
})
.detach();
self.connected = true;
log::info!("Connected to signaling server as '{}'", self.peer_id);
Ok(())
}
pub async fn send(&self, message: SignalingMessage) -> Result<()> {
let tx = self
.tx
.as_ref()
.ok_or_else(|| Error::network("Not connected".to_string()))?;
tx.send(message)
.await
.map_err(|e| Error::network(format!("Failed to send: {}", e)))
}
pub async fn recv(&self) -> Result<Option<SignalingMessage>> {
let rx = self
.rx
.as_ref()
.ok_or_else(|| Error::network("Not connected".to_string()))?;
match rx.try_recv() {
Ok(msg) => Ok(Some(msg)),
Err(smol::channel::TryRecvError::Empty) => Ok(None),
Err(smol::channel::TryRecvError::Closed) => {
Err(Error::network("Channel closed".to_string()))
}
}
}
pub fn is_connected(&self) -> bool {
self.connected
}
pub fn peer_id(&self) -> &str {
&self.peer_id
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_webrtc_config_default() {
let config = WebRtcConfig::default();
assert!(config.stun_server.contains("stun"));
assert!(config.turn_server.is_none());
assert_eq!(config.signaling_port, 19080);
}
#[test]
fn test_webrtc_config_with_stun() {
let config = WebRtcConfig::with_stun("stun:custom.server:3478");
assert_eq!(config.stun_server, "stun:custom.server:3478");
}
#[test]
fn test_webrtc_config_with_turn() {
let config = WebRtcConfig::default().with_turn("turn:relay.server:3478", "user", "pass");
assert!(config.turn_server.is_some());
assert_eq!(config.turn_username, Some("user".to_string()));
}
#[test]
fn test_peer_connection_new() {
let peer = PeerConnection::new("test-peer");
assert_eq!(peer.peer_id, "test-peer");
assert_eq!(peer.state, ConnectionState::New);
assert!(!peer.is_connected());
}
#[test]
fn test_peer_connection_connected() {
let mut peer = PeerConnection::new("test-peer");
peer.state = ConnectionState::Connected;
peer.connected_at = Some(Instant::now());
assert!(peer.is_connected());
assert!(peer.connection_duration().is_some());
}
#[test]
fn test_webrtc_server_creation() {
let config = WebRtcConfig::default();
let server = WebRtcServer::new(config);
assert!(!server.is_running());
assert_eq!(server.peer_count(), 0);
assert!(!server.local_peer_id().is_empty());
}
#[test]
fn test_connection_state_equality() {
assert_eq!(ConnectionState::New, ConnectionState::New);
assert_ne!(ConnectionState::Connected, ConnectionState::Disconnected);
}
#[test]
fn test_webrtc_stats_default() {
let stats = WebRtcStats::default();
assert_eq!(stats.messages_sent, 0);
assert_eq!(stats.bytes_sent, 0);
assert!(!stats.using_relay);
}
#[test]
fn test_signaling_config_default() {
let config = SignalingConfig::default();
assert_eq!(config.bind_addr, "0.0.0.0:19080");
assert_eq!(config.max_connections, 100);
assert_eq!(config.heartbeat_interval, Duration::from_secs(30));
}
#[test]
fn test_signaling_message_serialization() {
let offer = SignalingMessage::Offer {
from: "peer-a".to_string(),
to: "peer-b".to_string(),
sdp: "v=0\r\n...".to_string(),
};
let json = serde_json::to_string(&offer).unwrap();
let parsed: SignalingMessage = serde_json::from_str(&json).unwrap();
if let SignalingMessage::Offer { from, to, sdp } = parsed {
assert_eq!(from, "peer-a");
assert_eq!(to, "peer-b");
assert!(sdp.contains("v=0"));
} else {
panic!("Expected Offer message");
}
}
#[test]
fn test_signaling_message_ice_candidate() {
let ice = SignalingMessage::IceCandidate {
from: "peer-a".to_string(),
to: "peer-b".to_string(),
candidate: "candidate:1 1 UDP 2130706431 192.168.1.1 54321 typ host".to_string(),
sdp_mid: Some("0".to_string()),
sdp_mline_index: Some(0),
};
let json = serde_json::to_string(&ice).unwrap();
assert!(json.contains("IceCandidate"));
assert!(json.contains("candidate:1"));
}
#[test]
fn test_signaling_message_join_leave() {
let join = SignalingMessage::Join {
peer_id: "test-peer".to_string(),
};
let json = serde_json::to_string(&join).unwrap();
assert!(json.contains("test-peer"));
let leave = SignalingMessage::Leave {
peer_id: "test-peer".to_string(),
};
let json = serde_json::to_string(&leave).unwrap();
assert!(json.contains("Leave"));
}
}