enigma-rtc 0.1.0

WebRTC signaling and session management for Enigma Messenger
Documentation
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Mutex;

use tokio::sync::mpsc;

use crate::error::{EnigmaRtcError, RtcResult};
use crate::signaling::decode_signaling;
use crate::types::{RtcEvent, SignalingMessage};
use crate::webrtc::RtcEngine;

pub struct MockWebRtcEngine {
    next_offer: Mutex<String>,
    next_answer: Mutex<String>,
    applied_answers: Mutex<Vec<String>>,
    remote_candidates: Mutex<Vec<String>>,
    event_sender: Mutex<Option<mpsc::UnboundedSender<RtcEvent>>>,
    microphone_enabled: AtomicBool,
    camera_enabled: AtomicBool,
}

impl MockWebRtcEngine {
    pub fn new() -> Self {
        Self {
            next_offer: Mutex::new("mock-offer".to_string()),
            next_answer: Mutex::new("mock-answer".to_string()),
            applied_answers: Mutex::new(Vec::new()),
            remote_candidates: Mutex::new(Vec::new()),
            event_sender: Mutex::new(None),
            microphone_enabled: AtomicBool::new(true),
            camera_enabled: AtomicBool::new(true),
        }
    }

    pub fn set_offer(&self, sdp: impl Into<String>) {
        if let Ok(mut guard) = self.next_offer.lock() {
            *guard = sdp.into();
        }
    }

    pub fn set_answer(&self, sdp: impl Into<String>) {
        if let Ok(mut guard) = self.next_answer.lock() {
            *guard = sdp.into();
        }
    }

    pub fn emit_candidate(&self, candidate: &str) -> RtcResult<()> {
        let sender = self
            .event_sender
            .lock()
            .map_err(|_| EnigmaRtcError::ChannelClosed)?
            .clone();
        if let Some(tx) = sender {
            let message = SignalingMessage::IceCandidate {
                candidate: candidate.to_string(),
                sdp_mid: Some("0".to_string()),
                sdp_mline_index: Some(0),
            };
            tx.send(RtcEvent::LocalIceCandidate(message))
                .map_err(|_| EnigmaRtcError::ChannelClosed)?;
        }
        Ok(())
    }

    pub fn last_applied_answer(&self) -> Option<String> {
        self.applied_answers
            .lock()
            .ok()
            .and_then(|v| v.last().cloned())
    }

    pub fn candidates(&self) -> Vec<String> {
        self.remote_candidates
            .lock()
            .map(|v| v.clone())
            .unwrap_or_default()
    }

    pub fn microphone_enabled(&self) -> bool {
        self.microphone_enabled.load(Ordering::SeqCst)
    }

    pub fn camera_enabled(&self) -> bool {
        self.camera_enabled.load(Ordering::SeqCst)
    }
}

impl RtcEngine for MockWebRtcEngine {
    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> {
        self.next_offer
            .lock()
            .map(|s| s.clone())
            .map_err(|_| EnigmaRtcError::InvalidState)
    }

    fn create_answer(&self, offer_sdp: &str) -> RtcResult<String> {
        if offer_sdp.trim().is_empty() {
            return Err(EnigmaRtcError::InvalidSdp);
        }
        self.next_answer
            .lock()
            .map(|s| s.clone())
            .map_err(|_| EnigmaRtcError::InvalidState)
    }

    fn apply_answer(&self, answer_sdp: &str) -> RtcResult<()> {
        if answer_sdp.trim().is_empty() {
            return Err(EnigmaRtcError::InvalidSdp);
        }
        if let Ok(mut guard) = self.applied_answers.lock() {
            guard.push(answer_sdp.to_string());
            return Ok(());
        }
        Err(EnigmaRtcError::InvalidState)
    }

    fn add_remote_candidate(&self, candidate_json: &str) -> RtcResult<()> {
        match decode_signaling(candidate_json)? {
            SignalingMessage::IceCandidate { candidate, .. } => {
                if let Ok(mut guard) = self.remote_candidates.lock() {
                    guard.push(candidate);
                    return Ok(());
                }
                Err(EnigmaRtcError::InvalidState)
            }
            _ => 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(())
    }
}