pub mod messages; pub mod jitter; use crate::audio::{AudioBackend, pipewire_impl::PipeWireBackend}; use crate::codec::{AudioEncoder, opus_impl::OpusEncoder}; use crate::core::jitter::{JitterBuffer, FRAME_SAMPLES}; use crate::network::{ NetworkTransport, RoomState, PeerState, RoomEvent, PeerSpeakTicket, iroh_impl::IrohTransport, gossip::IrohGossipState, }; use crate::core::messages::{CoreCommand, UiEvent}; use crate::config::NetworkMode; use iroh::{Endpoint, EndpointId, RelayMode, endpoint::presets, protocol::Router}; use iroh_gossip::net::Gossip; use tokio::sync::{mpsc, Mutex}; use std::collections::HashMap; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; pub struct CoreController { cmd_tx: mpsc::Sender, } impl CoreController { pub fn new(ui_tx: mpsc::Sender) -> Self { let (cmd_tx, cmd_rx) = mpsc::channel(100); std::thread::spawn(move || { let rt = tokio::runtime::Runtime::new().expect("Failed to create Tokio runtime"); rt.block_on(async move { crate::log_msg("Starting core network loop in dedicated Tokio runtime"); if let Err(e) = run_core_loop(cmd_rx, ui_tx).await { crate::log_msg(&format!("App core loop failed: {:?}", e)); } }); }); Self { cmd_tx } } pub fn send(&self, cmd: CoreCommand) -> Result<(), mpsc::error::TrySendError> { self.cmd_tx.try_send(cmd) } } struct ActiveSession { endpoint: Endpoint, router: Router, room_state: Arc, capture_thread: std::thread::JoinHandle<()>, datagram_task: tokio::task::JoinHandle<()>, mixer_task: tokio::task::JoinHandle<()>, event_task: tokio::task::JoinHandle<()>, } impl ActiveSession { async fn shutdown(self, audio_backend: Arc) { crate::log_msg("ActiveSession::shutdown started"); self.datagram_task.abort(); self.mixer_task.abort(); self.event_task.abort(); crate::log_msg("Aborted tasks"); let audio_backend_clone = audio_backend.clone(); let _ = tokio::task::spawn_blocking(move || { crate::log_msg("Stopping audio backend..."); let _ = audio_backend_clone.stop(); crate::log_msg("Audio backend stopped"); }).await; crate::log_msg("Leaving room..."); let _ = self.room_state.leave().await; crate::log_msg("Room left"); crate::log_msg("Shutting down router..."); let _ = tokio::time::timeout(std::time::Duration::from_secs(1), self.router.shutdown()).await; crate::log_msg("Router shut down"); crate::log_msg("Joining capture thread..."); let _ = self.capture_thread.join(); crate::log_msg("ActiveSession::shutdown complete"); } } async fn run_core_loop( mut cmd_rx: mpsc::Receiver, ui_tx: mpsc::Sender, ) -> Result<(), anyhow::Error> { let memory_lookup = iroh::address_lookup::memory::MemoryLookup::new(); let secret_key = iroh::SecretKey::generate(); let audio_backend = Arc::new(PipeWireBackend::new()); let is_muted = Arc::new(AtomicBool::new(false)); let is_deafened = Arc::new(AtomicBool::new(false)); let ptt_mode = Arc::new(AtomicBool::new(false)); let ptt_active = Arc::new(AtomicBool::new(false)); let noise_gate_threshold = Arc::new(std::sync::atomic::AtomicU32::new(0.01f32.to_bits())); let peer_volumes = Arc::new(Mutex::new(HashMap::::new())); let mut current_name = "Anonymous".to_string(); let mut network_mode = NetworkMode::default(); let mut active_session: Option = None; while let Some(cmd) = cmd_rx.recv().await { match cmd { CoreCommand::Join { name, ticket, input_device, output_device } => { current_name = name.clone(); // Clean up any existing session if let Some(session) = active_session.take() { crate::log_msg("Shutting down existing active session"); session.shutdown(audio_backend.clone()).await; } // Build the endpoint per the configured relay/discovery posture. // All postures keep the in-memory address lookup (fed by tickets // and gossip); they differ in whether n0's relay and DNS presence // beacon are used. `Minimal` sets only the mandatory crypto // provider and deliberately omits the n0 DNS publish/resolve. let bind_result = match network_mode { NetworkMode::N0Full => { Endpoint::builder(presets::N0) .secret_key(secret_key.clone()) .address_lookup(memory_lookup.clone()) .bind() .await } NetworkMode::RelayNoDiscovery => { Endpoint::builder(presets::Minimal) .secret_key(secret_key.clone()) .relay_mode(RelayMode::Default) .address_lookup(memory_lookup.clone()) .bind() .await } NetworkMode::DirectOnly => { Endpoint::builder(presets::Minimal) .secret_key(secret_key.clone()) .relay_mode(RelayMode::Disabled) .address_lookup(memory_lookup.clone()) .bind() .await } }; let endpoint = match bind_result { Ok(ep) => ep, Err(e) => { let _ = ui_tx.send(UiEvent::Error(format!("Failed to bind endpoint: {}", e))).await; continue; } }; endpoint.online().await; // Determine target ticket let ticket_str = if ticket.trim().is_empty() || ticket == "create" { let topic_id: [u8; 32] = rand::random(); let host_addr = endpoint.addr(); crate::log_msg(&format!("Creating room. host_addr={:?}, topic_id={:?}", host_addr, topic_id)); let ticket = PeerSpeakTicket { host_addr, topic_id }; ticket.to_string() } else { let ticket_str = ticket.trim().to_string(); crate::log_msg(&format!("Joining room with existing ticket={}", ticket_str)); ticket_str }; // Initialize Gossip and Transport let gossip = Gossip::builder().spawn(endpoint.clone()); let (transport, audio_proto) = IrohTransport::new(endpoint.clone()); let transport = Arc::new(transport); // Start Router let router = iroh::protocol::Router::builder(endpoint.clone()) .accept(iroh_gossip::net::GOSSIP_ALPN, gossip.clone()) .accept(b"peerspeak-audio", audio_proto) .spawn(); let room_state = Arc::new(IrohGossipState::new( endpoint.clone(), gossip.clone(), memory_lookup.clone(), )); let self_state = PeerState { name: current_name.clone(), is_muted: is_muted.load(Ordering::Relaxed), addr: endpoint.addr(), }; crate::log_msg(&format!("Attempting room_state.join with self_state={:?}", self_state)); if let Err(e) = room_state.join(&ticket_str, self_state.clone()).await { crate::log_msg(&format!("Error room_state.join failed: {:?}", e)); let _ = ui_tx.send(UiEvent::Error(format!("Failed to join room: {}", e))).await; let _ = router.shutdown().await; continue; } crate::log_msg("Joined room successfully via room_state"); // Setup raw audio channels let (capture_tx, capture_rx) = std::sync::mpsc::channel(); let (playback_tx, playback_rx) = std::sync::mpsc::channel(); if let Err(e) = audio_backend.start_capture(capture_tx, input_device) { let _ = ui_tx.send(UiEvent::Error(format!("Failed to start capture: {}", e))).await; let _ = room_state.leave().await; let _ = router.shutdown().await; continue; } if let Err(e) = audio_backend.start_playback(playback_rx, output_device) { let _ = ui_tx.send(UiEvent::Error(format!("Failed to start playback: {}", e))).await; let _ = audio_backend.stop(); let _ = room_state.leave().await; let _ = router.shutdown().await; continue; } let jitter: Arc>> = Arc::new(Mutex::new(HashMap::new())); // 1. Capture & encoding thread let is_muted_clone = is_muted.clone(); let ptt_mode_clone = ptt_mode.clone(); let ptt_active_clone = ptt_active.clone(); let noise_gate_threshold_clone = noise_gate_threshold.clone(); let transport_clone = transport.clone(); let capture_thread = std::thread::spawn(move || { use opus::{Channels, Application}; let mut encoder = match OpusEncoder::new(48000, Channels::Mono, Application::Voip) { Ok(enc) => enc, Err(e) => { crate::log_msg(&format!("Capture thread error: {:?}", e)); return; } }; // Per-sender packet sequence number, prepended to every frame so // receivers can reorder and conceal loss. Wraps after ~years. let mut seq: u32 = 0; while let Ok(pcm) = capture_rx.recv() { if is_muted_clone.load(Ordering::Relaxed) { continue; } if ptt_mode_clone.load(Ordering::Relaxed) && !ptt_active_clone.load(Ordering::Relaxed) { continue; } let ng_bits = noise_gate_threshold_clone.load(Ordering::Relaxed); let ng_thresh = f32::from_bits(ng_bits); if ng_thresh > 0.0001 { let mut sum_sq = 0.0f32; for &sample in &pcm { let normalized = (sample as f32) / 32768.0; sum_sq += normalized * normalized; } let rms = (sum_sq / pcm.len() as f32).sqrt(); if rms < ng_thresh { continue; } } if let Ok(encoded) = encoder.encode(&pcm) { // Frame on the wire: [seq: u32 LE][opus payload]. let mut packet = Vec::with_capacity(4 + encoded.len()); packet.extend_from_slice(&seq.to_le_bytes()); packet.extend_from_slice(&encoded); seq = seq.wrapping_add(1); transport_clone.broadcast(bytes::Bytes::from(packet)); } } }); // 2. Receiver task: parse the sequence header and hand each packet // to that peer's jitter buffer. Decoding happens later, on the // playout side, so loss can be concealed at the right moment. let transport_recv = transport.clone(); let jitter_recv = jitter.clone(); let datagram_task = tokio::spawn(async move { let mut datagram_rx = match transport_recv.receive_datagrams().await { Ok(rx) => rx, Err(e) => { crate::log_msg(&format!("Receiver task error: {:?}", e)); return; } }; while let Some((from_peer, bytes)) = datagram_rx.recv().await { if bytes.len() < 4 { continue; // malformed: missing sequence header } let seq = u32::from_le_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]); let payload = bytes[4..].to_vec(); let mut guard = jitter_recv.lock().await; let buffer = match guard.entry(from_peer) { std::collections::hash_map::Entry::Occupied(entry) => entry.into_mut(), std::collections::hash_map::Entry::Vacant(entry) => { match JitterBuffer::new() { Ok(jb) => entry.insert(jb), Err(e) => { crate::log_msg(&format!("Failed to init jitter buffer for {:?}: {:?}", from_peer, e)); continue; } } } }; buffer.insert(seq, payload); } }); // 3. Mixing & level extraction loop task. Every 20ms, pull one // concealed frame per peer from its jitter buffer, apply // per-peer volume, sum, and hand the mix to playback. let jitter_mixer = jitter.clone(); let is_deafened_clone = is_deafened.clone(); let peer_volumes_mixer = peer_volumes.clone(); let ui_tx_mixer = ui_tx.clone(); let mixer_task = tokio::spawn(async move { let mut interval = tokio::time::interval(Duration::from_millis(20)); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { interval.tick().await; let current_volumes = peer_volumes_mixer.lock().await.clone(); let mut active_levels = Vec::new(); let mut peer_frames = Vec::new(); { let mut guard = jitter_mixer.lock().await; for (&peer_id, buffer) in guard.iter_mut() { // `None` means idle/buffering: contribute nothing // and report a zero level so the UI shows idle. let Some(mut frame) = buffer.pop_frame() else { active_levels.push((peer_id, 0.0)); continue; }; let vol = current_volumes.get(&peer_id).copied().unwrap_or(1.0); if (vol - 1.0).abs() > f32::EPSILON { for sample in frame.iter_mut() { *sample = (*sample as f32 * vol).clamp(i16::MIN as f32, i16::MAX as f32) as i16; } } let sum_sq: f32 = frame.iter().map(|&x| (x as f32).powi(2)).sum(); let rms = (sum_sq / frame.len().max(1) as f32).sqrt(); let level = (rms / 32768.0).clamp(0.0, 1.0); active_levels.push((peer_id, level)); peer_frames.push(frame); } } let mut mixed = vec![0i16; FRAME_SAMPLES]; if !peer_frames.is_empty() { for (i, out) in mixed.iter_mut().enumerate() { let sum: i32 = peer_frames .iter() .map(|f| f.get(i).copied().unwrap_or(0) as i32) .sum(); *out = sum.clamp(i16::MIN as i32, i16::MAX as i32) as i16; } } let frame_to_send = if is_deafened_clone.load(Ordering::Relaxed) { vec![0i16; FRAME_SAMPLES] } else { mixed }; if playback_tx.send(frame_to_send).is_err() { break; } let _ = ui_tx_mixer.send(UiEvent::AudioLevels(active_levels)).await; } }); // 4. Room event subscriber task let mut room_events = match room_state.subscribe_events().await { Ok(rx) => rx, Err(e) => { let _ = ui_tx.send(UiEvent::Error(format!("Failed to subscribe events: {}", e))).await; continue; } }; let ui_tx_events = ui_tx.clone(); let jitter_events = jitter.clone(); let transport_events = transport.clone(); let event_task = tokio::spawn(async move { while let Some(event) = room_events.recv().await { match event { RoomEvent::PeerJoined(peer_id, state) => { // Establish the audio connection as soon as the peer // is known (the transport dedupes the full-mesh race). transport_events.connect_peer(peer_id).await; let _ = ui_tx_events.send(UiEvent::PeerJoined { id: peer_id, state }).await; } RoomEvent::PeerLeft(peer_id) => { transport_events.disconnect_peer(peer_id).await; jitter_events.lock().await.remove(&peer_id); let _ = ui_tx_events.send(UiEvent::PeerLeft { id: peer_id }).await; } RoomEvent::PeerUpdated(peer_id, state) => { let _ = ui_tx_events.send(UiEvent::PeerUpdated { id: peer_id, state }).await; } } } }); let session = ActiveSession { endpoint: endpoint.clone(), router, room_state: room_state.clone(), capture_thread, datagram_task, mixer_task, event_task, }; let self_id = endpoint.id().to_string(); let _ = ui_tx.send(UiEvent::RoomJoined { ticket: ticket_str, self_id }).await; active_session = Some(session); } CoreCommand::Leave => { if let Some(session) = active_session.take() { session.shutdown(audio_backend.clone()).await; let _ = ui_tx.send(UiEvent::RoomLeft).await; } } CoreCommand::ToggleMute => { let current = is_muted.load(Ordering::Relaxed); let new_state = !current; is_muted.store(new_state, Ordering::Relaxed); if let Some(session) = &active_session { let self_state = PeerState { name: current_name.clone(), is_muted: new_state, addr: session.endpoint.addr(), }; let _ = session.room_state.update_self_state(self_state).await; } } CoreCommand::ToggleDeafen => { let current = is_deafened.load(Ordering::Relaxed); is_deafened.store(!current, Ordering::Relaxed); } CoreCommand::SetPttMode(enabled) => { ptt_mode.store(enabled, Ordering::Relaxed); } CoreCommand::SetPttActive(active) => { ptt_active.store(active, Ordering::Relaxed); } CoreCommand::SetPeerVolume(peer_id, vol) => { let mut guard = peer_volumes.lock().await; guard.insert(peer_id, vol); } CoreCommand::SetNoiseGateThreshold(threshold) => { noise_gate_threshold.store(threshold.to_bits(), Ordering::Relaxed); } CoreCommand::SetNetworkMode(mode) => { network_mode = mode; } } } Ok(()) }