enigma-rtc 0.1.0

WebRTC signaling and session management for Enigma Messenger
Documentation
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(())
    }
}