fix(audio): eliminate Android output stutter via Oboe config + lock-free callback

Phase 1 — Oboe configuration:
- Change output stream from Usage::VoiceCommunication to Usage::Game with
  ContentType::Sonification to avoid forcing the Legacy (OpenSL ES) data
  path on most devices (Oboe issue #2075)
- Switch output format from i16 Mono to f32 Stereo, matching Qint's proven
  configuration and eliminating per-callback downmix conversion
- Set buffer size to 2x burst after stream open, reducing default buffer
  from 8-20x burst to 2x burst for lower latency
- Remove scratch Mutex<Vec<f32>>; callback writes directly to Oboe buffer

Phase 2 — Lock-free output callback:
- Add audio_event_queue.rs: lock-free SPSC bridge using crossbeam ArrayQueue
  with separate packet (lossy) and control (reliable) channels
- OutputCallback now owns AudioHandler directly (no Arc<Mutex<>> on Android)
- Inbound forwarder pushes packets via AudioEventProducer (no mutex)
- set_client_volume pushes control commands via event queue on Android
- iOS/desktop Arc<Mutex<AudioHandler>> path unchanged
This commit is contained in:
Edison Jwa
2026-06-04 00:51:25 +09:00
parent d3bb199208
commit 484dad1072
7 changed files with 291 additions and 87 deletions
Generated
+22
View File
@@ -423,6 +423,8 @@ dependencies = [
"coreaudio-rs", "coreaudio-rs",
"cpal", "cpal",
"criterion", "criterion",
"crossbeam",
"crossbeam-utils",
"dhat", "dhat",
"dispatch2", "dispatch2",
"futures-util", "futures-util",
@@ -803,6 +805,17 @@ version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b"
[[package]]
name = "crossbeam"
version = "0.8.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1137cd7e7fc0fb5d3c5a8678be38ec56e819125d8d7907411fe24ccb943faca8"
dependencies = [
"crossbeam-epoch",
"crossbeam-queue",
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-channel" name = "crossbeam-channel"
version = "0.5.15" version = "0.5.15"
@@ -831,6 +844,15 @@ dependencies = [
"crossbeam-utils", "crossbeam-utils",
] ]
[[package]]
name = "crossbeam-queue"
version = "0.3.12"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0f58bbc28f91df819d0aa2a2c00cd19754769c2fad90579b3592b1c9ba7a3115"
dependencies = [
"crossbeam-utils",
]
[[package]] [[package]]
name = "crossbeam-utils" name = "crossbeam-utils"
version = "0.8.21" version = "0.8.21"
+2
View File
@@ -29,6 +29,8 @@ audiopus = "0.3.0-rc.0"
tsclientlib = { git = "https://github.com/ReSpeak/tsclientlib.git", rev = "04aa2491", default-features = false, features = ["audio"] } tsclientlib = { git = "https://github.com/ReSpeak/tsclientlib.git", rev = "04aa2491", default-features = false, features = ["audio"] }
tokio = { version = "1", features = ["sync", "rt", "macros", "time"] } tokio = { version = "1", features = ["sync", "rt", "macros", "time"] }
rustfft = "6.2.0" rustfft = "6.2.0"
crossbeam = { version = "0.8", default-features = false, features = ["alloc", "crossbeam-queue"] }
crossbeam-utils = { version = "0.8", default-features = false }
[target.'cfg(all(not(target_os = "android"), not(target_os = "ios"), not(target_os = "macos")))'.dependencies] [target.'cfg(all(not(target_os = "android"), not(target_os = "ios"), not(target_os = "macos")))'.dependencies]
# Desktop audio I/O for Windows capture/playback and Linux capture. # Desktop audio I/O for Windows capture/playback and Linux capture.
+84 -53
View File
@@ -44,6 +44,7 @@ use std::sync::{Arc, Mutex};
use audiopus::coder::Encoder as OpusEncoder; use audiopus::coder::Encoder as OpusEncoder;
use tracing::{debug, info, warn}; use tracing::{debug, info, warn};
use crate::audio_event_queue::{AudioCommand, AudioEventQueue};
use crate::mobile_voice_backend::{ use crate::mobile_voice_backend::{
clear_android_audio_diagnostics, latency_tier_for, next_input_preset_after, clear_android_audio_diagnostics, latency_tier_for, next_input_preset_after,
next_sharing_mode_after, publish_android_audio_diagnostics, AchievedInputPreset, next_sharing_mode_after, publish_android_audio_diagnostics, AchievedInputPreset,
@@ -63,7 +64,7 @@ use oboe::{
AudioInputCallback, AudioInputStreamSafe, AudioOutputCallback, AudioOutputStreamSafe, AudioInputCallback, AudioInputStreamSafe, AudioOutputCallback, AudioOutputStreamSafe,
AudioStream, AudioStreamAsync, AudioStreamBase, AudioStreamBuilder, AudioStreamSafe, AudioStream, AudioStreamAsync, AudioStreamBase, AudioStreamBuilder, AudioStreamSafe,
DataCallbackResult, Input as OboeInput, InputPreset, Mono, Output as OboeOutput, DataCallbackResult, Input as OboeInput, InputPreset, Mono, Output as OboeOutput,
PerformanceMode, SessionId, SharingMode, Usage, PerformanceMode, SessionId, SharingMode, Stereo, Usage,
}; };
use crate::processor::AudioProcessor; use crate::processor::AudioProcessor;
@@ -518,14 +519,14 @@ impl AudioInputCallback for InputCallback {
// //
// Mirrors the iOS VPIO render callback. Pulls mixed 48 kHz stereo f32 // Mirrors the iOS VPIO render callback. Pulls mixed 48 kHz stereo f32
// from `AudioHandler::fill_buffer`, applies output gain + mute, and // from `AudioHandler::fill_buffer`, applies output gain + mute, and
// writes mono i16 to the Oboe output buffer. // writes stereo f32 directly to the Oboe output buffer.
struct OutputCallback { struct OutputCallback {
handler: Arc<Mutex<AudioHandler<SessionAudioId>>>, handler: AudioHandler<SessionAudioId>,
event_queue: Arc<AudioEventQueue>,
output_gain: Arc<AtomicU32>, output_gain: Arc<AtomicU32>,
output_muted: Arc<AtomicBool>, output_muted: Arc<AtomicBool>,
event_tx: BackendEventTx, event_tx: BackendEventTx,
scratch: Arc<Mutex<Vec<f32>>>,
render_reference: Arc<RenderReferenceBuffer>, render_reference: Arc<RenderReferenceBuffer>,
audio_processing_stats: Arc<crate::SharedAudioProcessingStats>, audio_processing_stats: Arc<crate::SharedAudioProcessingStats>,
pending_render_ref: [f32; crate::frame::FRAME_10MS_SAMPLES], pending_render_ref: [f32; crate::frame::FRAME_10MS_SAMPLES],
@@ -533,51 +534,56 @@ struct OutputCallback {
} }
impl AudioOutputCallback for OutputCallback { impl AudioOutputCallback for OutputCallback {
type FrameType = (i16, Mono); type FrameType = (f32, Stereo);
fn on_audio_ready( fn on_audio_ready(
&mut self, &mut self,
_stream: &mut dyn AudioOutputStreamSafe, _stream: &mut dyn AudioOutputStreamSafe,
frames: &mut [i16], frames: &mut [(f32, f32)],
) -> DataCallbackResult { ) -> DataCallbackResult {
let _ = catch_unwind(AssertUnwindSafe(|| { let _ = catch_unwind(AssertUnwindSafe(|| {
let needed = frames.len() * 2; // stereo let buf: &mut [f32] = unsafe {
let scratch = &mut self.scratch.lock().unwrap(); std::slice::from_raw_parts_mut(frames.as_mut_ptr() as *mut f32, frames.len() * 2)
if scratch.len() < needed { };
scratch.resize(needed, 0.0); for s in buf.iter_mut() {
} else { *s = 0.0;
for s in &mut scratch[..needed] { }
*s = 0.0; let consumer = AudioEventQueue::consumer(&self.event_queue);
for cmd in 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);
}
} }
} }
match self.handler.try_lock() {
Ok(mut h) => { for pkt in consumer.drain_packets(50) {
let _ = h.fill_buffer(&mut scratch[..needed]); if let Err(e) = self.handler.handle_packet(pkt.client_id, pkt.data) {
} debug!(target: "chanora_audio", error = %e, "decode failed");
Err(std::sync::TryLockError::WouldBlock) => {}
Err(std::sync::TryLockError::Poisoned(e)) => {
warn!(
target: "chanora_audio",
"AudioHandler mutex poisoned: {}",
e
);
} }
} }
let _ = self.handler.fill_buffer(buf);
let gain = f32::from_bits(self.output_gain.load(Ordering::Relaxed)); let gain = f32::from_bits(self.output_gain.load(Ordering::Relaxed));
let muted = self.output_muted.load(Ordering::Relaxed); let muted = self.output_muted.load(Ordering::Relaxed);
let _ = crate::voice_render::downmix_stereo_f32_to_mono_i16( if muted {
&scratch[..needed], for s in buf.iter_mut() {
frames, *s = 0.0;
gain, }
muted, } else if gain != 1.0 {
); for s in buf.iter_mut() {
*s *= gain;
}
}
self.audio_processing_stats self.audio_processing_stats
.update_render(crate::frame::dbfs(&scratch[..needed]), frames.len() as u32); .update_render(crate::frame::dbfs(buf), frames.len() as u32);
// Accumulate the full render callback into 10 ms mono chunks so for chunk in buf.chunks_exact(2) {
// AEC sees consistent reference timing even when output callbacks
// are shorter or longer than 10 ms.
for chunk in scratch[..needed].chunks_exact(2) {
self.pending_render_ref[self.pending_render_ref_len] = (chunk[0] + chunk[1]) * 0.5; self.pending_render_ref[self.pending_render_ref_len] = (chunk[0] + chunk[1]) * 0.5;
self.pending_render_ref_len += 1; self.pending_render_ref_len += 1;
if self.pending_render_ref_len == crate::frame::FRAME_10MS_SAMPLES { if self.pending_render_ref_len == crate::frame::FRAME_10MS_SAMPLES {
@@ -677,8 +683,6 @@ impl AndroidVoiceUnit {
.map_err(|e| BackendError::OpenFailed(format!("capture state init: {e}")))?, .map_err(|e| BackendError::OpenFailed(format!("capture state init: {e}")))?,
)); ));
let scratch = Arc::new(Mutex::new(Vec::with_capacity(8192)));
// --- Open input stream (SDD-112) --------------------------- // --- Open input stream (SDD-112) ---------------------------
let input_builder = AudioStreamBuilder::default() let input_builder = AudioStreamBuilder::default()
.set_direction::<OboeInput>() .set_direction::<OboeInput>()
@@ -766,8 +770,8 @@ impl AndroidVoiceUnit {
let output_builder = AudioStreamBuilder::default() let output_builder = AudioStreamBuilder::default()
.set_direction::<OboeOutput>() .set_direction::<OboeOutput>()
.set_sample_rate(cfg.sample_rate as i32) .set_sample_rate(cfg.sample_rate as i32)
.set_channel_count::<Mono>() .set_channel_count::<Stereo>()
.set_format::<i16>() .set_format::<f32>()
.set_performance_mode(if cfg.request_low_latency { .set_performance_mode(if cfg.request_low_latency {
PerformanceMode::LowLatency PerformanceMode::LowLatency
} else { } else {
@@ -778,16 +782,21 @@ impl AndroidVoiceUnit {
} else { } else {
SharingMode::Shared SharingMode::Shared
}) })
.set_usage(Usage::VoiceCommunication) // Usage::Game avoids forcing the Legacy (OpenSL ES) data path
.set_content_type(oboe::ContentType::Speech); // that Usage::VoiceCommunication triggers on most devices.
// Android audio routing is already handled by
// AudioManager.MODE_IN_COMMUNICATION on the Flutter side.
.set_usage(Usage::Game)
.set_content_type(oboe::ContentType::Sonification);
let render_ref_for_output = render_ref_buf.clone(); let render_ref_for_output = render_ref_buf.clone();
let event_queue = params.event_producer.queue();
let output_cb = OutputCallback { let output_cb = OutputCallback {
handler: params.handler.clone(), handler: params.handler,
event_queue: event_queue.clone(),
output_gain: params.output_gain.clone(), output_gain: params.output_gain.clone(),
output_muted: params.output_muted.clone(), output_muted: params.output_muted.clone(),
event_tx: event_tx.clone(), event_tx: event_tx.clone(),
scratch: scratch.clone(),
render_reference: render_ref_for_output, render_reference: render_ref_for_output,
audio_processing_stats: audio_processing_stats.clone(), audio_processing_stats: audio_processing_stats.clone(),
pending_render_ref: [0.0_f32; crate::frame::FRAME_10MS_SAMPLES], pending_render_ref: [0.0_f32; crate::frame::FRAME_10MS_SAMPLES],
@@ -806,20 +815,41 @@ impl AndroidVoiceUnit {
Self::open_output_fallback( Self::open_output_fallback(
cfg, cfg,
&event_tx, &event_tx,
params.handler.clone(), AudioHandler::new(),
event_queue.clone(),
params.output_gain.clone(), params.output_gain.clone(),
params.output_muted.clone(), params.output_muted.clone(),
audio_processing_stats.clone(), audio_processing_stats.clone(),
scratch.clone(),
render_ref_buf, render_ref_buf,
)? )?
} }
}; };
let output_frames_per_burst = output_stream.get_frames_per_burst();
if output_frames_per_burst > 0 {
let desired = output_frames_per_burst * 2;
match output_stream.set_buffer_size_in_frames(desired) {
Ok(actual) => {
debug!(
target: "chanora_audio",
desired,
actual,
"android: output buffer size tuned"
);
}
Err(e) => {
warn!(
target: "chanora_audio",
error = ?e,
"android: output buffer size tuning failed; using device default"
);
}
}
}
let output_perf = perf_from_oboe(output_stream.get_performance_mode()); let output_perf = perf_from_oboe(output_stream.get_performance_mode());
let output_share = share_from_oboe(output_stream.get_sharing_mode()); let output_share = share_from_oboe(output_stream.get_sharing_mode());
let output_sample_rate = output_stream.get_sample_rate(); let output_sample_rate = output_stream.get_sample_rate();
let output_frames_per_burst = output_stream.get_frames_per_burst();
// SDD-112 / SRS-210: structured "stream opened" event with // SDD-112 / SRS-210: structured "stream opened" event with
// achieved values. No PII; only platform-reported scalars. // achieved values. No PII; only platform-reported scalars.
@@ -1038,19 +1068,19 @@ impl AndroidVoiceUnit {
fn open_output_fallback( fn open_output_fallback(
cfg: &AndroidVoiceStreamConfig, cfg: &AndroidVoiceStreamConfig,
event_tx: &BackendEventTx, event_tx: &BackendEventTx,
handler: Arc<Mutex<AudioHandler<SessionAudioId>>>, handler: AudioHandler<SessionAudioId>,
event_queue: Arc<AudioEventQueue>,
output_gain: Arc<AtomicU32>, output_gain: Arc<AtomicU32>,
output_muted: Arc<AtomicBool>, output_muted: Arc<AtomicBool>,
audio_processing_stats: Arc<crate::SharedAudioProcessingStats>, audio_processing_stats: Arc<crate::SharedAudioProcessingStats>,
scratch: Arc<Mutex<Vec<f32>>>,
render_reference: Arc<RenderReferenceBuffer>, render_reference: Arc<RenderReferenceBuffer>,
) -> Result<AudioStreamAsync<OboeOutput, OutputCallback>, BackendError> { ) -> Result<AudioStreamAsync<OboeOutput, OutputCallback>, BackendError> {
let cb = OutputCallback { let cb = OutputCallback {
handler, handler,
event_queue,
output_gain, output_gain,
output_muted, output_muted,
event_tx: event_tx.clone(), event_tx: event_tx.clone(),
scratch,
render_reference, render_reference,
audio_processing_stats, audio_processing_stats,
pending_render_ref: [0.0_f32; crate::frame::FRAME_10MS_SAMPLES], pending_render_ref: [0.0_f32; crate::frame::FRAME_10MS_SAMPLES],
@@ -1059,12 +1089,13 @@ impl AndroidVoiceUnit {
let builder = AudioStreamBuilder::default() let builder = AudioStreamBuilder::default()
.set_direction::<OboeOutput>() .set_direction::<OboeOutput>()
.set_sample_rate(cfg.sample_rate as i32) .set_sample_rate(cfg.sample_rate as i32)
.set_channel_count::<Mono>() .set_channel_count::<Stereo>()
.set_format::<i16>() .set_format::<f32>()
.set_performance_mode(PerformanceMode::LowLatency) .set_performance_mode(PerformanceMode::LowLatency)
.set_sharing_mode(SharingMode::Shared) .set_sharing_mode(SharingMode::Shared)
.set_usage(Usage::VoiceCommunication) // Same Usage::Game rationale as primary output builder above.
.set_content_type(oboe::ContentType::Speech) .set_usage(Usage::Game)
.set_content_type(oboe::ContentType::Sonification)
.set_callback(cb); .set_callback(cb);
builder builder
.open_stream() .open_stream()
@@ -0,0 +1,118 @@
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use chanora_protocol::InAudioBuf;
use crossbeam::queue::ArrayQueue;
use crate::engine::SessionAudioId;
const PACKET_QUEUE_CAPACITY: usize = 100;
const CONTROL_QUEUE_CAPACITY: usize = 32;
/// A raw inbound voice packet waiting to be inserted into AudioHandler.
pub struct AudioPacket {
/// Client whose TeamSpeak audio packet this belongs to.
pub client_id: SessionAudioId,
/// Raw inbound TeamSpeak audio payload accepted by AudioHandler::handle_packet.
pub data: InAudioBuf,
}
/// Control commands from the main thread to the audio callback.
pub enum AudioCommand {
/// Set a client's output volume.
SetVolume(SessionAudioId, f32),
/// Remove a client's decode queue.
RemoveClient(SessionAudioId),
}
/// Lock-free bridge between the inbound forwarder / main thread and the
/// audio callback. The callback owns the consumer halves.
pub struct AudioEventQueue {
/// Bounded lossy queue for raw voice packets. On overflow, the push
/// fails and the packet is dropped (counted via `packets_dropped`).
/// Capacity: 100 packets (~2 seconds at 50pps, far more than needed).
pub packet_queue: ArrayQueue<AudioPacket>,
/// Bounded reliable queue for control commands (volume, client removal).
/// On overflow, the caller retries. Capacity: 32 commands.
pub control_queue: ArrayQueue<AudioCommand>,
/// Atomic counter for dropped packets (for diagnostics).
pub packets_dropped: AtomicU64,
}
impl AudioEventQueue {
/// Create the Android audio event bridge with fixed queue capacities.
pub fn new() -> Arc<Self> {
Arc::new(Self {
packet_queue: ArrayQueue::new(PACKET_QUEUE_CAPACITY),
control_queue: ArrayQueue::new(CONTROL_QUEUE_CAPACITY),
packets_dropped: AtomicU64::new(0),
})
}
/// Create a producer handle sharing this queue.
pub fn producer(queue: &Arc<Self>) -> AudioEventProducer {
AudioEventProducer {
queue: Arc::clone(queue),
}
}
/// Create a consumer handle sharing this queue.
pub fn consumer(queue: &Arc<Self>) -> AudioEventConsumer {
AudioEventConsumer {
queue: Arc::clone(queue),
}
}
}
/// Producer side used by the inbound forwarder and engine control methods.
#[derive(Clone)]
pub struct AudioEventProducer {
queue: Arc<AudioEventQueue>,
}
impl AudioEventProducer {
/// Push a raw voice packet, incrementing the drop counter if full.
pub fn push_packet(&self, packet: AudioPacket) -> Result<(), AudioPacket> {
self.queue.packet_queue.push(packet).map_err(|packet| {
self.queue.packets_dropped.fetch_add(1, Ordering::Relaxed);
packet
})
}
/// Push a control command, returning it unchanged if the queue is full.
pub fn push_control(&self, cmd: AudioCommand) -> Result<(), AudioCommand> {
self.queue.control_queue.push(cmd)
}
/// Shared queue backing this producer.
pub fn queue(&self) -> Arc<AudioEventQueue> {
Arc::clone(&self.queue)
}
}
/// Consumer side used by the Android output callback.
pub struct AudioEventConsumer {
queue: Arc<AudioEventQueue>,
}
impl AudioEventConsumer {
/// Pop up to `cap` queued packets.
pub fn drain_packets(&self, cap: usize) -> impl Iterator<Item = AudioPacket> + '_ {
let mut drained = 0;
std::iter::from_fn(move || {
if drained >= cap {
return None;
}
let packet = self.queue.packet_queue.pop();
if packet.is_some() {
drained += 1;
}
packet
})
}
/// Pop all currently queued controls.
pub fn drain_controls(&self) -> impl Iterator<Item = AudioCommand> + '_ {
std::iter::from_fn(move || self.queue.control_queue.pop())
}
}
+52 -33
View File
@@ -52,6 +52,8 @@ use tsclientlib::audio::AudioHandler;
use chanora_protocol::{InboundVoice, OutPacket}; use chanora_protocol::{InboundVoice, OutPacket};
#[cfg(target_os = "android")]
use crate::audio_event_queue::{AudioCommand, AudioEventQueue, AudioPacket};
use crate::AudioError; use crate::AudioError;
#[cfg(all( #[cfg(all(
@@ -305,12 +307,15 @@ pub struct AudioEngine {
output_muted: Arc<AtomicBool>, output_muted: Arc<AtomicBool>,
audio_processing_config: Arc<Mutex<crate::AudioProcessingConfig>>, audio_processing_config: Arc<Mutex<crate::AudioProcessingConfig>>,
audio_processing_stats: Arc<crate::SharedAudioProcessingStats>, audio_processing_stats: Arc<crate::SharedAudioProcessingStats>,
#[cfg(not(target_os = "android"))]
audio_handler: Arc<Mutex<AudioHandler<SessionAudioId>>>, audio_handler: Arc<Mutex<AudioHandler<SessionAudioId>>>,
#[cfg(any(target_os = "ios", target_os = "macos", target_os = "android"))] #[cfg(target_os = "android")]
audio_event_producer: crate::audio_event_queue::AudioEventProducer,
#[cfg(target_os = "android")]
voice_out_tx: mpsc::Sender<OutPacket>, voice_out_tx: mpsc::Sender<OutPacket>,
#[cfg(any(target_os = "ios", target_os = "macos", target_os = "android"))] #[cfg(target_os = "android")]
voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>, voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>,
#[cfg(any(target_os = "ios", target_os = "macos", target_os = "android"))] #[cfg(target_os = "android")]
mic_gain: f32, mic_gain: f32,
// Streams must be dropped to stop audio. Both are `!Send` because // Streams must be dropped to stop audio. Both are `!Send` because
// cpal's Stream isn't Send on some backends; we keep them in an // cpal's Stream isn't Send on some backends; we keep them in an
@@ -482,7 +487,7 @@ impl AudioEngine {
voice_out_tx: mpsc::Sender<OutPacket>, voice_out_tx: mpsc::Sender<OutPacket>,
transmit_gate: crate::ptt::AudioTransmitGate, transmit_gate: crate::ptt::AudioTransmitGate,
frames_sent: Arc<AtomicU32>, frames_sent: Arc<AtomicU32>,
audio_handler: Arc<Mutex<AudioHandler<SessionAudioId>>>, event_producer: crate::audio_event_queue::AudioEventProducer,
output_gain: Arc<AtomicU32>, output_gain: Arc<AtomicU32>,
output_muted: Arc<AtomicBool>, output_muted: Arc<AtomicBool>,
voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>, voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>,
@@ -552,7 +557,8 @@ impl AudioEngine {
transmit_active: transmit_gate.flag_arc(), transmit_active: transmit_gate.flag_arc(),
frames_sent: frames_sent.clone(), frames_sent: frames_sent.clone(),
mic_gain, mic_gain,
handler: audio_handler.clone(), handler: AudioHandler::new(),
event_producer: event_producer.clone(),
output_gain: output_gain.clone(), output_gain: output_gain.clone(),
output_muted: output_muted.clone(), output_muted: output_muted.clone(),
voice_activity_selector: voice_activity_selector.clone(), voice_activity_selector: voice_activity_selector.clone(),
@@ -572,7 +578,7 @@ impl AudioEngine {
voice_out_tx.clone(), voice_out_tx.clone(),
transmit_gate.clone(), transmit_gate.clone(),
frames_sent.clone(), frames_sent.clone(),
audio_handler.clone(), event_producer.clone(),
output_gain.clone(), output_gain.clone(),
output_muted.clone(), output_muted.clone(),
voice_activity_selector.clone(), voice_activity_selector.clone(),
@@ -1007,8 +1013,8 @@ impl AudioEngine {
let audio_processing_config = Arc::new(Mutex::new(crate::AudioProcessingConfig::default())); let audio_processing_config = Arc::new(Mutex::new(crate::AudioProcessingConfig::default()));
let audio_processing_stats = Arc::new(crate::SharedAudioProcessingStats::default()); let audio_processing_stats = Arc::new(crate::SharedAudioProcessingStats::default());
let audio_handler: Arc<Mutex<AudioHandler<SessionAudioId>>> = let event_queue = AudioEventQueue::new();
Arc::new(Mutex::new(AudioHandler::new())); let event_producer = AudioEventQueue::producer(&event_queue);
if !cfg.mobile_voice_preset { if !cfg.mobile_voice_preset {
return Err(AudioError::Backend( return Err(AudioError::Backend(
@@ -1069,7 +1075,8 @@ impl AudioEngine {
transmit_active: transmit_flag_for_capture, transmit_active: transmit_flag_for_capture,
frames_sent: frames_sent.clone(), frames_sent: frames_sent.clone(),
mic_gain: cfg.mic_gain, mic_gain: cfg.mic_gain,
handler: audio_handler.clone(), handler: AudioHandler::new(),
event_producer: event_producer.clone(),
output_gain: output_gain.clone(), output_gain: output_gain.clone(),
output_muted: output_muted.clone(), output_muted: output_muted.clone(),
voice_activity_selector: cfg.voice_activity_selector.clone(), voice_activity_selector: cfg.voice_activity_selector.clone(),
@@ -1103,7 +1110,7 @@ impl AudioEngine {
voice_out_tx.clone(), voice_out_tx.clone(),
transmit_gate.clone(), transmit_gate.clone(),
frames_sent.clone(), frames_sent.clone(),
audio_handler.clone(), event_producer.clone(),
output_gain.clone(), output_gain.clone(),
output_muted.clone(), output_muted.clone(),
cfg.voice_activity_selector.clone(), cfg.voice_activity_selector.clone(),
@@ -1151,7 +1158,7 @@ impl AudioEngine {
let capture_active = true; let capture_active = true;
let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel(); let (shutdown_tx, mut shutdown_rx) = tokio::sync::oneshot::channel();
let handler_for_task = audio_handler.clone(); let event_producer_for_task = event_producer.clone();
let frames_received_for_task = frames_received.clone(); let frames_received_for_task = frames_received.clone();
tokio::spawn(async move { tokio::spawn(async move {
loop { loop {
@@ -1164,10 +1171,8 @@ impl AudioEngine {
match item { match item {
Some(v) => { Some(v) => {
let id = SessionAudioId(v.from_client); let id = SessionAudioId(v.from_client);
let mut h = handler_for_task.lock().unwrap(); let packet = AudioPacket { client_id: id, data: v.packet };
if let Err(e) = h.handle_packet(id, v.packet) { if event_producer_for_task.push_packet(packet).is_ok() {
debug!(target: "chanora_audio", error = %e, "decode failed");
} else {
frames_received_for_task.fetch_add(1, Ordering::Relaxed); frames_received_for_task.fetch_add(1, Ordering::Relaxed);
} }
} }
@@ -1186,7 +1191,7 @@ impl AudioEngine {
output_muted, output_muted,
audio_processing_config, audio_processing_config,
audio_processing_stats, audio_processing_stats,
audio_handler, audio_event_producer: event_producer,
voice_out_tx, voice_out_tx,
voice_activity_selector: cfg.voice_activity_selector.clone(), voice_activity_selector: cfg.voice_activity_selector.clone(),
mic_gain: cfg.mic_gain, mic_gain: cfg.mic_gain,
@@ -1308,9 +1313,6 @@ impl AudioEngine {
audio_processing_config, audio_processing_config,
audio_processing_stats, audio_processing_stats,
audio_handler, audio_handler,
voice_out_tx,
voice_activity_selector: cfg.voice_activity_selector.clone(),
mic_gain: cfg.mic_gain,
_ios_voice_backend: Mutex::new(Some(ios_voice_backend)), _ios_voice_backend: Mutex::new(Some(ios_voice_backend)),
shutdown_tx: Some(shutdown_tx), shutdown_tx: Some(shutdown_tx),
capture_active, capture_active,
@@ -1470,7 +1472,8 @@ impl AudioEngine {
transmit_active: self.transmit_gate.flag_arc(), transmit_active: self.transmit_gate.flag_arc(),
frames_sent: self.frames_sent.clone(), frames_sent: self.frames_sent.clone(),
mic_gain: self.mic_gain, mic_gain: self.mic_gain,
handler: self.audio_handler.clone(), handler: AudioHandler::new(),
event_producer: self.audio_event_producer.clone(),
output_gain: self.output_gain.clone(), output_gain: self.output_gain.clone(),
output_muted: self.output_muted.clone(), output_muted: self.output_muted.clone(),
voice_activity_selector: self.voice_activity_selector.clone(), voice_activity_selector: self.voice_activity_selector.clone(),
@@ -1491,7 +1494,7 @@ impl AudioEngine {
self.voice_out_tx.clone(), self.voice_out_tx.clone(),
self.transmit_gate.clone(), self.transmit_gate.clone(),
self.frames_sent.clone(), self.frames_sent.clone(),
self.audio_handler.clone(), self.audio_event_producer.clone(),
self.output_gain.clone(), self.output_gain.clone(),
self.output_muted.clone(), self.output_muted.clone(),
self.voice_activity_selector.clone(), self.voice_activity_selector.clone(),
@@ -1658,20 +1661,36 @@ impl AudioEngine {
/// `0.0..4.0`. /// `0.0..4.0`.
pub fn set_client_volume(&self, client_id: u64, volume: f32) { pub fn set_client_volume(&self, client_id: u64, volume: f32) {
let clamped = volume.clamp(0.0, 4.0); let clamped = volume.clamp(0.0, 4.0);
match self.audio_handler.lock() { #[cfg(target_os = "android")]
Ok(mut h) => { {
if let Some(q) = h.get_mut_queues().get_mut(&SessionAudioId(client_id)) { let mut cmd = AudioCommand::SetVolume(SessionAudioId(client_id), clamped);
q.volume = clamped; loop {
match self.audio_event_producer.push_control(cmd) {
Ok(()) => break,
Err(returned) => {
cmd = returned;
std::thread::yield_now();
}
} }
} }
Err(e) => { }
tracing::warn!( #[cfg(not(target_os = "android"))]
target: "chanora_audio", {
client_id, match self.audio_handler.lock() {
volume = clamped, Ok(mut h) => {
error = %e, if let Some(q) = h.get_mut_queues().get_mut(&SessionAudioId(client_id)) {
"set_client_volume: audio_handler lock poisoned — volume not applied" q.volume = clamped;
); }
}
Err(e) => {
tracing::warn!(
target: "chanora_audio",
client_id,
volume = clamped,
error = %e,
"set_client_volume: audio_handler lock poisoned — volume not applied"
);
}
} }
} }
} }
+2
View File
@@ -28,6 +28,8 @@
#![warn(missing_docs)] #![warn(missing_docs)]
#[cfg(target_os = "android")]
mod audio_event_queue;
pub mod audio_processing; pub mod audio_processing;
pub mod debug_wav; pub mod debug_wav;
mod engine; mod engine;
@@ -22,6 +22,9 @@ use std::sync::atomic::{AtomicBool, AtomicU32};
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
use tokio::sync::mpsc; use tokio::sync::mpsc;
#[cfg(not(target_os = "android"))]
use tsclientlib::audio::AudioHandler;
#[cfg(target_os = "android")]
use tsclientlib::audio::AudioHandler; use tsclientlib::audio::AudioHandler;
use crate::engine::SessionAudioId; use crate::engine::SessionAudioId;
@@ -67,7 +70,7 @@ pub type BackendEventTx = mpsc::UnboundedSender<BackendEvent>;
pub type AudioSessionId = i32; pub type AudioSessionId = i32;
/// Engine-owned state shared with mobile voice audio callbacks. /// Engine-owned state shared with mobile voice audio callbacks.
#[derive(Clone)] #[cfg_attr(not(target_os = "android"), derive(Clone))]
pub(crate) struct VoiceAudioParams { pub(crate) struct VoiceAudioParams {
/// Opus-encoded voice packets sent on this channel toward the /// Opus-encoded voice packets sent on this channel toward the
/// protocol layer. /// protocol layer.
@@ -78,8 +81,15 @@ pub(crate) struct VoiceAudioParams {
pub frames_sent: Arc<AtomicU32>, pub frames_sent: Arc<AtomicU32>,
/// Pre-encode amplitude scale (1.0 = unity). /// Pre-encode amplitude scale (1.0 = unity).
pub mic_gain: f32, pub mic_gain: f32,
/// AudioHandler owned by the Android output callback.
#[cfg(target_os = "android")]
pub handler: AudioHandler<SessionAudioId>,
/// Producer used by Android engine tasks to feed the output callback.
#[cfg(target_os = "android")]
pub event_producer: crate::audio_event_queue::AudioEventProducer,
/// AudioHandler that inbound decode+mix feeds into; the output /// AudioHandler that inbound decode+mix feeds into; the output
/// callback pulls mixed stereo f32 from it. /// callback pulls mixed stereo f32 from it.
#[cfg(not(target_os = "android"))]
pub handler: Arc<Mutex<AudioHandler<SessionAudioId>>>, pub handler: Arc<Mutex<AudioHandler<SessionAudioId>>>,
/// Master output gain (f32 bits stored in AtomicU32 for lock-free /// Master output gain (f32 bits stored in AtomicU32 for lock-free
/// cross-thread read from the realtime audio callback). /// cross-thread read from the realtime audio callback).