//! Optional raw iOS RemoteIO path for the WebRTC APM experimental mode. //! //! Provides an alternative to `ios_voice_unit.rs` for the //! `SonoraExperimental` processing mode. Instead of //! `kAudioUnitSubType_VoiceProcessingIO` (which owns AEC/NS/AGC), it //! opens `kAudioUnitSubType_RemoteIO` with voice processing explicitly //! disabled so WebRTC APM can own the full signal path. //! //! ## Hard invariants enforced here //! //! * INV_009: Rust AEC only active when platform AEC is disabled. //! * INV_010: VoiceProcessingIO and WebRTC APM AEC are mutually exclusive. //! * INV_011: Software AEC backend receives both capture and render-reference. //! * INV_012: Render reference is copied from decoded/mixed remote PCM //! before playout. //! //! ## Fallback //! //! If RemoteIO construction fails, the caller falls back to `IosVoiceUnit` //! (VPIO) and logs the error. //! //! ## Status //! //! Experimental / disabled by default. Only activated when the user //! explicitly selects `SonoraExperimental` mode via the bridge API. //! //! ## Platform //! //! `kAudioUnitSubType_RemoteIO` is only available in the iOS SDK. //! This module is gated to `target_os = "ios"`. #[cfg(target_os = "ios")] pub use inner::IosRawUnit; #[cfg(target_os = "ios")] mod inner { use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::sync::{Arc, Mutex}; use audiopus::coder::Encoder as OpusEncoder; use coreaudio::audio_unit::audio_format::LinearPcmFlags; use coreaudio::audio_unit::render_callback::{self, data}; use coreaudio::audio_unit::IOType; use coreaudio::audio_unit::{AudioUnit, Element, SampleFormat, Scope, StreamFormat}; use tokio::sync::mpsc; use tracing::{info, warn}; use crate::mobile_voice_backend::VoiceAudioParams; use crate::processor::AudioProcessor; use crate::AudioError; use chanora_protocol::OutPacket; const SAMPLE_RATE_HZ: f64 = 48_000.0; // ------------------------------------------------------------------ // // Render-reference ring buffer // // ------------------------------------------------------------------ // /// 4-slot ring buffer shared between the render callback (writer) and /// the capture callback (reader for Sonora AEC3). Capacity: 4 × 10 ms /// = 40 ms of headroom. /// /// If the capture callback runs before the render callback has written /// a frame it reads zeros (silence reference), which is safe — Sonora /// AEC3 simply skips cancellation for that frame. struct RenderReferenceBuffer { buf: Box<[[f32; 480]; 4]>, write_idx: std::sync::atomic::AtomicUsize, } impl RenderReferenceBuffer { fn new() -> Arc { Arc::new(Self { buf: Box::new([[0.0; 480]; 4]), write_idx: std::sync::atomic::AtomicUsize::new(0), }) } /// Write one 10 ms render-reference frame. Realtime-safe. fn write(&self, frame: &[f32; 480]) { let idx = self.write_idx.load(Ordering::Relaxed); // SAFETY: only one writer (render callback); torn reads // are bounded to one frame of AEC degradation. unsafe { let slot = &self.buf[idx] as *const [f32; 480] as *mut [f32; 480]; (*slot).copy_from_slice(frame); } self.write_idx.store((idx + 1) % 4, Ordering::Relaxed); } /// Read the most recently completed render-reference frame. fn read_latest(&self) -> [f32; 480] { let wi = self.write_idx.load(Ordering::Relaxed); let ri = (wi + 3) % 4; self.buf[ri] } } // SAFETY: accessed from two audio callback threads; data races are // bounded to one frame of AEC quality degradation. unsafe impl Send for RenderReferenceBuffer {} unsafe impl Sync for RenderReferenceBuffer {} // ------------------------------------------------------------------ // // Capture pipeline state // // ------------------------------------------------------------------ // struct RawCaptureState { encoder: OpusEncoder, pcm_accum: Vec, opus_out: [u8; crate::opus_voice::MAX_OPUS_FRAME], voice_out_tx: mpsc::Sender, transmit_active: Arc, output_muted: Arc, frames_sent: Arc, mic_gain: f32, voice_activity_selector: Option>, vad_detector: crate::vad::WebRtcFallbackVad, silero_coreml_worker: Option, current_vad_backend: crate::VadBackend, capture_frame_seq: u64, vad_state: crate::voice_activity::VoiceActivityStateMachine, /// Processing config — retained for route-change reloads. audio_processing_config: Arc>, webrtc_apm_processor: crate::processor::WebRtcApmProcessor, audio_processing_stats: Arc, render_reference: Arc, pending_10ms: [i16; crate::frame::FRAME_10MS_SAMPLES], pending_10ms_len: usize, fallback_warned_backend: Option, wav_recorder: Option>, } impl RawCaptureState { fn new( params: &VoiceAudioParams, render_reference: Arc, ) -> Result { let encoder = crate::opus_voice::new_voip_encoder("ios raw")?; let webrtc_apm_config = params .audio_processing_config .lock() .map(|cfg| crate::processor::webrtc_apm::WebRtcApmConfig::from_audio_config(&cfg)) .unwrap_or_default(); Ok(Self { encoder, pcm_accum: Vec::with_capacity(crate::frame::FRAME_20MS_SAMPLES * 2), opus_out: [0u8; crate::opus_voice::MAX_OPUS_FRAME], voice_out_tx: params.voice_out_tx.clone(), transmit_active: params.transmit_active.clone(), output_muted: params.output_muted.clone(), frames_sent: params.frames_sent.clone(), mic_gain: params.mic_gain, voice_activity_selector: params.voice_activity_selector.clone(), vad_detector: crate::vad::WebRtcFallbackVad::default(), silero_coreml_worker: None, current_vad_backend: crate::VadBackend::WebrtcVad, capture_frame_seq: 0, vad_state: crate::voice_activity::VoiceActivityStateMachine::default(), audio_processing_config: params.audio_processing_config.clone(), webrtc_apm_processor: crate::processor::WebRtcApmProcessor::with_config( webrtc_apm_config, )?, audio_processing_stats: params.audio_processing_stats.clone(), render_reference, pending_10ms: [0_i16; crate::frame::FRAME_10MS_SAMPLES], pending_10ms_len: 0, fallback_warned_backend: None, wav_recorder: None, }) } fn mark_vad_fallback_active(&mut self, failed_backend: crate::VadBackend) { if self.fallback_warned_backend != Some(failed_backend) { self.fallback_warned_backend = Some(failed_backend); if self.capture_frame_seq < 128 { tracing::info!( target: "chanora_audio", backend = failed_backend.as_str(), seq = self.capture_frame_seq, "VAD backend warming up; using WebRTC fallback" ); } else { tracing::warn!( target: "chanora_audio", backend = failed_backend.as_str(), "VAD backend unavailable; using WebRTC fallback for runtime detection" ); } } } fn ingest_i16(&mut self, samples: &[i16]) { // Accumulate into 10 ms frames for VAD / Sonora processing. let mut offset = 0; while offset < samples.len() { let remaining = crate::frame::FRAME_10MS_SAMPLES - self.pending_10ms_len; let take = remaining.min(samples.len() - offset); self.pending_10ms[self.pending_10ms_len..self.pending_10ms_len + take] .copy_from_slice(&samples[offset..offset + take]); self.pending_10ms_len += take; offset += take; if self.pending_10ms_len == crate::frame::FRAME_10MS_SAMPLES { let frame = self.pending_10ms; self.process_10ms_capture_frame(&frame); self.pending_10ms_len = 0; } } if !self.transmit_active.load(Ordering::Relaxed) { self.pcm_accum.clear(); return; } // Encode complete 20 ms Opus frames. while self.pcm_accum.len() >= crate::frame::FRAME_20MS_SAMPLES { let mut frame = [0i16; crate::frame::FRAME_20MS_SAMPLES]; frame.copy_from_slice(&self.pcm_accum[..crate::frame::FRAME_20MS_SAMPLES]); self.pcm_accum.drain(..crate::frame::FRAME_20MS_SAMPLES); match self.encoder.encode(&frame, &mut self.opus_out[..]) { Ok(len) => { crate::opus_voice::send_voip_frame( &self.voice_out_tx, &self.frames_sent, &self.opus_out, len, || { warn!( target: "chanora_audio", "ios raw: voice_out queue full; dropping frame" ); }, || {}, ); } Err(e) => { tracing::error!(target: "chanora_audio", error = %e, "ios raw opus encode failed"); } } } } fn process_10ms_capture_frame( &mut self, samples: &[i16; crate::frame::FRAME_10MS_SAMPLES], ) { let mut frame = [0.0_f32; crate::frame::FRAME_10MS_SAMPLES]; for (dst, src) in frame.iter_mut().zip(samples.iter().copied()) { *dst = crate::frame::i16_to_f32(src); } let input_dbfs = crate::frame::dbfs(&frame); // WAV tap: raw mic (before processing). if let Some(ref rec) = self.wav_recorder { rec.push_raw_mic(&frame); } // Feed render reference to WebRTC APM before capture so AEC can adapt. let render_ref = self.render_reference.read_latest(); self.webrtc_apm_processor.process_render(&render_ref); self.webrtc_apm_processor.process_capture(&mut frame); // WAV tap: processed mic (after WebRTC APM). if let Some(ref rec) = self.wav_recorder { rec.push_processed_mic(&frame); } let voice_activity_mode = self .voice_activity_selector .as_ref() .map(|selector| selector.mode() == crate::TransmitMode::VoiceActivity) .unwrap_or(false); if !voice_activity_mode { self.silero_coreml_worker = None; self.current_vad_backend = crate::VadBackend::Disabled; self.fallback_warned_backend = None; self.audio_processing_stats.set_vad_fallback_active(false); } let (vad_backend, vad_hangover) = self .audio_processing_config .try_lock() .map(|cfg| (cfg.vad_backend, cfg.vad_hangover_ms)) .unwrap_or(( crate::VadBackend::WebrtcVad, crate::voice_activity::VAD_HANGOVER_MS, )); if voice_activity_mode { self.vad_state.configure( crate::voice_activity::VAD_OPEN_AFTER_MS, vad_hangover, crate::voice_activity::VAD_MIN_TX_MS, ); } if voice_activity_mode && vad_backend != self.current_vad_backend { self.current_vad_backend = vad_backend; self.fallback_warned_backend = None; if vad_backend == crate::VadBackend::SileroOnnx { self.silero_coreml_worker = crate::vad::apple_coreml::AppleCoreMlVadWorker::try_new(); if self.silero_coreml_worker.is_none() { self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx); self.audio_processing_stats.set_vad_fallback_active(true); } else { self.audio_processing_stats.set_vad_fallback_active(false); } } else { self.silero_coreml_worker = None; self.audio_processing_stats.set_vad_fallback_active(false); } self.vad_state.reset(); } let (vad_probability, active) = if voice_activity_mode { self.capture_frame_seq = self.capture_frame_seq.wrapping_add(1); let capture_seq = self.capture_frame_seq; let mut used_fallback_vad = false; let vad = if vad_backend == crate::VadBackend::Disabled { crate::vad::VadOutput { probability: 1.0, speech: true, } } else if vad_backend == crate::VadBackend::SileroOnnx { if let Some(worker) = self.silero_coreml_worker.as_ref() { let enqueued = worker.try_send(capture_seq, &frame); if !worker.is_stale(capture_seq) { let p = worker.latest_probability(); crate::vad::VadOutput { probability: p, speech: p >= 0.5, } } else if enqueued { crate::vad::VadOutput { probability: 0.0, speech: false, } } else { used_fallback_vad = true; self.mark_vad_fallback_active(vad_backend); crate::vad::VoiceActivityDetector::process_10ms( &mut self.vad_detector, &frame, ) } } else { used_fallback_vad = true; self.mark_vad_fallback_active(vad_backend); crate::vad::VoiceActivityDetector::process_10ms( &mut self.vad_detector, &frame, ) } } else { crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, &frame) }; self.audio_processing_stats .set_vad_fallback_active(used_fallback_vad); (vad.probability, self.vad_state.update(vad.speech)) } else { (0.0, false) }; if let Some(sel) = &self.voice_activity_selector { sel.set_voice_activity_open(voice_activity_mode && active); } self.audio_processing_stats.update_capture( input_dbfs, crate::frame::dbfs(&frame), vad_probability, voice_activity_mode && active, self.transmit_active.load(Ordering::Relaxed), ); if !self.transmit_active.load(Ordering::Relaxed) { return; } let gain = self.mic_gain; if (gain - 1.0).abs() < f32::EPSILON { self.pcm_accum .extend(frame.iter().copied().map(crate::frame::f32_to_i16)); } else { self.pcm_accum.extend(frame.iter().copied().map(|s| { let scaled = crate::frame::f32_to_i16(s) as f32 * gain; scaled.clamp(i16::MIN as f32, i16::MAX as f32) as i16 })); } } } // ------------------------------------------------------------------ // // IosRawUnit // // ------------------------------------------------------------------ // /// Raw iOS RemoteIO audio unit for the Sonora experimental path. pub struct IosRawUnit { unit: AudioUnit, } impl IosRawUnit { /// Open a RemoteIO AudioUnit, install render + input callbacks, start. pub(crate) fn start(params: VoiceAudioParams) -> Result { // INV_010: reject if config requests VPIO (that's IosVoiceUnit's job). { let cfg = params.audio_processing_config.lock().unwrap(); if cfg.ios_mode == crate::IosVoiceProcessingMode::PlatformVoiceProcessing { return Err(AudioError::InvalidAudioProcessingConfig( "IosRawUnit requires raw WebRTC APM mode".to_string(), )); } } let mut unit = AudioUnit::new_uninitialized(IOType::RemoteIO) .map_err(|e| AudioError::Backend(format!("remoteio new: {e}")))?; // Enable input on bus 1. const ENABLE_IO: u32 = 2003; let enable: u32 = 1; unit.set_property(ENABLE_IO, Scope::Input, Element::Input, Some(&enable)) .map_err(|e| AudioError::Backend(format!("remoteio enable input: {e}")))?; // 48 kHz Int16 mono on both buses. let fmt = StreamFormat { sample_rate: SAMPLE_RATE_HZ, sample_format: SampleFormat::I16, flags: LinearPcmFlags::IS_SIGNED_INTEGER | LinearPcmFlags::IS_PACKED, channels: 1, }; unit.set_stream_format(fmt, Scope::Input, Element::Output) .map_err(|e| AudioError::StreamConfig(format!("remoteio fmt output: {e}")))?; unit.set_stream_format(fmt, Scope::Output, Element::Input) .map_err(|e| AudioError::StreamConfig(format!("remoteio fmt input: {e}")))?; // Shared render-reference buffer (INV_011 / INV_012). let render_ref_buf = RenderReferenceBuffer::new(); let render_ref_for_capture = render_ref_buf.clone(); let mut capture_state = RawCaptureState::new(¶ms, render_ref_for_capture)?; unit.set_input_callback(move |args: render_callback::Args>| { capture_state.ingest_i16(args.data.buffer); Ok(()) }) .map_err(|e| AudioError::Backend(format!("remoteio input cb: {e}")))?; let mut scratch: Vec = Vec::with_capacity(2048); let handler = params.handler.clone(); let output_gain = params.output_gain.clone(); let output_muted = params.output_muted.clone(); let stats_render = params.audio_processing_stats.clone(); unit.set_render_callback(move |args: render_callback::Args>| { let out = args.data.buffer; let n = out.len(); let stereo_n = n * 2; if scratch.len() < stereo_n { scratch.resize(stereo_n, 0.0); } scratch[..stereo_n].fill(0.0); match handler.try_lock() { Ok(mut h) => { let _ = h.fill_buffer(&mut scratch[..stereo_n]); } Err(std::sync::TryLockError::WouldBlock) => { stats_render.increment_callback_xrun(); } Err(std::sync::TryLockError::Poisoned(e)) => { warn!(target: "chanora_audio", "AudioHandler poisoned (raw render): {e}"); } } // INV_012: copy render reference BEFORE playout. let mono_n = n.min(480); let mut ref_frame = [0.0_f32; 480]; crate::voice_render::downmix_stereo_f32_to_mono_f32( &scratch[..stereo_n], &mut ref_frame[..mono_n], ); render_ref_buf.write(&ref_frame); let gain = f32::from_bits(output_gain.load(Ordering::Relaxed)); let muted = output_muted.load(Ordering::Relaxed); let mix_stats = crate::voice_render::downmix_stereo_f32_to_mono_i16( &scratch[..stereo_n], out, gain, muted, ); if mix_stats.clipped_samples > 0 { stats_render.add_clipped_samples(mix_stats.clipped_samples); } stats_render.update_render(crate::frame::dbfs(&scratch[..stereo_n]), n as u32); Ok(()) }) .map_err(|e| AudioError::Backend(format!("remoteio render cb: {e}")))?; unit.initialize() .map_err(|e| AudioError::Backend(format!("remoteio init: {e}")))?; unit.start() .map_err(|e| AudioError::Backend(format!("remoteio start: {e}")))?; info!( target: "chanora_audio", sample_rate_hz = SAMPLE_RATE_HZ, "ios RemoteIO (Sonora experimental) started" ); Ok(Self { unit }) } /// Restart the unit after a route change (stop → uninit → init → start). pub fn restart(&mut self) -> Result<(), AudioError> { self.unit .stop() .map_err(|e| AudioError::Backend(format!("remoteio restart stop: {e}")))?; self.unit .uninitialize() .map_err(|e| AudioError::Backend(format!("remoteio restart uninit: {e}")))?; self.unit .initialize() .map_err(|e| AudioError::Backend(format!("remoteio restart init: {e}")))?; self.unit .start() .map_err(|e| AudioError::Backend(format!("remoteio restart start: {e}")))?; info!(target: "chanora_audio", "ios RemoteIO restarted"); Ok(()) } /// Pause the unit during an AVAudioSession interruption. pub fn pause(&mut self) -> Result<(), AudioError> { self.unit .stop() .map_err(|e| AudioError::Backend(format!("remoteio pause: {e}"))) } /// Resume the unit after an interruption ends. pub fn resume(&mut self) -> Result<(), AudioError> { self.unit .start() .map_err(|e| AudioError::Backend(format!("remoteio resume: {e}"))) } } impl Drop for IosRawUnit { fn drop(&mut self) { if let Err(e) = self.unit.stop() { warn!(target: "chanora_audio", error = %e, "ios RemoteIO stop on drop failed"); } else { info!(target: "chanora_audio", "ios RemoteIO stopped"); } } } }