diff --git a/crates/chanora_audio/src/android_render_ring.rs b/crates/chanora_audio/src/android_render_ring.rs new file mode 100644 index 0000000..668c9e7 --- /dev/null +++ b/crates/chanora_audio/src/android_render_ring.rs @@ -0,0 +1,116 @@ +use std::sync::Arc; + +use crossbeam::queue::ArrayQueue; + +/// Fixed-capacity PCM handoff from the Android render producer task to +/// the Oboe output callback. +pub(crate) struct AndroidRenderRing { + frames: Arc>, +} + +impl AndroidRenderRing { + pub(crate) fn new(capacity: usize) -> Self { + Self { + frames: Arc::new(ArrayQueue::new((capacity / 2).max(1))), + } + } + + pub(crate) fn producer(&self) -> AndroidRenderRingProducer { + AndroidRenderRingProducer { + frames: Arc::clone(&self.frames), + } + } + + pub(crate) fn consumer(&self) -> AndroidRenderRingConsumer { + AndroidRenderRingConsumer { + frames: Arc::clone(&self.frames), + } + } +} + +pub(crate) struct AndroidRenderRingProducer { + frames: Arc>, +} + +impl AndroidRenderRingProducer { + pub(crate) fn push_frame_lossy(&self, samples: &[f32]) { + for frame in samples.chunks_exact(2) { + let stereo_frame = [frame[0], frame[1]]; + if self.frames.push(stereo_frame).is_err() { + let _ = self.frames.pop(); + let _ = self.frames.push(stereo_frame); + } + } + } +} + +pub(crate) struct AndroidRenderRingConsumer { + frames: Arc>, +} + +impl AndroidRenderRingConsumer { + pub(crate) fn drain_into_zero_filling(&self, out: &mut [f32]) { + let mut chunks = out.chunks_exact_mut(2); + for frame_out in &mut chunks { + let frame = self.frames.pop().unwrap_or([0.0, 0.0]); + frame_out.copy_from_slice(&frame); + } + for sample in chunks.into_remainder() { + *sample = 0.0; + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn producer_drops_oldest_samples_when_ring_is_full() { + let ring = AndroidRenderRing::new(4); + let producer = ring.producer(); + let consumer = ring.consumer(); + + producer.push_frame_lossy(&[1.0, 2.0, 3.0, 4.0]); + producer.push_frame_lossy(&[5.0, 6.0]); + + let mut out = [0.0; 4]; + consumer.drain_into_zero_filling(&mut out); + + assert_eq!(out, [3.0, 4.0, 5.0, 6.0]); + } + + #[test] + fn overflow_after_partial_consumer_drain_preserves_stereo_pairing() { + let ring = AndroidRenderRing::new(4); + let producer = ring.producer(); + let consumer = ring.consumer(); + + producer.push_frame_lossy(&[1.0, 10.0, 2.0, 20.0]); + + let mut odd_out = [9.0]; + consumer.drain_into_zero_filling(&mut odd_out); + assert_eq!(odd_out, [0.0]); + + producer.push_frame_lossy(&[3.0, 30.0]); + + let mut out = [0.0; 4]; + consumer.drain_into_zero_filling(&mut out); + + assert_eq!(out, [2.0, 20.0, 3.0, 30.0]); + } + + #[test] + fn consumer_zero_fills_tail_on_underrun() { + let ring = AndroidRenderRing::new(4); + let producer = ring.producer(); + let consumer = ring.consumer(); + + producer.push_frame_lossy(&[0.25, -0.25]); + + let mut out = [9.0; 4]; + consumer.drain_into_zero_filling(&mut out); + + assert_eq!(out, [0.25, -0.25, 0.0, 0.0]); + } +} diff --git a/crates/chanora_audio/src/android_voice_unit.rs b/crates/chanora_audio/src/android_voice_unit.rs index 43bd976..6f0d35d 100644 --- a/crates/chanora_audio/src/android_voice_unit.rs +++ b/crates/chanora_audio/src/android_voice_unit.rs @@ -53,7 +53,6 @@ use crate::mobile_voice_backend::{ BackendEventTx, EffectEngagement, EffectEngine, InputPresetChoice, MobileVoiceAudioBackend, SharingModeChoice, VoiceAudioParams, }; -use chanora_protocol::OutPacket; use tsclientlib::audio::AudioHandler; use crate::{engine::SessionAudioId, AudioError}; @@ -86,40 +85,11 @@ use crate::processor::AudioProcessor; const RENDER_REF_SLOTS: usize = 4; const RENDER_REF_SAMPLES: usize = crate::frame::FRAME_10MS_SAMPLES; +const ANDROID_RENDER_PULL_SAMPLES: usize = crate::frame::FRAME_20MS_SAMPLES * 2; +const ANDROID_RENDER_RING_CAPACITY: usize = ANDROID_RENDER_PULL_SAMPLES * 5; -struct RenderReferenceBuffer { - buf: Box<[[f32; RENDER_REF_SAMPLES]; RENDER_REF_SLOTS]>, - write_idx: std::sync::atomic::AtomicUsize, -} - -impl RenderReferenceBuffer { - fn new() -> Arc { - Arc::new(Self { - buf: Box::new([[0.0_f32; RENDER_REF_SAMPLES]; RENDER_REF_SLOTS]), - write_idx: std::sync::atomic::AtomicUsize::new(0), - }) - } - - fn write(&self, frame: &[f32; RENDER_REF_SAMPLES]) { - let idx = self.write_idx.load(Ordering::Relaxed); - unsafe { - let slot = &self.buf[idx] as *const [f32; RENDER_REF_SAMPLES] - as *mut [f32; RENDER_REF_SAMPLES]; - (*slot).copy_from_slice(frame); - } - self.write_idx - .store((idx + 1) % RENDER_REF_SLOTS, Ordering::Relaxed); - } - - fn read_latest(&self) -> [f32; RENDER_REF_SAMPLES] { - let wi = self.write_idx.load(Ordering::Relaxed); - let ri = (wi + RENDER_REF_SLOTS - 1) % RENDER_REF_SLOTS; - self.buf[ri] - } -} - -unsafe impl Send for RenderReferenceBuffer {} -unsafe impl Sync for RenderReferenceBuffer {} +type RenderReferenceBuffer = + crate::render_reference::RenderReferenceBuffer; // --- Capture state for Oboe input callback (SDD-111 / SDD-120) ---- // @@ -138,9 +108,8 @@ struct AndroidCaptureState { encoder: OpusEncoder, pcm_accum: Vec, opus_out: [u8; crate::opus_voice::MAX_OPUS_FRAME], - voice_out_tx: mpsc::Sender, + voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender, transmit_active: Arc, - frames_sent: Arc, mic_gain: f32, voice_activity_selector: Option>, vad_detector: crate::vad::WebRtcFallbackVad, @@ -164,7 +133,7 @@ struct AndroidCaptureState { impl AndroidCaptureState { fn new( - voice_out_tx: mpsc::Sender, + voice_out_tx: mpsc::Sender, transmit_active: Arc, frames_sent: Arc, mic_gain: f32, @@ -188,9 +157,12 @@ impl AndroidCaptureState { encoder, pcm_accum: Vec::with_capacity(crate::frame::FRAME_20MS_SAMPLES * 2), opus_out: [0u8; crate::opus_voice::MAX_OPUS_FRAME], - voice_out_tx, + voice_out_tx: crate::opus_voice::start_out_packet_worker( + voice_out_tx, + frames_sent.clone(), + "android", + )?, transmit_active, - frames_sent, mic_gain, voice_activity_selector, vad_detector: crate::vad::WebRtcFallbackVad::default(), @@ -221,8 +193,10 @@ impl AndroidCaptureState { self.audio_processing_stats .record_callback_frames(samples.len() as u64); if self.input_sample_rate_hz != crate::frame::SAMPLE_RATE_HZ { - let resampled = self.resample_capture_to_48k(samples); + self.resample_capture_to_48k(samples); + let resampled = std::mem::take(&mut self.resample_scratch); self.ingest_48k_i16(&resampled); + self.resample_scratch = resampled; return; } self.ingest_48k_i16(samples); @@ -241,6 +215,7 @@ impl AndroidCaptureState { if self.pending_10ms_len == crate::frame::FRAME_10MS_SAMPLES { let frame = self.pending_10ms; self.process_10ms_capture_frame(&frame); + self.encode_complete_20ms_frames(); self.pending_10ms_len = 0; } } @@ -250,6 +225,10 @@ impl AndroidCaptureState { return; } + self.encode_complete_20ms_frames(); + } + + fn encode_complete_20ms_frames(&mut self) { 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]); @@ -258,7 +237,6 @@ impl AndroidCaptureState { Ok(len) => { crate::opus_voice::send_voip_frame( &self.voice_out_tx, - &self.frames_sent, &self.opus_out, len, || { @@ -286,35 +264,18 @@ impl AndroidCaptureState { } } - fn resample_capture_to_48k(&mut self, samples: &[i16]) -> Vec { - if samples.is_empty() { - return Vec::new(); + fn resample_capture_to_48k(&mut self, samples: &[i16]) -> usize { + let result = crate::capture_resampler::resample_capture_to_48k( + samples, + self.input_sample_rate_hz, + &mut self.resample_pos, + &mut self.resample_last, + &mut self.resample_scratch, + ); + if result.dropped { + self.audio_processing_stats.increment_callback_xrun(); } - self.resample_scratch.clear(); - let ratio = self.input_sample_rate_hz as f64 / crate::frame::SAMPLE_RATE_HZ as f64; - let mut pos = self.resample_pos; - while pos < samples.len() as f64 { - let i = pos.floor() as isize; - let frac = pos - i as f64; - let a = if i <= 0 { - self.resample_last as f64 - } else { - samples[(i - 1) as usize] as f64 - }; - let b = if i < samples.len() as isize { - samples[i as usize] as f64 - } else { - a - }; - let value = (a + frac * (b - a)) - .round() - .clamp(i16::MIN as f64, i16::MAX as f64) as i16; - self.resample_scratch.push(value); - pos += ratio; - } - self.resample_pos = pos - samples.len() as f64; - self.resample_last = *samples.last().unwrap_or(&self.resample_last); - self.resample_scratch.clone() + result.output_len } fn set_input_sample_rate_hz(&mut self, sample_rate_hz: u32) { @@ -393,15 +354,9 @@ impl AndroidCaptureState { self.fallback_warned_backend = None; match vad_backend { crate::VadBackend::SileroOnnx => { - let path = crate::vad::silero_model_bundle_path(); - self.silero_vad_worker = - crate::vad::silero_onnx::SileroOnnxVadWorker::try_new(&path); - if self.silero_vad_worker.is_none() { - warn!( - target: "chanora_audio", - "android: Silero VAD model not found at {path}; falling back to WebRTC VAD" - ); - } + self.silero_vad_worker = None; + self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx); + self.audio_processing_stats.set_vad_fallback_active(true); } _ => { self.silero_vad_worker = None; @@ -420,20 +375,38 @@ impl AndroidCaptureState { speech: true, } } else if vad_backend == crate::VadBackend::SileroOnnx { - if let Some(worker) = self.silero_vad_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, + match crate::vad::callback_vad_worker_policy( + voice_activity_mode, + vad_backend, + self.silero_vad_worker.is_some(), + ) { + crate::vad::VadWorkerPolicy::UseWorker => { + let worker = self + .silero_vad_worker + .as_ref() + .expect("policy checked worker"); + 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 if enqueued { - crate::vad::VadOutput { - probability: 0.0, - speech: false, - } - } else { + } + crate::vad::VadWorkerPolicy::UseFallback => { used_fallback_vad = true; self.mark_vad_fallback_active(vad_backend); crate::vad::VoiceActivityDetector::process_10ms( @@ -441,10 +414,10 @@ impl AndroidCaptureState { &frame, ) } - } else { - used_fallback_vad = true; - self.mark_vad_fallback_active(vad_backend); - crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, &frame) + crate::vad::VadWorkerPolicy::NotModelBacked => crate::vad::VadOutput { + probability: 1.0, + speech: true, + }, } } else { crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, &frame) @@ -472,21 +445,19 @@ impl AndroidCaptureState { 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 - })); + if crate::capture_accumulator::append_processed_i16_bounded( + &mut self.pcm_accum, + &frame, + self.mic_gain, + ) { + self.audio_processing_stats.increment_callback_xrun(); } } } struct InputCallback { state: Arc>, + audio_processing_stats: Arc, event_tx: BackendEventTx, } @@ -498,9 +469,13 @@ impl AudioInputCallback for InputCallback { _stream: &mut dyn AudioInputStreamSafe, frames: &[i16], ) -> DataCallbackResult { - let _ = catch_unwind(AssertUnwindSafe(|| { - if let Ok(mut s) = self.state.lock() { - s.ingest_i16(frames); + let _ = catch_unwind(AssertUnwindSafe(|| match self.state.try_lock() { + Ok(mut s) => s.ingest_i16(frames), + Err(std::sync::TryLockError::WouldBlock) => { + self.audio_processing_stats.increment_callback_xrun(); + } + Err(std::sync::TryLockError::Poisoned(e)) => { + warn!(target: "chanora_audio", "android: capture state poisoned: {e}"); } })); DataCallbackResult::Continue @@ -522,8 +497,7 @@ impl AudioInputCallback for InputCallback { // writes stereo f32 directly to the Oboe output buffer. struct OutputCallback { - handler: AudioHandler, - event_consumer: crate::audio_event_queue::AudioEventConsumer, + pcm_consumer: crate::android_render_ring::AndroidRenderRingConsumer, output_gain: Arc, output_muted: Arc, event_tx: BackendEventTx, @@ -542,31 +516,8 @@ impl AudioOutputCallback for OutputCallback { frames: &mut [(f32, f32)], ) -> DataCallbackResult { let _ = catch_unwind(AssertUnwindSafe(|| { - let buf: &mut [f32] = - bytemuck::cast_slice_mut::<(f32, f32), f32>(frames); - for s in buf.iter_mut() { - *s = 0.0; - } - for cmd in self.event_consumer.drain_controls() { - match cmd { - AudioCommand::SetVolume(id, vol) => { - if let Some(q) = self.handler.get_mut_queues().get_mut(&id) { - q.volume = vol; - } - } - AudioCommand::RemoveClient(id) => { - self.handler.get_mut_queues().remove(&id); - } - } - } - - for pkt in self.event_consumer.drain_packets(50) { - if let Err(e) = self.handler.handle_packet(pkt.client_id, pkt.data) { - debug!(target: "chanora_audio", error = %e, "decode failed"); - } - } - - let _ = self.handler.fill_buffer(buf); + let buf: &mut [f32] = bytemuck::cast_slice_mut::<(f32, f32), f32>(frames); + self.pcm_consumer.drain_into_zero_filling(buf); let gain = f32::from_bits(self.output_gain.load(Ordering::Relaxed)); let muted = self.output_muted.load(Ordering::Relaxed); if muted { @@ -614,6 +565,7 @@ impl AudioOutputCallback for OutputCallback { pub struct AndroidVoiceUnit { input: Option>, output: Option>, + render_producer_shutdown: Arc, // Recorded achieved values (SDD-112). input_perf: AchievedPerformanceMode, @@ -706,6 +658,7 @@ impl AndroidVoiceUnit { let input_cb = InputCallback { state: capture_state.clone(), + audio_processing_stats: audio_processing_stats.clone(), event_tx: event_tx.clone(), }; let input_builder = input_builder.set_callback(input_cb); @@ -722,7 +675,12 @@ impl AndroidVoiceUnit { error = ?e, "android: primary input stream open failed; entering fallback ladder" ); - match Self::open_input_fallback(cfg, &event_tx, capture_state.clone()) { + match Self::open_input_fallback( + cfg, + &event_tx, + capture_state.clone(), + audio_processing_stats.clone(), + ) { Ok(s) => Some(s), Err(fallback_err) => { warn!( @@ -789,9 +747,10 @@ impl AndroidVoiceUnit { let render_ref_for_output = render_ref_buf.clone(); let event_queue = params.event_producer.queue(); + let render_ring = + crate::android_render_ring::AndroidRenderRing::new(ANDROID_RENDER_RING_CAPACITY); let output_cb = OutputCallback { - handler: params.handler, - event_consumer: AudioEventQueue::consumer(&event_queue), + pcm_consumer: render_ring.consumer(), output_gain: params.output_gain.clone(), output_muted: params.output_muted.clone(), event_tx: event_tx.clone(), @@ -813,8 +772,7 @@ impl AndroidVoiceUnit { Self::open_output_fallback( cfg, &event_tx, - AudioHandler::new(), - AudioEventQueue::consumer(&event_queue), + render_ring.consumer(), params.output_gain.clone(), params.output_muted.clone(), audio_processing_stats.clone(), @@ -822,6 +780,11 @@ impl AndroidVoiceUnit { )? } }; + let render_producer_shutdown = Self::spawn_render_producer( + params.handler, + AudioEventQueue::consumer(&event_queue), + render_ring.producer(), + ); let output_frames_per_burst = output_stream.get_frames_per_burst(); if output_frames_per_burst > 0 { @@ -978,6 +941,7 @@ impl AndroidVoiceUnit { Ok(Self { input: input_stream, output: Some(output_stream), + render_producer_shutdown, input_perf, input_share, output_perf, @@ -995,6 +959,7 @@ impl AndroidVoiceUnit { cfg: &AndroidVoiceStreamConfig, event_tx: &BackendEventTx, capture_state: Arc>, + audio_processing_stats: Arc, ) -> Result, BackendError> { // SDD-112 items 6 & 7: explore (preset × sharing) independently // via the pure helpers in `mobile_voice_backend`. Primary @@ -1031,6 +996,7 @@ impl AndroidVoiceUnit { }; let cb = InputCallback { state: capture_state.clone(), + audio_processing_stats: audio_processing_stats.clone(), event_tx: event_tx.clone(), }; let builder = AudioStreamBuilder::default() @@ -1066,16 +1032,14 @@ impl AndroidVoiceUnit { fn open_output_fallback( cfg: &AndroidVoiceStreamConfig, event_tx: &BackendEventTx, - handler: AudioHandler, - event_consumer: crate::audio_event_queue::AudioEventConsumer, + pcm_consumer: crate::android_render_ring::AndroidRenderRingConsumer, output_gain: Arc, output_muted: Arc, audio_processing_stats: Arc, render_reference: Arc, ) -> Result, BackendError> { let cb = OutputCallback { - handler, - event_consumer, + pcm_consumer, output_gain, output_muted, event_tx: event_tx.clone(), @@ -1100,6 +1064,50 @@ impl AndroidVoiceUnit { .map_err(|e| BackendError::OpenFailed(format!("output fallback: {e:?}"))) } + fn spawn_render_producer( + mut handler: AudioHandler, + event_consumer: crate::audio_event_queue::AudioEventConsumer, + pcm_producer: crate::android_render_ring::AndroidRenderRingProducer, + ) -> Arc { + let shutdown = Arc::new(AtomicBool::new(false)); + let shutdown_for_task = shutdown.clone(); + tokio::spawn(async move { + let mut pull_scratch = vec![0.0_f32; ANDROID_RENDER_PULL_SAMPLES]; + let mut interval = tokio::time::interval(std::time::Duration::from_millis(20)); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + loop { + interval.tick().await; + if shutdown_for_task.load(Ordering::Relaxed) { + break; + } + + for cmd in event_consumer.drain_controls() { + match cmd { + AudioCommand::SetVolume(id, vol) => { + if let Some(q) = handler.get_mut_queues().get_mut(&id) { + q.volume = vol; + } + } + AudioCommand::RemoveClient(id) => { + handler.get_mut_queues().remove(&id); + } + } + } + + for pkt in event_consumer.drain_packets(50) { + if let Err(e) = handler.handle_packet(pkt.client_id, pkt.data) { + debug!(target: "chanora_audio", error = %e, "decode failed"); + } + } + + pull_scratch.fill(0.0); + let _ = handler.fill_buffer(&mut pull_scratch); + pcm_producer.push_frame_lossy(&pull_scratch); + } + }); + shutdown + } + /// Clone of the event sender, for JNI focus / SCO listeners /// registered on the engine's behalf. pub fn event_sender(&self) -> BackendEventTx { @@ -1151,6 +1159,7 @@ impl MobileVoiceAudioBackend for AndroidVoiceUnit { fn close(&mut self) -> Result<(), BackendError> { // SDD-115 reverse order: release hardware effects FIRST, // then close streams. + self.render_producer_shutdown.store(true, Ordering::Relaxed); release_hardware_effects(&mut self.hw_effects); self.stop().ok(); // Dropping the Option drops the underlying AudioStreamAsync @@ -1208,6 +1217,7 @@ impl Drop for AndroidVoiceUnit { // Wrap in catch_unwind so a panic during Drop cannot unwind // into the JVM (SDD-115 callback safety). let _ = catch_unwind(AssertUnwindSafe(|| { + self.render_producer_shutdown.store(true, Ordering::Relaxed); release_hardware_effects(&mut self.hw_effects); // SDD-116: clear the diagnostics slot on Drop too. clear_android_audio_diagnostics(); diff --git a/crates/chanora_audio/src/capture_accumulator.rs b/crates/chanora_audio/src/capture_accumulator.rs new file mode 100644 index 0000000..183f0b2 --- /dev/null +++ b/crates/chanora_audio/src/capture_accumulator.rs @@ -0,0 +1,118 @@ +pub(crate) fn append_processed_i16_bounded( + pcm_accum: &mut Vec, + frame: &[f32], + gain: f32, +) -> bool { + if (gain - 1.0).abs() < f32::EPSILON { + for src in frame.iter().copied() { + if pcm_accum.len() == pcm_accum.capacity() { + return true; + } + pcm_accum.push(crate::frame::f32_to_i16(src)); + } + } else { + for src in frame.iter().copied() { + if pcm_accum.len() == pcm_accum.capacity() { + return true; + } + let scaled = (crate::frame::f32_to_i16(src) as f32) * gain; + pcm_accum.push(scaled.clamp(i16::MIN as f32, i16::MAX as f32) as i16); + } + } + false +} + +pub(crate) fn append_i16_bounded(pcm_accum: &mut Vec, frame: &[i16]) -> bool { + for src in frame.iter().copied() { + if pcm_accum.len() == pcm_accum.capacity() { + return true; + } + pcm_accum.push(src); + } + false +} + +#[cfg(test)] +mod tests { + use super::{append_i16_bounded, append_processed_i16_bounded}; + + #[test] + fn append_processed_i16_bounded_does_not_grow_when_full() { + let frame = [0.25_f32; crate::frame::FRAME_10MS_SAMPLES]; + let mut accum = Vec::with_capacity(crate::frame::FRAME_10MS_SAMPLES / 2); + let warmed_capacity = accum.capacity(); + let warmed_ptr = accum.as_ptr(); + + let dropped = append_processed_i16_bounded(&mut accum, &frame, 1.0); + + assert!(dropped); + assert_eq!(accum.len(), warmed_capacity); + assert_eq!(accum.capacity(), warmed_capacity); + assert_eq!(accum.as_ptr(), warmed_ptr); + } + + #[test] + fn append_processed_i16_bounded_preserves_expected_10ms_append() { + let frame = [0.25_f32; crate::frame::FRAME_10MS_SAMPLES]; + let mut accum = Vec::with_capacity(crate::frame::FRAME_20MS_SAMPLES * 2); + let warmed_capacity = accum.capacity(); + let warmed_ptr = accum.as_ptr(); + + let dropped = append_processed_i16_bounded(&mut accum, &frame, 1.0); + + assert!(!dropped); + assert_eq!(accum.len(), crate::frame::FRAME_10MS_SAMPLES); + assert_eq!(accum.capacity(), warmed_capacity); + assert_eq!(accum.as_ptr(), warmed_ptr); + } + + #[test] + fn append_i16_bounded_does_not_grow_when_preroll_exceeds_capacity() { + let frame = [7_i16; crate::frame::FRAME_10MS_SAMPLES]; + let mut accum = Vec::with_capacity(crate::frame::FRAME_10MS_SAMPLES / 2); + let warmed_capacity = accum.capacity(); + let warmed_ptr = accum.as_ptr(); + + let dropped = append_i16_bounded(&mut accum, &frame); + + assert!(dropped); + assert_eq!(accum.len(), warmed_capacity); + assert_eq!(accum.capacity(), warmed_capacity); + assert_eq!(accum.as_ptr(), warmed_ptr); + } + + #[test] + fn append_i16_bounded_preserves_expected_10ms_append() { + let frame = [7_i16; crate::frame::FRAME_10MS_SAMPLES]; + let mut accum = Vec::with_capacity(crate::frame::FRAME_20MS_SAMPLES * 2); + let warmed_capacity = accum.capacity(); + let warmed_ptr = accum.as_ptr(); + + let dropped = append_i16_bounded(&mut accum, &frame); + + assert!(!dropped); + assert_eq!(accum.len(), crate::frame::FRAME_10MS_SAMPLES); + assert_eq!(accum.capacity(), warmed_capacity); + assert_eq!(accum.as_ptr(), warmed_ptr); + } + + #[test] + fn append_i16_bounded_preserves_full_vad_preroll_window() { + let frame = [7_i16; crate::frame::FRAME_10MS_SAMPLES]; + let mut accum = Vec::with_capacity(crate::frame::FRAME_10MS_SAMPLES * 16); + let warmed_capacity = accum.capacity(); + let warmed_ptr = accum.as_ptr(); + + for _ in 0..16 { + assert!(!append_i16_bounded(&mut accum, &frame)); + } + + assert_eq!(accum.len(), crate::frame::FRAME_10MS_SAMPLES * 16); + assert_eq!(accum.capacity(), warmed_capacity); + assert_eq!(accum.as_ptr(), warmed_ptr); + assert!(append_i16_bounded(&mut accum, &frame)); + assert_eq!(accum.len(), warmed_capacity); + assert_eq!(accum.capacity(), warmed_capacity); + assert_eq!(accum.as_ptr(), warmed_ptr); + } +} diff --git a/crates/chanora_audio/src/capture_resampler.rs b/crates/chanora_audio/src/capture_resampler.rs new file mode 100644 index 0000000..c9d6f4e --- /dev/null +++ b/crates/chanora_audio/src/capture_resampler.rs @@ -0,0 +1,105 @@ +pub(crate) struct CaptureResampleResult { + pub(crate) output_len: usize, + pub(crate) dropped: bool, +} + +pub(crate) fn resample_capture_to_48k( + samples: &[i16], + input_sample_rate_hz: u32, + resample_pos: &mut f64, + resample_last: &mut i16, + scratch: &mut Vec, +) -> CaptureResampleResult { + scratch.clear(); + if samples.is_empty() { + return CaptureResampleResult { + output_len: 0, + dropped: false, + }; + } + + let ratio = input_sample_rate_hz.max(1) as f64 / crate::frame::SAMPLE_RATE_HZ as f64; + let mut pos = *resample_pos; + let mut dropped = false; + while pos < samples.len() as f64 { + let i = pos.floor() as isize; + let frac = pos - i as f64; + let a = if i <= 0 { + *resample_last as f64 + } else { + samples[(i - 1) as usize] as f64 + }; + let b = if i < samples.len() as isize { + samples[i as usize] as f64 + } else { + a + }; + let value = (a + frac * (b - a)) + .round() + .clamp(i16::MIN as f64, i16::MAX as f64) as i16; + if scratch.len() < scratch.capacity() { + scratch.push(value); + } else { + dropped = true; + } + pos += ratio; + } + *resample_pos = pos - samples.len() as f64; + *resample_last = *samples.last().unwrap_or(resample_last); + CaptureResampleResult { + output_len: scratch.len(), + dropped, + } +} + +#[cfg(test)] +mod android_voice_unit_resampler_tests { + use super::resample_capture_to_48k; + + #[test] + fn android_voice_unit_resampler_reuses_scratch_without_capacity_growth() { + let samples: Vec = (0..882).map(|i| i as i16).collect(); + let mut pos = 0.0; + let mut last = 0_i16; + let mut scratch = Vec::with_capacity(960); + + let first = resample_capture_to_48k(&samples, 44_100, &mut pos, &mut last, &mut scratch); + let first_len = first.output_len; + assert_eq!(first_len, 960); + assert!(!first.dropped); + assert_eq!(scratch.len(), first_len); + let warmed_capacity = scratch.capacity(); + let warmed_ptr = scratch.as_ptr(); + + for _ in 0..8 { + let result = + resample_capture_to_48k(&samples, 44_100, &mut pos, &mut last, &mut scratch); + let len = result.output_len; + assert_eq!(len, scratch.len()); + assert!(len >= 959 && len <= 960); + assert!(!result.dropped); + assert_eq!(scratch.capacity(), warmed_capacity); + assert_eq!(scratch.as_ptr(), warmed_ptr); + } + } + + #[test] + fn android_voice_unit_resampler_truncates_oversized_burst_without_capacity_growth() { + let samples: Vec = (0..4_800).map(|i| i as i16).collect(); + let mut pos = 0.0; + let mut last = 0_i16; + let mut scratch = Vec::with_capacity(960); + let warmed_capacity = scratch.capacity(); + let warmed_ptr = scratch.as_ptr(); + + let result = resample_capture_to_48k(&samples, 48_000, &mut pos, &mut last, &mut scratch); + + assert_eq!(result.output_len, warmed_capacity); + assert!(result.dropped); + assert_eq!(scratch.len(), warmed_capacity); + assert_eq!(scratch.capacity(), warmed_capacity); + assert_eq!(scratch.as_ptr(), warmed_ptr); + assert_eq!(pos, 0.0); + assert_eq!(last, *samples.last().unwrap()); + } +} diff --git a/crates/chanora_audio/src/debug_wav.rs b/crates/chanora_audio/src/debug_wav.rs index 1865141..613803f 100644 --- a/crates/chanora_audio/src/debug_wav.rs +++ b/crates/chanora_audio/src/debug_wav.rs @@ -314,6 +314,15 @@ mod tests { rec.stop(); } + #[test] + fn ios_raw_debug_wav_does_not_push_from_realtime_callback() { + let src = include_str!("ios_raw_unit.rs"); + assert!( + !src.contains("push_raw_mic") && !src.contains("push_processed_mic"), + "ios raw callbacks must not call WavDebugRecorder push_*_mic until it has a preallocated handoff" + ); + } + #[test] fn wav_header_is_44_bytes() { // Write to a temp file to test the header. diff --git a/crates/chanora_audio/src/engine.rs b/crates/chanora_audio/src/engine.rs index ec61510..40d79e7 100644 --- a/crates/chanora_audio/src/engine.rs +++ b/crates/chanora_audio/src/engine.rs @@ -981,7 +981,6 @@ impl AudioEngine { Ok(Self { transmit_gate, - frames_sent, frames_received, output_gain, output_muted, @@ -1085,11 +1084,25 @@ impl AudioEngine { audio_processing_stats: audio_processing_stats.clone(), }; let mut android_voice_unit = - crate::android_voice_unit::AndroidVoiceUnit::open(&cfg_av, params).map_err(|e| { - AudioError::Backend(format!("android: failed to open Oboe voice unit: {e}")) - })?; + match crate::android_voice_unit::AndroidVoiceUnit::open(&cfg_av, params) { + Ok(unit) => unit, + Err(e) => { + Self::rollback_android_startup_resources(&mut audio_mode_stack); + return Err(AudioError::Backend(format!( + "android: failed to open Oboe voice unit: {e}" + ))); + } + }; if let Err(e) = android_voice_unit.start() { + if let Err(close_err) = android_voice_unit.close() { + warn!( + target: "chanora_audio", + error = %close_err, + "android: AndroidVoiceUnit::close failed during startup rollback" + ); + } + Self::rollback_android_startup_resources(&mut audio_mode_stack); return Err(AudioError::Backend(format!( "android: failed to start Oboe voice unit: {e}" ))); @@ -1326,6 +1339,43 @@ impl AudioEngine { }) } + #[cfg(target_os = "android")] + fn rollback_android_startup_resources(audio_mode_stack: &mut crate::mode_stack::ModeStack) { + crate::android_voice_unit::chanora_android_stop_bluetooth_sco(); + crate::android_voice_unit::chanora_android_abandon_audio_focus(); + match release_android_audio_mode_for_startup_rollback(audio_mode_stack) { + crate::mode_stack::ModeRelease::LastRelease { prior } => { + match android_set_audio_mode(prior) { + Ok(()) => info!( + target: "chanora_audio", + restored_mode = prior, + "android: AudioManager mode restored during startup rollback" + ), + Err(e) => warn!( + target: "chanora_audio", + error = %e, + restored_mode = prior, + "android: failed to restore AudioManager mode during startup rollback" + ), + } + } + crate::mode_stack::ModeRelease::StillHeld => { + info!( + target: "chanora_audio", + "android: audio mode still held during startup rollback" + ); + } + crate::mode_stack::ModeRelease::AlreadyReleased => {} + } + + if crate::android_voice_unit::chanora_android_stop_voice_service() { + info!( + target: "chanora_audio", + "android: voice foreground service stopped during startup rollback" + ); + } + } + /// Stop the engine. Idempotent. pub fn stop(&mut self) { if let Some(tx) = self.shutdown_tx.take() { @@ -1720,6 +1770,33 @@ impl Drop for AudioEngine { } } +#[cfg(any(test, target_os = "android"))] +fn release_android_audio_mode_for_startup_rollback( + audio_mode_stack: &mut crate::mode_stack::ModeStack, +) -> crate::mode_stack::ModeRelease { + audio_mode_stack.release() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn android_startup_rollback_releases_acquired_mode_snapshot() { + let mut stack = crate::mode_stack::ModeStack::new(); + let _ = stack.acquire(7); + + let release = release_android_audio_mode_for_startup_rollback(&mut stack); + + assert_eq!( + release, + crate::mode_stack::ModeRelease::LastRelease { prior: 7 } + ); + assert_eq!(stack.refcount(), 0); + assert_eq!(stack.snapshot(), None); + } +} + // ---------- Capture pipeline ---------- #[cfg(not(any(target_os = "ios", target_os = "macos", target_os = "android")))] @@ -1762,9 +1839,12 @@ fn try_open_capture( in_sample_rate, in_channels, mic_gain, - voice_out_tx, + crate::opus_voice::start_out_packet_worker( + voice_out_tx, + frames_sent.clone(), + "cpal-capture", + )?, transmit_active, - frames_sent, audio_processing_stats, ))); @@ -1799,11 +1879,10 @@ struct CaptureState { /// Linux ALSA defaults). resample_last: f32, opus_out: [u8; crate::opus_voice::MAX_OPUS_FRAME], - voice_out_tx: mpsc::Sender, + voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender, /// The PTT transmission gate. Read once per outbound frame; the /// CaptureState never mutates this flag. transmit_active: Arc, - frames_sent: Arc, /// Pre-allocated mono downmix buffer. Resized in-place each /// callback; `clear()` retains capacity. SDD-094 realtime-thread /// invariant: this avoids the heap allocation that the prior fix @@ -1836,8 +1915,7 @@ struct CaptureState { /// Time-based gating is robust to cpal buffer-size and sample-rate /// changes that a fixed callback-count would not be. #[cfg(not(any(target_os = "ios", target_os = "macos", target_os = "android")))] -const LEVEL_METER_INTERVAL: std::time::Duration = - std::time::Duration::from_millis(33); +const LEVEL_METER_INTERVAL: std::time::Duration = std::time::Duration::from_millis(33); #[cfg(not(any(target_os = "ios", target_os = "macos", target_os = "android")))] impl CaptureState { @@ -1846,9 +1924,8 @@ impl CaptureState { in_sample_rate: u32, in_channels: usize, mic_gain: f32, - voice_out_tx: mpsc::Sender, + voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender, transmit_active: Arc, - frames_sent: Arc, audio_processing_stats: Arc, ) -> Self { Self { @@ -1862,7 +1939,6 @@ impl CaptureState { opus_out: [0u8; crate::opus_voice::MAX_OPUS_FRAME], voice_out_tx, transmit_active, - frames_sent, mono_scratch: Vec::with_capacity(4096), frame_scratch: Vec::with_capacity(FRAME_SAMPLES), audio_processing_stats, @@ -1959,7 +2035,6 @@ impl CaptureState { Ok(len) => { crate::opus_voice::send_voip_frame( &self.voice_out_tx, - &self.frames_sent, &self.opus_out, len, || { diff --git a/crates/chanora_audio/src/ios_raw_unit.rs b/crates/chanora_audio/src/ios_raw_unit.rs index c774ebd..54cb367 100644 --- a/crates/chanora_audio/src/ios_raw_unit.rs +++ b/crates/chanora_audio/src/ios_raw_unit.rs @@ -42,13 +42,11 @@ mod inner { 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; @@ -63,43 +61,10 @@ mod inner { /// 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 {} + type RenderReferenceBuffer = crate::render_reference::RenderReferenceBuffer<480, 4>; + type RenderReferenceFrameAccumulator = + crate::render_reference::RenderReferenceFrameAccumulator<480>; + const RAW_RENDER_SCRATCH_FRAMES: usize = 1024; // ------------------------------------------------------------------ // // Capture pipeline state // @@ -109,10 +74,9 @@ mod inner { encoder: OpusEncoder, pcm_accum: Vec, opus_out: [u8; crate::opus_voice::MAX_OPUS_FRAME], - voice_out_tx: mpsc::Sender, + voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender, transmit_active: Arc, output_muted: Arc, - frames_sent: Arc, mic_gain: f32, voice_activity_selector: Option>, vad_detector: crate::vad::WebRtcFallbackVad, @@ -128,7 +92,6 @@ mod inner { pending_10ms: [i16; crate::frame::FRAME_10MS_SAMPLES], pending_10ms_len: usize, fallback_warned_backend: Option, - wav_recorder: Option>, } impl RawCaptureState { @@ -146,10 +109,13 @@ mod inner { 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(), + voice_out_tx: crate::opus_voice::start_out_packet_worker( + params.voice_out_tx.clone(), + params.frames_sent.clone(), + "ios-raw", + )?, 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(), @@ -166,7 +132,6 @@ mod inner { pending_10ms: [0_i16; crate::frame::FRAME_10MS_SAMPLES], pending_10ms_len: 0, fallback_warned_backend: None, - wav_recorder: None, }) } @@ -204,6 +169,7 @@ mod inner { if self.pending_10ms_len == crate::frame::FRAME_10MS_SAMPLES { let frame = self.pending_10ms; self.process_10ms_capture_frame(&frame); + self.encode_complete_20ms_frames(); self.pending_10ms_len = 0; } } @@ -213,7 +179,10 @@ mod inner { return; } - // Encode complete 20 ms Opus frames. + self.encode_complete_20ms_frames(); + } + + fn encode_complete_20ms_frames(&mut self) { 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]); @@ -223,7 +192,6 @@ mod inner { Ok(len) => { crate::opus_voice::send_voip_frame( &self.voice_out_tx, - &self.frames_sent, &self.opus_out, len, || { @@ -253,20 +221,18 @@ mod inner { } 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); - } + // Debug WAV mic taps are intentionally unavailable on iOS raw + // realtime callbacks until WavDebugRecorder supports a + // preallocated handoff path; the current recorder push path + // allocates per 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); - } + // Processed-mic debug WAV capture is disabled for the same + // realtime allocation reason as the raw-mic tap above. let voice_activity_mode = self .voice_activity_selector @@ -300,14 +266,9 @@ mod inner { 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); - } + self.silero_coreml_worker = None; + self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx); + self.audio_processing_stats.set_vad_fallback_active(true); } else { self.silero_coreml_worker = None; self.audio_processing_stats.set_vad_fallback_active(false); @@ -325,20 +286,38 @@ mod inner { 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, + match crate::vad::callback_vad_worker_policy( + voice_activity_mode, + vad_backend, + self.silero_coreml_worker.is_some(), + ) { + crate::vad::VadWorkerPolicy::UseWorker => { + let worker = self + .silero_coreml_worker + .as_ref() + .expect("policy checked worker"); + 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 if enqueued { - crate::vad::VadOutput { - probability: 0.0, - speech: false, - } - } else { + } + crate::vad::VadWorkerPolicy::UseFallback => { used_fallback_vad = true; self.mark_vad_fallback_active(vad_backend); crate::vad::VoiceActivityDetector::process_10ms( @@ -346,13 +325,10 @@ mod inner { &frame, ) } - } else { - used_fallback_vad = true; - self.mark_vad_fallback_active(vad_backend); - crate::vad::VoiceActivityDetector::process_10ms( - &mut self.vad_detector, - &frame, - ) + crate::vad::VadWorkerPolicy::NotModelBacked => crate::vad::VadOutput { + probability: 1.0, + speech: true, + }, } } else { crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, &frame) @@ -378,15 +354,12 @@ mod inner { 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 - })); + if crate::capture_accumulator::append_processed_i16_bounded( + &mut self.pcm_accum, + &frame, + self.mic_gain, + ) { + self.audio_processing_stats.increment_callback_xrun(); } } } @@ -446,7 +419,9 @@ mod inner { }) .map_err(|e| AudioError::Backend(format!("remoteio input cb: {e}")))?; - let mut scratch: Vec = Vec::with_capacity(2048); + let mut scratch = [0.0_f32; RAW_RENDER_SCRATCH_FRAMES * 2]; + let mut mono = [0.0_f32; RAW_RENDER_SCRATCH_FRAMES]; + let mut render_ref_accum = RenderReferenceFrameAccumulator::new(); let handler = params.handler.clone(); let output_gain = params.output_gain.clone(); let output_muted = params.output_muted.clone(); @@ -455,9 +430,10 @@ mod inner { 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); + let process_n = n.min(RAW_RENDER_SCRATCH_FRAMES); + let stereo_n = process_n * 2; + if n > RAW_RENDER_SCRATCH_FRAMES { + stats_render.increment_callback_xrun(); } scratch[..stereo_n].fill(0.0); @@ -475,22 +451,26 @@ mod inner { } // 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], + &mut mono[..process_n], ); - render_ref_buf.write(&ref_frame); + render_ref_accum.push_mono_samples(&mono[..process_n], |frame| { + render_ref_buf.write(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( + let mix_stats = crate::voice_render::downmix_stereo_f32_to_interleaved_i16( &scratch[..stereo_n], - out, + &mut out[..process_n], + 1, gain, muted, ); + if process_n < n { + out[process_n..].fill(0); + } if mix_stats.clipped_samples > 0 { stats_render.add_clipped_samples(mix_stats.clipped_samples); } diff --git a/crates/chanora_audio/src/ios_voice_unit.rs b/crates/chanora_audio/src/ios_voice_unit.rs index cbc0129..44a4e2d 100644 --- a/crates/chanora_audio/src/ios_voice_unit.rs +++ b/crates/chanora_audio/src/ios_voice_unit.rs @@ -66,7 +66,7 @@ //! * AVAudioSession category / mode configuration — Swift owns the //! session (it must be set up before Flutter loads). -use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use audiopus::coder::Encoder as OpusEncoder; @@ -75,12 +75,10 @@ use coreaudio::audio_unit::render_callback::{self, data}; use coreaudio::audio_unit::IOType; use coreaudio::audio_unit::{AudioUnit, Element, SampleFormat, Scope, StreamFormat}; use crossbeam::queue::ArrayQueue; -use tokio::sync::mpsc; use tracing::{debug, error, info, warn}; use crate::mobile_voice_backend::VoiceAudioParams; use crate::AudioError; -use chanora_protocol::OutPacket; /// Sample rate every layer above us assumes. Matches the Opus /// encoder rate, the `tsclientlib::AudioHandler` mix rate, and the @@ -105,6 +103,15 @@ const INPUT_BUS: Element = Element::Input; /// when the VAD gate opens (VAD_004 / pre_roll_ms=160). const PRE_ROLL_FRAMES: usize = 16; +/// Enough room for the 160 ms VAD pre-roll plus a few jitter frames, without +/// growing inside the input callback. +const CAPTURE_ACCUM_CAPACITY_SAMPLES: usize = crate::frame::FRAME_10MS_SAMPLES * 20; + +/// Fixed iOS render scratch capacity. Larger callback requests are truncated +/// to this capacity and the remaining output is silence. +#[cfg_attr(not(target_os = "ios"), allow(dead_code))] +const IOS_RENDER_SCRATCH_FRAMES: usize = 4096; + /// Capture pipeline state owned by the VPIO input callback. The /// AudioUnit hands us 48 kHz signed-int16 mono PCM directly (no /// downmix or resample needed — VPIO's hardware-side mix-down @@ -132,10 +139,9 @@ struct IosCaptureState { /// jitter without reallocating. pcm_accum: Vec, opus_out: [u8; crate::opus_voice::MAX_OPUS_FRAME], - voice_out_tx: mpsc::Sender, + voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender, transmit_active: Arc, output_muted: Arc, - frames_sent: Arc, mic_gain: f32, voice_activity_selector: Option>, vad_detector: crate::vad::WebRtcFallbackVad, @@ -154,7 +160,6 @@ struct IosCaptureState { pre_roll_count: usize, pre_roll_flushed: bool, capture_frame_seq: u64, - wav_recorder: Arc>>>, } impl IosCaptureState { @@ -162,20 +167,20 @@ impl IosCaptureState { /// Encoder configuration is the same as cpal-side /// `try_open_capture` (engine.rs) so audio quality is platform- /// neutral. - fn new( - params: &VoiceAudioParams, - wav_recorder: Arc>>>, - ) -> Result { + fn new(params: &VoiceAudioParams) -> Result { let encoder = crate::opus_voice::new_voip_encoder("ios VPIO")?; Ok(Self { encoder, - pcm_accum: Vec::with_capacity(crate::frame::FRAME_20MS_SAMPLES * 2), + pcm_accum: Vec::with_capacity(CAPTURE_ACCUM_CAPACITY_SAMPLES), opus_out: [0u8; crate::opus_voice::MAX_OPUS_FRAME], - voice_out_tx: params.voice_out_tx.clone(), + voice_out_tx: crate::opus_voice::start_out_packet_worker( + params.voice_out_tx.clone(), + params.frames_sent.clone(), + "ios-vpio", + )?, 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(), @@ -193,7 +198,6 @@ impl IosCaptureState { pre_roll_count: 0, pre_roll_flushed: false, capture_frame_seq: 0, - wav_recorder, }) } @@ -267,7 +271,6 @@ impl IosCaptureState { Ok(len) => { crate::opus_voice::send_voip_frame( &self.voice_out_tx, - &self.frames_sent, &self.opus_out, len, || { @@ -298,25 +301,9 @@ impl IosCaptureState { } let input_dbfs = crate::frame::dbfs(&frame); - // WAV tap: raw mic (before processing, DIAG_002). - if let Ok(guard) = self.wav_recorder.try_lock() { - if let Some(rec) = guard.as_ref() { - rec.push_raw_mic(&frame); - } - } - // Read config once per frame (try_lock: non-blocking, falls back to // last-known values if the lock is contended — safe to miss one frame). - let ( - run_ns, - run_agc, - run_hpf, - vad_backend, - vad_hangover, - debug_wav_dump_enabled, - route, - processing_backend, - ) = self + let (run_ns, run_agc, run_hpf, vad_backend, vad_hangover, debug_wav_dump_enabled) = self .audio_processing_config .try_lock() .map(|cfg| { @@ -332,8 +319,6 @@ impl IosCaptureState { cfg.vad_backend, cfg.vad_hangover_ms, cfg.debug_wav_dump_enabled, - cfg.route, - cfg.processing_backend, ) }) .unwrap_or(( @@ -343,10 +328,13 @@ impl IosCaptureState { crate::VadBackend::WebrtcVad, crate::voice_activity::VAD_HANGOVER_MS, false, - crate::AudioRoute::Unknown, - crate::AudioBackend::PlatformVoiceProcessing, )); + // VPIO realtime callbacks cannot use WavDebugRecorder today: its push + // path allocates per frame. Debug WAV capture is intentionally disabled + // here until the recorder can hand off preallocated frames. + let _ = debug_wav_dump_enabled; + let voice_activity_mode = self .voice_activity_selector .as_ref() @@ -359,32 +347,13 @@ impl IosCaptureState { self.audio_processing_stats.set_vad_fallback_active(false); } - // Switch VAD backend only while VoiceActivity mode is active. - if let Ok(mut recorder_guard) = self.wav_recorder.try_lock() { - if debug_wav_dump_enabled { - if recorder_guard.is_none() { - *recorder_guard = Some(crate::debug_wav::WavDebugRecorder::start( - route, - processing_backend, - )); - } - } else if let Some(recorder) = recorder_guard.take() { - recorder.stop(); - } - } - 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); - } + self.silero_coreml_worker = None; + self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx); + self.audio_processing_stats.set_vad_fallback_active(true); } else { self.silero_coreml_worker = None; self.audio_processing_stats.set_vad_fallback_active(false); @@ -437,20 +406,38 @@ impl IosCaptureState { 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, + match crate::vad::callback_vad_worker_policy( + voice_activity_mode, + vad_backend, + self.silero_coreml_worker.is_some(), + ) { + crate::vad::VadWorkerPolicy::UseWorker => { + let worker = self + .silero_coreml_worker + .as_ref() + .expect("policy checked worker"); + 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(crate::VadBackend::SileroOnnx); + crate::vad::VoiceActivityDetector::process_10ms( + &mut self.vad_detector, + &frame, + ) } - } else if enqueued { - crate::vad::VadOutput { - probability: 0.0, - speech: false, - } - } else { + } + crate::vad::VadWorkerPolicy::UseFallback => { used_fallback_vad = true; self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx); crate::vad::VoiceActivityDetector::process_10ms( @@ -458,10 +445,10 @@ impl IosCaptureState { &frame, ) } - } else { - used_fallback_vad = true; - self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx); - crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, &frame) + crate::vad::VadWorkerPolicy::NotModelBacked => crate::vad::VadOutput { + probability: 1.0, + speech: true, + }, } } else { crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, &frame) @@ -484,13 +471,6 @@ impl IosCaptureState { transmit_active, ); - // WAV tap: processed mic (after Rust DSP, DIAG_002). - if let Ok(guard) = self.wav_recorder.try_lock() { - if let Some(rec) = guard.as_ref() { - rec.push_processed_mic(&frame); - } - } - // Convert to i16 for accumulation. let mut pcm_frame = [0_i16; crate::frame::FRAME_10MS_SAMPLES]; if (self.mic_gain - 1.0).abs() < f32::EPSILON { @@ -528,7 +508,13 @@ impl IosCaptureState { let pre_roll_to_emit = self.pre_roll_count.saturating_sub(1); for i in 0..pre_roll_to_emit { let idx = (oldest + i) % PRE_ROLL_FRAMES; - self.pcm_accum.extend_from_slice(&self.pre_roll_buf[idx]); + if crate::capture_accumulator::append_i16_bounded( + &mut self.pcm_accum, + &self.pre_roll_buf[idx], + ) { + self.audio_processing_stats.increment_callback_xrun(); + break; + } } } else if !transmit_active { // Gate closed — reset the flush flag so pre-roll fires again @@ -540,7 +526,9 @@ impl IosCaptureState { return; } - self.pcm_accum.extend_from_slice(&pcm_frame); + if crate::capture_accumulator::append_i16_bounded(&mut self.pcm_accum, &pcm_frame) { + self.audio_processing_stats.increment_callback_xrun(); + } } } @@ -694,9 +682,7 @@ impl IosVoiceUnit { Element::Output, Some(&ducking_config), ) { - tracing::debug!( - "vpio set OtherAudioDuckingConfiguration failed (older OS?): {e}" - ); + tracing::debug!("vpio set OtherAudioDuckingConfiguration failed (older OS?): {e}"); } // Note: we keep VPIO's voice processing chain ENABLED @@ -748,18 +734,7 @@ impl IosVoiceUnit { // scratch are owned by the closure — no Mutex needed // because the input callback is the sole writer/reader on // the audio thread. - let wav_recorder = Arc::new(Mutex::new({ - let cfg = params.audio_processing_config.lock().unwrap().clone(); - if cfg.debug_wav_dump_enabled { - Some(crate::debug_wav::WavDebugRecorder::start( - cfg.route, - cfg.processing_backend, - )) - } else { - None - } - })); - let mut capture_state = IosCaptureState::new(¶ms, wav_recorder.clone())?; + let mut capture_state = IosCaptureState::new(¶ms)?; unit.set_input_callback(move |args: render_callback::Args>| { // VPIO with our pinned stream format delivers @@ -879,11 +854,8 @@ impl IosVoiceUnit { tokio::spawn(async move { let mut pull_scratch: Vec = vec![0.0; PULL_SAMPLES]; - let mut interval = - tokio::time::interval(std::time::Duration::from_millis(20)); - interval.set_missed_tick_behavior( - tokio::time::MissedTickBehavior::Delay, - ); + let mut interval = tokio::time::interval(std::time::Duration::from_millis(20)); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); loop { interval.tick().await; if producer_shutdown_for_task.load(Ordering::Relaxed) { @@ -909,16 +881,16 @@ impl IosVoiceUnit { unit.set_render_callback(move |args: render_callback::Args>| { let render_callback::Args { - data, - num_frames, - .. + data, num_frames, .. } = args; let out: &mut [i16] = data.buffer; let out_channels = data.channels; let needed = num_frames * out_channels; if pcm_ring_consumer.len() < PREBUFFER_SAMPLES { - for sample in &mut out[..needed] { *sample = 0; } + for sample in &mut out[..needed] { + *sample = 0; + } return Ok(()); } @@ -939,10 +911,8 @@ impl IosVoiceUnit { let mono = (l_lim + r_lim) * 0.5; out[base] = (mono.clamp(-1.0, 1.0) * i16::MAX as f32) as i16; } else { - out[base] = - (l_lim.clamp(-1.0, 1.0) * i16::MAX as f32) as i16; - out[base + 1] = - (r_lim.clamp(-1.0, 1.0) * i16::MAX as f32) as i16; + out[base] = (l_lim.clamp(-1.0, 1.0) * i16::MAX as f32) as i16; + out[base + 1] = (r_lim.clamp(-1.0, 1.0) * i16::MAX as f32) as i16; } written_frames += 1; } @@ -951,17 +921,18 @@ impl IosVoiceUnit { let remaining = num_frames - written_frames; for f in 0..remaining { let base = (written_frames + f) * out_channels; - for c in 0..out_channels { out[base + c] = 0; } + for c in 0..out_channels { + out[base + c] = 0; + } } } - let gain = f32::from_bits( - output_gain_for_render.load(Ordering::Relaxed), - ); - let muted = - output_muted_for_render.load(Ordering::Relaxed); + let gain = f32::from_bits(output_gain_for_render.load(Ordering::Relaxed)); + let muted = output_muted_for_render.load(Ordering::Relaxed); if muted { - for sample in &mut out[..needed] { *sample = 0; } + for sample in &mut out[..needed] { + *sample = 0; + } } else if gain != 1.0 { for sample in &mut out[..needed] { *sample = (((*sample as f32) * gain) @@ -972,9 +943,7 @@ impl IosVoiceUnit { Ok(()) }) - .map_err(|e| AudioError::Backend(format!( - "audio unit set render callback: {e}" - )))?; + .map_err(|e| AudioError::Backend(format!("audio unit set render callback: {e}")))?; } // iOS path: direct fill_buffer in callback. iOS VPIO @@ -984,144 +953,112 @@ impl IosVoiceUnit { // producer-task path above. #[cfg(target_os = "ios")] { - let mut scratch_stereo: Vec = vec![0.0; 4096 * 2]; - let handler_for_render = params.handler.clone(); - let output_gain_for_render = params.output_gain.clone(); - let output_muted_for_render = params.output_muted.clone(); - let audio_processing_stats_for_render = params.audio_processing_stats.clone(); - let wav_recorder_for_render = wav_recorder.clone(); - // Level meter decimation: the render callback fires ~93 - // times/sec, but the bridge consumer reads at ~30 Hz. - let mut render_level_decimation: u32 = 0; - // Render-side reference recorder state. Accumulates downmixed - // mono samples until a 10 ms frame is full, then pushes to the - // recorder. `render_recorder_active` tracks whether the - // recorder is currently armed so we can reset the accumulator - // when it goes from off→on (avoids stitching pre-stop tail - // into the post-start head). - let mut render_recorder_active: bool = false; - let mut render_ref_len: usize = 0; - let mut render_ref_accum: [f32; crate::frame::FRAME_10MS_SAMPLES] = - [0.0; crate::frame::FRAME_10MS_SAMPLES]; - // Diagnostic counters sampled every 100 callbacks. - let mut cb_count: u64 = 0; - let mut last_num_frames: usize = 0; - let mut num_frames_changes: u64 = 0; - let mut callbacks_with_audio: u64 = 0; - let mut callbacks_with_silence: u64 = 0; - unit.set_render_callback(move |args: render_callback::Args>| { - let render_callback::Args { - data, - num_frames, - .. - } = args; - let out: &mut [i16] = data.buffer; - let out_channels = data.channels; - // AudioHandler produces 48 kHz stereo f32 (= num_frames * 2 floats). - let needed = num_frames * 2; - if scratch_stereo.len() < needed { - scratch_stereo.resize(needed, 0.0); - } - // Zero the live slice. AudioHandler::fill_buffer is - // additive (does NOT clear); residual values from - // earlier callbacks (when scratch was bigger) would - // leak through otherwise. - scratch_stereo[..needed].fill(0.0); - match handler_for_render.try_lock() { - Ok(mut h) => { - let _ = h.fill_buffer(&mut scratch_stereo[..needed]); - } - Err(std::sync::TryLockError::WouldBlock) => { + let mut scratch_stereo: Vec = vec![0.0; IOS_RENDER_SCRATCH_FRAMES * 2]; + let handler_for_render = params.handler.clone(); + let output_gain_for_render = params.output_gain.clone(); + let output_muted_for_render = params.output_muted.clone(); + let audio_processing_stats_for_render = params.audio_processing_stats.clone(); + // Level meter decimation: the render callback fires ~93 + // times/sec, but the bridge consumer reads at ~30 Hz. + let mut render_level_decimation: u32 = 0; + // Debug WAV render-reference capture is intentionally unavailable + // on iOS VPIO callbacks until WavDebugRecorder supports a + // preallocated handoff; its current push path allocates per frame. + // Diagnostic counters sampled every 100 callbacks. + let mut cb_count: u64 = 0; + let mut last_num_frames: usize = 0; + let mut num_frames_changes: u64 = 0; + let mut callbacks_with_audio: u64 = 0; + let mut callbacks_with_silence: u64 = 0; + unit.set_render_callback(move |args: render_callback::Args>| { + let render_callback::Args { + data, num_frames, .. + } = args; + let out: &mut [i16] = data.buffer; + let out_channels = data.channels; + let process_frames = num_frames.min(IOS_RENDER_SCRATCH_FRAMES); + if process_frames < num_frames { audio_processing_stats_for_render.increment_callback_xrun(); - // scratch_stereo is already zeroed above. } - Err(std::sync::TryLockError::Poisoned(e)) => { - // Never panic on the realtime IO thread. - warn!(target: "chanora_audio", "AudioHandler mutex poisoned: {e}"); - } - } - // Peak limiter — multi-client mixes can sum past 0 dBFS; - // without this the downmix helper would hard-clip to i16::MAX. - crate::voice_render::limit_peak_inplace(&mut scratch_stereo[..needed], 0.99); - let gain = f32::from_bits(output_gain_for_render.load(Ordering::Relaxed)); - let muted = output_muted_for_render.load(Ordering::Relaxed); - let mix_stats = crate::voice_render::downmix_stereo_f32_to_interleaved_i16( - &scratch_stereo[..needed], - out, - out_channels, - gain, - muted, - ); - if mix_stats.clipped_samples > 0 { - audio_processing_stats_for_render.add_clipped_samples(mix_stats.clipped_samples); - } - render_level_decimation = render_level_decimation.wrapping_add(1); - if render_level_decimation % 3 == 0 { - audio_processing_stats_for_render.update_render( - crate::frame::dbfs(&scratch_stereo[..needed]), - num_frames as u32, - ); - } - - if let Ok(guard) = wav_recorder_for_render.try_lock() { - if let Some(rec) = guard.as_ref() { - if !render_recorder_active { - render_ref_len = 0; - render_ref_accum.fill(0.0); - render_recorder_active = true; + // AudioHandler produces 48 kHz stereo f32 (= frames * 2 floats). + let needed = process_frames * 2; + // Zero the live slice. AudioHandler::fill_buffer is + // additive (does NOT clear); residual values from + // earlier callbacks (when scratch was bigger) would + // leak through otherwise. + scratch_stereo[..needed].fill(0.0); + match handler_for_render.try_lock() { + Ok(mut h) => { + let _ = h.fill_buffer(&mut scratch_stereo[..needed]); } - let mut idx = 0; - while idx + 1 < needed { - let mono = (scratch_stereo[idx] + scratch_stereo[idx + 1]) * 0.5; - render_ref_accum[render_ref_len] = mono; - render_ref_len += 1; - idx += 2; - if render_ref_len == crate::frame::FRAME_10MS_SAMPLES { - rec.push_render_reference(&render_ref_accum); - render_ref_len = 0; - } + Err(std::sync::TryLockError::WouldBlock) => { + audio_processing_stats_for_render.increment_callback_xrun(); + // scratch_stereo is already zeroed above. + } + Err(std::sync::TryLockError::Poisoned(e)) => { + // Never panic on the realtime IO thread. + warn!(target: "chanora_audio", "AudioHandler mutex poisoned: {e}"); } - } else { - render_recorder_active = false; } - } else { - render_recorder_active = false; - } - // Track audio-vs-silence for the diagnostic. - if mix_stats.peak_i16 > 0 { - callbacks_with_audio = callbacks_with_audio.wrapping_add(1); - } else { - callbacks_with_silence = callbacks_with_silence.wrapping_add(1); - // Muted output writes intentional silence (peak_i16 == 0 by - // design), not a starved render path. Gate on !muted to avoid - // counting deliberate silence as an output underrun. - if !muted { - audio_processing_stats_for_render.increment_output_underrun(); - } - } - - // Diagnostic sampling. - if last_num_frames != 0 && last_num_frames != num_frames { - num_frames_changes = num_frames_changes.wrapping_add(1); - } - last_num_frames = num_frames; - cb_count = cb_count.wrapping_add(1); - if cb_count.is_multiple_of(100) { - debug!( - target: "chanora_audio", - cb = cb_count, - num_frames, - frames_changes = num_frames_changes, - callbacks_with_audio, - callbacks_with_silence, - peak_out_i16 = mix_stats.peak_i16, + // Peak limiter — multi-client mixes can sum past 0 dBFS; + // without this the downmix helper would hard-clip to i16::MAX. + crate::voice_render::limit_peak_inplace(&mut scratch_stereo[..needed], 0.99); + let gain = f32::from_bits(output_gain_for_render.load(Ordering::Relaxed)); + let muted = output_muted_for_render.load(Ordering::Relaxed); + let mix_stats = crate::voice_render::downmix_stereo_f32_to_interleaved_i16( + &scratch_stereo[..needed], + out, + out_channels, gain, - "ios audio unit render callback diagnostic sample (direct fill_buffer)" + muted, ); - } - Ok(()) - }) - .map_err(|e| AudioError::Backend(format!("audio unit set render callback: {e}")))?; + if mix_stats.clipped_samples > 0 { + audio_processing_stats_for_render + .add_clipped_samples(mix_stats.clipped_samples); + } + render_level_decimation = render_level_decimation.wrapping_add(1); + if render_level_decimation % 3 == 0 { + audio_processing_stats_for_render.update_render( + crate::frame::dbfs(&scratch_stereo[..needed]), + num_frames as u32, + ); + } + + // Track audio-vs-silence for the diagnostic. + if mix_stats.peak_i16 > 0 { + callbacks_with_audio = callbacks_with_audio.wrapping_add(1); + } else { + callbacks_with_silence = callbacks_with_silence.wrapping_add(1); + // Muted output writes intentional silence (peak_i16 == 0 by + // design), not a starved render path. Gate on !muted to avoid + // counting deliberate silence as an output underrun. + if !muted { + audio_processing_stats_for_render.increment_output_underrun(); + } + } + + // Diagnostic sampling. + if last_num_frames != 0 && last_num_frames != num_frames { + num_frames_changes = num_frames_changes.wrapping_add(1); + } + last_num_frames = num_frames; + cb_count = cb_count.wrapping_add(1); + if cb_count.is_multiple_of(100) { + debug!( + target: "chanora_audio", + cb = cb_count, + num_frames, + frames_changes = num_frames_changes, + callbacks_with_audio, + callbacks_with_silence, + peak_out_i16 = mix_stats.peak_i16, + gain, + "ios audio unit render callback diagnostic sample (direct fill_buffer)" + ); + } + Ok(()) + }) + .map_err(|e| AudioError::Backend(format!("audio unit set render callback: {e}")))?; } // end #[cfg(target_os = "ios")] block // Finalise the unit — allocates internal buffers per the diff --git a/crates/chanora_audio/src/lib.rs b/crates/chanora_audio/src/lib.rs index 12d36be..8fe91e1 100644 --- a/crates/chanora_audio/src/lib.rs +++ b/crates/chanora_audio/src/lib.rs @@ -28,10 +28,17 @@ #![warn(missing_docs)] +#[cfg(any(target_os = "android", test))] +#[cfg_attr(not(target_os = "android"), allow(dead_code))] +mod android_render_ring; #[cfg(any(target_os = "android", test))] #[cfg_attr(not(target_os = "android"), allow(dead_code))] mod audio_event_queue; pub mod audio_processing; +#[cfg_attr(not(target_os = "android"), allow(dead_code))] +mod capture_accumulator; +#[cfg_attr(not(target_os = "android"), allow(dead_code))] +mod capture_resampler; pub mod debug_wav; mod engine; pub mod frame; @@ -42,6 +49,11 @@ pub mod processor; pub mod ptt; pub mod ptt_backends; pub mod release_tail; +#[cfg_attr( + not(any(target_os = "android", target_os = "ios", test)), + allow(dead_code) +)] +pub(crate) mod render_reference; pub mod route_policy; pub mod transmit_mode; pub mod transmit_selector; diff --git a/crates/chanora_audio/src/opus_voice.rs b/crates/chanora_audio/src/opus_voice.rs index c62fe11..d6289e1 100644 --- a/crates/chanora_audio/src/opus_voice.rs +++ b/crates/chanora_audio/src/opus_voice.rs @@ -3,15 +3,18 @@ use audiopus::{ Application as OpusApp, Bitrate as OpusBitrate, Channels as OpusChannels, SampleRate as OpusSampleRate, }; -use std::sync::atomic::{AtomicU32, Ordering}; +use crossbeam::queue::ArrayQueue; +use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; +use std::sync::Arc; use tokio::sync::mpsc; -use tracing::{info, warn}; +use tracing::{debug, info, warn}; use chanora_protocol::{AudioData, CodecType, OutAudio, OutPacket}; use crate::AudioError; pub(crate) const MAX_OPUS_FRAME: usize = 1275; +const VOICE_FRAME_QUEUE_CAPACITY: usize = 64; const VOIP_BITRATE_BPS: i32 = 32_000; const VOIP_COMPLEXITY: u8 = 10; @@ -54,10 +57,129 @@ pub(crate) fn tune_voip_encoder(encoder: &mut OpusEncoder, context: &str) { ); } +pub(crate) struct EncodedVoiceFrame { + data: [u8; MAX_OPUS_FRAME], + len: usize, +} + +pub(crate) struct EncodedVoiceFrameSender { + queue: Arc>, + open: Arc, +} + +enum EncodedVoiceFrameSendError { + Full, + Closed, +} + +impl EncodedVoiceFrameSender { + fn new(capacity: usize) -> Self { + Self { + queue: Arc::new(ArrayQueue::new(capacity)), + open: Arc::new(AtomicBool::new(true)), + } + } + + fn worker_queue(&self) -> Arc> { + Arc::clone(&self.queue) + } + + fn worker_open_flag(&self) -> Arc { + Arc::clone(&self.open) + } + + fn push(&self, frame: EncodedVoiceFrame) -> Result<(), EncodedVoiceFrameSendError> { + if !self.open.load(Ordering::Relaxed) { + return Err(EncodedVoiceFrameSendError::Closed); + } + self.queue + .push(frame) + .map_err(|_| EncodedVoiceFrameSendError::Full) + } +} + +pub(crate) fn start_out_packet_worker( + voice_out_tx: mpsc::Sender, + frames_sent: Arc, + context: &'static str, +) -> Result { + start_out_packet_worker_with_spawner(voice_out_tx, frames_sent, context, |name, worker| { + std::thread::Builder::new() + .name(name) + .spawn(worker) + .map(|_| ()) + }) +} + +fn start_out_packet_worker_with_spawner( + voice_out_tx: mpsc::Sender, + frames_sent: Arc, + context: &'static str, + spawn: S, +) -> Result +where + S: FnOnce(String, Box) -> std::io::Result<()>, +{ + let tx = EncodedVoiceFrameSender::new(VOICE_FRAME_QUEUE_CAPACITY); + let rx = tx.worker_queue(); + let worker_open = tx.worker_open_flag(); + spawn( + format!("chanora-{context}-voice-packets"), + Box::new(move || { + loop { + let Some(frame) = rx.pop() else { + if Arc::strong_count(&rx) == 1 { + break; + } + std::thread::sleep(std::time::Duration::from_millis(1)); + continue; + }; + let packet = OutAudio::new(&AudioData::C2S { + id: 0, + codec: CodecType::OpusVoice, + data: frame.as_slice(), + }); + match voice_out_tx.try_send(packet) { + Ok(()) => { + frames_sent.fetch_add(1, Ordering::Relaxed); + } + Err(mpsc::error::TrySendError::Full(_)) => { + warn!(target: "chanora_audio", context = %context, "voice_out queue full; dropping frame"); + } + Err(mpsc::error::TrySendError::Closed(_)) => { + debug!(target: "chanora_audio", context = %context, "voice_out closed; voice packet worker stopping"); + worker_open.store(false, Ordering::Relaxed); + break; + } + } + } + }), + ) + .map_err(|e| { + tx.open.store(false, Ordering::Relaxed); + AudioError::Backend(format!("voice packet worker spawn ({context}): {e}")) + })?; + Ok(tx) +} + +impl EncodedVoiceFrame { + fn try_from_opus(opus_out: &[u8], len: usize) -> Option { + if len > opus_out.len() || len > MAX_OPUS_FRAME { + return None; + } + let mut data = [0u8; MAX_OPUS_FRAME]; + data[..len].copy_from_slice(&opus_out[..len]); + Some(Self { data, len }) + } + + fn as_slice(&self) -> &[u8] { + &self.data[..self.len] + } +} + /// Encode-scope send helper for a freshly encoded Opus voice frame. pub(crate) fn send_voip_frame( - voice_out_tx: &mpsc::Sender, - frames_sent: &AtomicU32, + voice_out_tx: &EncodedVoiceFrameSender, opus_out: &[u8], len: usize, on_full: F, @@ -66,16 +188,77 @@ pub(crate) fn send_voip_frame( F: FnOnce(), G: FnOnce(), { - let packet = OutAudio::new(&AudioData::C2S { - id: 0, - codec: CodecType::OpusVoice, - data: &opus_out[..len], - }); - match voice_out_tx.try_send(packet) { - Ok(()) => { - frames_sent.fetch_add(1, Ordering::Relaxed); - } - Err(mpsc::error::TrySendError::Full(_)) => on_full(), - Err(mpsc::error::TrySendError::Closed(_)) => on_closed(), + let Some(frame) = EncodedVoiceFrame::try_from_opus(opus_out, len) else { + on_full(); + return; + }; + match voice_out_tx.push(frame) { + Ok(()) => {} + Err(EncodedVoiceFrameSendError::Full) => on_full(), + Err(EncodedVoiceFrameSendError::Closed) => on_closed(), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn encoded_voice_frame_copies_into_fixed_storage() { + let source = [7u8; MAX_OPUS_FRAME]; + + let frame = EncodedVoiceFrame::try_from_opus(&source, MAX_OPUS_FRAME).unwrap(); + + assert_eq!(frame.as_slice().len(), MAX_OPUS_FRAME); + assert!(frame.as_slice().iter().all(|byte| *byte == 7)); + } + + #[test] + fn encoded_voice_frame_rejects_lengths_beyond_fixed_storage() { + let source = [0u8; MAX_OPUS_FRAME]; + + assert!(EncodedVoiceFrame::try_from_opus(&source, MAX_OPUS_FRAME + 1).is_none()); + } + + #[test] + fn encoded_voice_frame_sender_reports_full_without_blocking() { + let sender = EncodedVoiceFrameSender::new(1); + let source = [3u8; MAX_OPUS_FRAME]; + let first = EncodedVoiceFrame::try_from_opus(&source, 4).unwrap(); + let second = EncodedVoiceFrame::try_from_opus(&source, 4).unwrap(); + + assert!(sender.push(first).is_ok()); + assert!(sender.push(second).is_err()); + } + + #[test] + fn encoded_voice_frame_sender_reports_closed_without_queueing() { + let sender = EncodedVoiceFrameSender::new(1); + sender.open.store(false, Ordering::Relaxed); + let source = [3u8; MAX_OPUS_FRAME]; + let frame = EncodedVoiceFrame::try_from_opus(&source, 4).unwrap(); + + assert!(matches!( + sender.push(frame), + Err(EncodedVoiceFrameSendError::Closed) + )); + assert_eq!(sender.queue.len(), 0); + } + + #[test] + fn encoded_voice_frame_sender_reports_spawn_failure() { + let (voice_out_tx, _voice_out_rx) = mpsc::channel(1); + let frames_sent = Arc::new(AtomicU32::new(0)); + + let result = start_out_packet_worker_with_spawner( + voice_out_tx, + frames_sent, + "test", + |_name, _worker| Err(std::io::Error::other("spawn failed")), + ); + + assert!( + matches!(result, Err(AudioError::Backend(message)) if message.contains("spawn failed")) + ); } } diff --git a/crates/chanora_audio/src/render_reference.rs b/crates/chanora_audio/src/render_reference.rs new file mode 100644 index 0000000..face3af --- /dev/null +++ b/crates/chanora_audio/src/render_reference.rs @@ -0,0 +1,241 @@ +use std::array; +use std::sync::atomic::{AtomicU32, AtomicUsize, Ordering}; +use std::sync::Arc; + +const NO_LATEST_SLOT: usize = usize::MAX; + +struct Slot { + version: AtomicUsize, + samples: [AtomicU32; SAMPLES], + #[cfg(test)] + bump_after_first_sample_read: std::sync::atomic::AtomicBool, +} + +impl Slot { + fn new() -> Self { + Self { + version: AtomicUsize::new(0), + samples: array::from_fn(|_| AtomicU32::new(0.0_f32.to_bits())), + #[cfg(test)] + bump_after_first_sample_read: std::sync::atomic::AtomicBool::new(false), + } + } +} + +pub(crate) struct RenderReferenceFrameAccumulator { + pending: [f32; SAMPLES], + pending_len: usize, +} + +impl RenderReferenceFrameAccumulator { + pub(crate) fn new() -> Self { + assert!( + SAMPLES > 0, + "RenderReferenceFrameAccumulator requires at least one sample" + ); + Self { + pending: [0.0; SAMPLES], + pending_len: 0, + } + } + + pub(crate) fn push_mono_samples( + &mut self, + mut samples: &[f32], + mut publish: impl FnMut(&[f32; SAMPLES]), + ) { + while !samples.is_empty() { + let needed = SAMPLES - self.pending_len; + let take = needed.min(samples.len()); + self.pending[self.pending_len..self.pending_len + take] + .copy_from_slice(&samples[..take]); + self.pending_len += take; + samples = &samples[take..]; + + if self.pending_len == SAMPLES { + publish(&self.pending); + self.pending_len = 0; + } + } + } + + #[cfg(test)] + fn pending_len(&self) -> usize { + self.pending_len + } +} + +pub(crate) struct RenderReferenceBuffer { + slots: Box<[Slot; SLOTS]>, + write_idx: AtomicUsize, + latest_slot: AtomicUsize, +} + +impl RenderReferenceBuffer { + pub(crate) fn new() -> Arc { + assert!( + SLOTS > 0, + "RenderReferenceBuffer requires at least one slot" + ); + Arc::new(Self { + slots: Box::new(array::from_fn(|_| Slot::new())), + write_idx: AtomicUsize::new(0), + latest_slot: AtomicUsize::new(NO_LATEST_SLOT), + }) + } + + pub(crate) fn write(&self, frame: &[f32; SAMPLES]) { + let idx = self.write_idx.load(Ordering::Relaxed) % SLOTS; + let slot = &self.slots[idx]; + // The acquire half keeps payload stores after the odd in-progress marker. + let version = slot.version.fetch_add(1, Ordering::AcqRel); + debug_assert_eq!(version & 1, 0, "single writer should only enter even slots"); + for (sample, value) in slot.samples.iter().zip(frame.iter().copied()) { + sample.store(value.to_bits(), Ordering::Relaxed); + } + slot.version + .store(version.wrapping_add(2) & !1, Ordering::Release); + self.latest_slot.store(idx, Ordering::Release); + self.write_idx.store((idx + 1) % SLOTS, Ordering::Relaxed); + } + + pub(crate) fn read_latest(&self) -> [f32; SAMPLES] { + let mut out = [0.0_f32; SAMPLES]; + self.read_latest_into(&mut out); + out + } + + pub(crate) fn read_latest_into(&self, out: &mut [f32; SAMPLES]) { + let idx = self.latest_slot.load(Ordering::Acquire); + if idx == NO_LATEST_SLOT { + out.fill(0.0); + return; + } + + let slot = &self.slots[idx]; + let before = slot.version.load(Ordering::Acquire); + if before & 1 == 1 { + out.fill(0.0); + return; + } + + #[cfg(not(test))] + for (dst, sample) in out.iter_mut().zip(slot.samples.iter()) { + *dst = f32::from_bits(sample.load(Ordering::Relaxed)); + } + + #[cfg(test)] + for (idx, (dst, sample)) in out.iter_mut().zip(slot.samples.iter()).enumerate() { + *dst = f32::from_bits(sample.load(Ordering::Relaxed)); + if idx == 0 + && slot + .bump_after_first_sample_read + .swap(false, Ordering::Relaxed) + { + slot.version.fetch_add(2, Ordering::Release); + } + } + + let after = slot.version.load(Ordering::Acquire); + if before != after || after & 1 == 1 { + out.fill(0.0); + } + } + + #[cfg(test)] + fn mark_latest_slot_in_progress_for_test(&self) { + let idx = self.latest_slot.load(Ordering::Acquire); + assert_ne!(idx, NO_LATEST_SLOT); + self.slots[idx].version.fetch_or(1, Ordering::Release); + } + + #[cfg(test)] + fn bump_latest_slot_version_after_first_sample_for_test(&self) { + let idx = self.latest_slot.load(Ordering::Acquire); + assert_ne!(idx, NO_LATEST_SLOT); + self.slots[idx] + .bump_after_first_sample_read + .store(true, Ordering::Relaxed); + } +} + +#[cfg(test)] +mod tests { + use super::{RenderReferenceBuffer, RenderReferenceFrameAccumulator}; + + #[test] + fn render_reference_reads_zero_before_first_publish() { + let buffer = RenderReferenceBuffer::<4, 2>::new(); + + assert_eq!(buffer.read_latest(), [0.0; 4]); + } + + #[test] + fn render_reference_reader_gets_latest_complete_frame() { + let buffer = RenderReferenceBuffer::<4, 3>::new(); + + buffer.write(&[1.0, 2.0, 3.0, 4.0]); + buffer.write(&[5.0, 6.0, 7.0, 8.0]); + + assert_eq!(buffer.read_latest(), [5.0, 6.0, 7.0, 8.0]); + } + + #[test] + fn render_reference_writes_wrap_without_returning_stale_frame() { + let buffer = RenderReferenceBuffer::<2, 2>::new(); + + buffer.write(&[1.0, 2.0]); + buffer.write(&[3.0, 4.0]); + buffer.write(&[5.0, 6.0]); + + assert_eq!(buffer.read_latest(), [5.0, 6.0]); + } + + #[test] + fn render_reference_accumulator_publishes_only_complete_frames() { + let mut accum = RenderReferenceFrameAccumulator::<4>::new(); + let mut frames = Vec::new(); + + accum.push_mono_samples(&[1.0, 2.0], |frame| frames.push(*frame)); + + assert!(frames.is_empty()); + assert_eq!(accum.pending_len(), 2); + + accum.push_mono_samples(&[3.0, 4.0, 5.0, 6.0, 7.0], |frame| frames.push(*frame)); + + assert_eq!(frames, vec![[1.0, 2.0, 3.0, 4.0]]); + assert_eq!(accum.pending_len(), 3); + + accum.push_mono_samples(&[8.0], |frame| frames.push(*frame)); + + assert_eq!(frames, vec![[1.0, 2.0, 3.0, 4.0], [5.0, 6.0, 7.0, 8.0]]); + assert_eq!(accum.pending_len(), 0); + } + + #[test] + fn render_reference_reader_rejects_in_progress_slot() { + let buffer = RenderReferenceBuffer::<2, 1>::new(); + + buffer.write(&[1.0, 2.0]); + buffer.mark_latest_slot_in_progress_for_test(); + + assert_eq!(buffer.read_latest(), [0.0, 0.0]); + } + + #[test] + fn render_reference_reader_rejects_stale_slot_changed_during_read() { + let buffer = RenderReferenceBuffer::<2, 1>::new(); + + buffer.write(&[1.0, 2.0]); + + buffer.bump_latest_slot_version_after_first_sample_for_test(); + + assert_eq!(buffer.read_latest(), [0.0, 0.0]); + } + + #[test] + #[should_panic(expected = "RenderReferenceBuffer requires at least one slot")] + fn render_reference_rejects_zero_slots() { + let _ = RenderReferenceBuffer::<2, 0>::new(); + } +} diff --git a/crates/chanora_audio/src/vad/mod.rs b/crates/chanora_audio/src/vad/mod.rs index d8b6000..060cff3 100644 --- a/crates/chanora_audio/src/vad/mod.rs +++ b/crates/chanora_audio/src/vad/mod.rs @@ -15,7 +15,7 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{OnceLock, RwLock}; use crate::frame::{f32_to_i16, i16_to_f32}; -use crate::AudioError; +use crate::{AudioError, VadBackend}; use resampler::{Downsampler48to16, INPUT_FRAME_10MS}; #[cfg(not(target_os = "ios"))] @@ -108,6 +108,36 @@ pub fn process_i16_10ms(detector: &mut dyn VoiceActivityDetector, samples: &[i16 detector.process_10ms(&frame) } +/// Callback-side policy for optional model-backed VAD workers. +#[derive(Debug, Clone, Copy, Eq, PartialEq)] +pub(crate) enum VadWorkerPolicy { + /// Keep using the already-available model worker. + UseWorker, + /// No worker may be constructed on the callback thread; use WebRTC fallback. + UseFallback, + /// This backend does not need a model worker. + NotModelBacked, +} + +/// Decide whether a realtime callback may use a model-backed VAD worker. +/// +/// Model/worker construction is intentionally absent from this policy: if a +/// worker is not already present, callbacks must stay nonblocking and fall back. +pub(crate) fn callback_vad_worker_policy( + voice_activity_mode: bool, + backend: VadBackend, + worker_available: bool, +) -> VadWorkerPolicy { + if !voice_activity_mode || backend != VadBackend::SileroOnnx { + return VadWorkerPolicy::NotModelBacked; + } + if worker_available { + VadWorkerPolicy::UseWorker + } else { + VadWorkerPolicy::UseFallback + } +} + static SILERO_MODEL_PATH_OVERRIDE: OnceLock>> = OnceLock::new(); static SILERO_MODEL_EPOCH: AtomicU64 = AtomicU64::new(0); @@ -252,4 +282,32 @@ mod tests { assert_eq!(silero_model_bundle_path(), path.to_string_lossy()); let _ = std::fs::remove_file(path); } + + #[test] + fn callback_policy_uses_existing_model_worker_only() { + assert_eq!( + callback_vad_worker_policy(true, VadBackend::SileroOnnx, true), + VadWorkerPolicy::UseWorker + ); + assert_eq!( + callback_vad_worker_policy(true, VadBackend::SileroOnnx, false), + VadWorkerPolicy::UseFallback + ); + } + + #[test] + fn callback_policy_keeps_disabled_and_webrtc_paths_worker_free() { + assert_eq!( + callback_vad_worker_policy(false, VadBackend::SileroOnnx, false), + VadWorkerPolicy::NotModelBacked + ); + assert_eq!( + callback_vad_worker_policy(true, VadBackend::Disabled, false), + VadWorkerPolicy::NotModelBacked + ); + assert_eq!( + callback_vad_worker_policy(true, VadBackend::WebrtcVad, false), + VadWorkerPolicy::NotModelBacked + ); + } }