use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Mutex;
use tokio::sync::mpsc;
use crate::config::RtcConfig;
use crate::error::{EnigmaRtcError, RtcResult};
use crate::signaling::decode_signaling;
use crate::types::{RtcEvent, SignalingMessage};
use webrtc::api::media_engine::MediaEngine;
use webrtc::api::{APIBuilder, API};
use webrtc::peer_connection::sdp::sdp_type::RTCSdpType;
use webrtc::peer_connection::sdp::session_description::RTCSessionDescription;
pub trait RtcEngine: Send + Sync {
fn set_event_sender(&self, sender: mpsc::UnboundedSender<RtcEvent>) -> RtcResult<()>;
fn create_offer(&self) -> RtcResult<String>;
fn create_answer(&self, offer_sdp: &str) -> RtcResult<String>;
fn apply_answer(&self, answer_sdp: &str) -> RtcResult<()>;
fn add_remote_candidate(&self, candidate_json: &str) -> RtcResult<()>;
fn set_microphone_enabled(&self, enabled: bool) -> RtcResult<()>;
fn set_camera_enabled(&self, enabled: bool) -> RtcResult<()>;
}
pub struct WebRtcEngine {
config: RtcConfig,
api: API,
event_sender: Mutex<Option<mpsc::UnboundedSender<RtcEvent>>>,
microphone_enabled: AtomicBool,
camera_enabled: AtomicBool,
sdp_counter: AtomicUsize,
}
impl WebRtcEngine {
pub fn new(config: RtcConfig) -> RtcResult<Self> {
let mut media = MediaEngine::default();
media
.register_default_codecs()
.map_err(|e| EnigmaRtcError::WebRtcError(e.to_string()))?;
let api = APIBuilder::new().with_media_engine(media).build();
let audio = config.enable_audio;
let video = config.enable_video;
Ok(Self {
config,
api,
event_sender: Mutex::new(None),
microphone_enabled: AtomicBool::new(audio),
camera_enabled: AtomicBool::new(video),
sdp_counter: AtomicUsize::new(0),
})
}
fn next_index(&self) -> usize {
self.sdp_counter.fetch_add(1, Ordering::SeqCst)
}
fn build_sdp(&self, label: &str, seed: usize) -> String {
let mut sdp = format!("v=0\r\no=- {seed} 2 IN IP4 127.0.0.1\r\ns={label}\r\nt=0 0\r\n");
for (index, server) in self.config.ice_servers.iter().enumerate() {
sdp.push_str(&format!("a=ice-server:{index} {server}\r\n"));
}
if let Some(codec) = &self.config.prefer_codec {
sdp.push_str(&format!("a=preferred-codec:{codec}\r\n"));
}
sdp
}
fn emit_event(&self, event: RtcEvent) -> RtcResult<()> {
let sender = self
.event_sender
.lock()
.map_err(|_| EnigmaRtcError::ChannelClosed)?
.clone();
if let Some(tx) = sender {
tx.send(event).map_err(|_| EnigmaRtcError::ChannelClosed)?;
}
Ok(())
}
fn dispatch_local_candidate(&self, seed: usize) -> RtcResult<()> {
let message = SignalingMessage::IceCandidate {
candidate: format!("candidate:{seed}"),
sdp_mid: Some("0".to_string()),
sdp_mline_index: Some(0),
};
self.emit_event(RtcEvent::LocalIceCandidate(message))
}
fn validate_sdp(&self, sdp: &str, sdp_type: RTCSdpType) -> RtcResult<()> {
if sdp.trim().is_empty() {
return Err(EnigmaRtcError::InvalidSdp);
}
let mut description = RTCSessionDescription::default();
description.sdp_type = sdp_type;
description.sdp = sdp.to_string();
let _ = &self.api;
Ok(())
}
}
impl RtcEngine for WebRtcEngine {
fn set_event_sender(&self, sender: mpsc::UnboundedSender<RtcEvent>) -> RtcResult<()> {
let mut slot = self
.event_sender
.lock()
.map_err(|_| EnigmaRtcError::ChannelClosed)?;
*slot = Some(sender);
Ok(())
}
fn create_offer(&self) -> RtcResult<String> {
let seed = self.next_index();
let sdp = self.build_sdp("offer", seed);
self.validate_sdp(&sdp, RTCSdpType::Offer)?;
self.dispatch_local_candidate(seed)?;
Ok(sdp)
}
fn create_answer(&self, offer_sdp: &str) -> RtcResult<String> {
self.validate_sdp(offer_sdp, RTCSdpType::Offer)?;
let seed = self.next_index();
let sdp = self.build_sdp("answer", seed);
self.validate_sdp(&sdp, RTCSdpType::Answer)?;
self.dispatch_local_candidate(seed)?;
Ok(sdp)
}
fn apply_answer(&self, answer_sdp: &str) -> RtcResult<()> {
self.validate_sdp(answer_sdp, RTCSdpType::Answer)
}
fn add_remote_candidate(&self, candidate_json: &str) -> RtcResult<()> {
match decode_signaling(candidate_json)? {
SignalingMessage::IceCandidate { candidate, .. } => {
if candidate.trim().is_empty() {
return Err(EnigmaRtcError::InvalidCandidate);
}
Ok(())
}
_ => Err(EnigmaRtcError::InvalidCandidate),
}
}
fn set_microphone_enabled(&self, enabled: bool) -> RtcResult<()> {
self.microphone_enabled.store(enabled, Ordering::SeqCst);
Ok(())
}
fn set_camera_enabled(&self, enabled: bool) -> RtcResult<()> {
self.camera_enabled.store(enabled, Ordering::SeqCst);
Ok(())
}
}