refactor(audio): split engine.rs into engine/ module directory (TODO-022)
Split monolithic engine.rs (3,244 lines) into focused modules: mod.rs (1,106), capture.rs (515), render.rs (176), lifecycle.rs (1,149). Capture and render are desktop-only (cfg-gated). No behavioral changes. All public API paths preserved.
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,515 @@
|
|||||||
|
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
|
use cpal::traits::StreamTrait;
|
||||||
|
use cpal::{SampleFormat, SizedSample};
|
||||||
|
use tracing::{debug, error, warn, info};
|
||||||
|
|
||||||
|
use chanora_protocol::OutPacket;
|
||||||
|
|
||||||
|
use audiopus::coder::Encoder as OpusEncoder;
|
||||||
|
|
||||||
|
use crate::AudioError;
|
||||||
|
|
||||||
|
use super::SAMPLE_RATE;
|
||||||
|
use super::FRAME_SAMPLES;
|
||||||
|
|
||||||
|
const LEVEL_METER_INTERVAL: std::time::Duration = std::time::Duration::from_millis(33);
|
||||||
|
|
||||||
|
pub(super) fn try_open_capture(
|
||||||
|
in_dev: &cpal::Device,
|
||||||
|
voice_out_tx: tokio::sync::mpsc::Sender<OutPacket>,
|
||||||
|
transmit_active: Arc<AtomicBool>,
|
||||||
|
frames_sent: Arc<AtomicU32>,
|
||||||
|
mic_gain: f32,
|
||||||
|
voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>,
|
||||||
|
audio_processing_config: Arc<Mutex<crate::AudioProcessingConfig>>,
|
||||||
|
silero_vad_worker: Arc<Mutex<Option<crate::vad::silero_onnx::SileroOnnxVadWorker>>>,
|
||||||
|
audio_processing_stats: Arc<crate::SharedAudioProcessingStats>,
|
||||||
|
) -> Result<cpal::Stream, AudioError> {
|
||||||
|
let in_cfg = in_dev
|
||||||
|
.default_input_config()
|
||||||
|
.map_err(|e| AudioError::StreamConfig(format!("input default: {e}")))?;
|
||||||
|
let in_sample_rate = in_cfg.sample_rate();
|
||||||
|
let in_channels = in_cfg.channels() as usize;
|
||||||
|
let in_format = in_cfg.sample_format();
|
||||||
|
let mut in_stream_cfg: cpal::StreamConfig = in_cfg.into();
|
||||||
|
#[cfg(target_os = "windows")]
|
||||||
|
{
|
||||||
|
in_stream_cfg.buffer_size = cpal::BufferSize::Fixed(2048);
|
||||||
|
}
|
||||||
|
#[cfg(not(target_os = "windows"))]
|
||||||
|
{
|
||||||
|
in_stream_cfg.buffer_size = cpal::BufferSize::Default;
|
||||||
|
}
|
||||||
|
|
||||||
|
let opus_enc = crate::opus_voice::new_voip_encoder("cpal capture")?;
|
||||||
|
|
||||||
|
let capture_state = Arc::new(Mutex::new(CaptureState::new(
|
||||||
|
opus_enc,
|
||||||
|
in_sample_rate,
|
||||||
|
in_channels,
|
||||||
|
mic_gain,
|
||||||
|
crate::opus_voice::start_out_packet_worker(
|
||||||
|
voice_out_tx,
|
||||||
|
frames_sent.clone(),
|
||||||
|
"cpal-capture",
|
||||||
|
)?,
|
||||||
|
transmit_active,
|
||||||
|
voice_activity_selector,
|
||||||
|
audio_processing_config,
|
||||||
|
silero_vad_worker,
|
||||||
|
audio_processing_stats,
|
||||||
|
)));
|
||||||
|
|
||||||
|
let stream = match in_format {
|
||||||
|
SampleFormat::F32 => build_input_stream::<f32>(in_dev, &in_stream_cfg, capture_state)?,
|
||||||
|
SampleFormat::I16 => build_input_stream::<i16>(in_dev, &in_stream_cfg, capture_state)?,
|
||||||
|
SampleFormat::U16 => build_input_stream::<u16>(in_dev, &in_stream_cfg, capture_state)?,
|
||||||
|
other => {
|
||||||
|
return Err(AudioError::StreamConfig(format!(
|
||||||
|
"unsupported input format: {other:?}"
|
||||||
|
)))
|
||||||
|
}
|
||||||
|
};
|
||||||
|
Ok(stream)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(super) struct CaptureState {
|
||||||
|
encoder: OpusEncoder,
|
||||||
|
pub(super) pcm_accum: Vec<f32>,
|
||||||
|
pub(super) pending_10ms: [f32; crate::frame::FRAME_10MS_SAMPLES],
|
||||||
|
pub(super) pending_10ms_len: usize,
|
||||||
|
pub(super) capture_frame_seq: u64,
|
||||||
|
in_sample_rate: u32,
|
||||||
|
in_channels: usize,
|
||||||
|
mic_gain: f32,
|
||||||
|
resample_pos: f64,
|
||||||
|
resample_last: f32,
|
||||||
|
opus_out: [u8; crate::opus_voice::MAX_OPUS_FRAME],
|
||||||
|
voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender,
|
||||||
|
transmit_active: Arc<AtomicBool>,
|
||||||
|
voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>,
|
||||||
|
vad_detector: crate::vad::WebRtcFallbackVad,
|
||||||
|
silero_vad_worker: Arc<Mutex<Option<crate::vad::silero_onnx::SileroOnnxVadWorker>>>,
|
||||||
|
silero_model_epoch: u64,
|
||||||
|
current_vad_backend: crate::VadBackend,
|
||||||
|
fallback_warned_backend: Option<crate::VadBackend>,
|
||||||
|
vad_state: crate::voice_activity::VoiceActivityStateMachine,
|
||||||
|
audio_processing_config: Arc<Mutex<crate::AudioProcessingConfig>>,
|
||||||
|
mono_scratch: Vec<f32>,
|
||||||
|
frame_scratch: Vec<f32>,
|
||||||
|
audio_processing_stats: Arc<crate::SharedAudioProcessingStats>,
|
||||||
|
last_level_emit: std::time::Instant,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl CaptureState {
|
||||||
|
pub(super) fn new(
|
||||||
|
encoder: OpusEncoder,
|
||||||
|
in_sample_rate: u32,
|
||||||
|
in_channels: usize,
|
||||||
|
mic_gain: f32,
|
||||||
|
voice_out_tx: crate::opus_voice::EncodedVoiceFrameSender,
|
||||||
|
transmit_active: Arc<AtomicBool>,
|
||||||
|
voice_activity_selector: Option<Arc<crate::TransmitModeSelector>>,
|
||||||
|
audio_processing_config: Arc<Mutex<crate::AudioProcessingConfig>>,
|
||||||
|
silero_vad_worker: Arc<Mutex<Option<crate::vad::silero_onnx::SileroOnnxVadWorker>>>,
|
||||||
|
audio_processing_stats: Arc<crate::SharedAudioProcessingStats>,
|
||||||
|
) -> Self {
|
||||||
|
Self {
|
||||||
|
encoder,
|
||||||
|
in_sample_rate,
|
||||||
|
in_channels,
|
||||||
|
mic_gain,
|
||||||
|
pcm_accum: Vec::with_capacity(FRAME_SAMPLES * 2),
|
||||||
|
resample_pos: 0.0,
|
||||||
|
resample_last: 0.0,
|
||||||
|
opus_out: [0u8; crate::opus_voice::MAX_OPUS_FRAME],
|
||||||
|
voice_out_tx,
|
||||||
|
transmit_active,
|
||||||
|
voice_activity_selector,
|
||||||
|
vad_detector: crate::vad::WebRtcFallbackVad::default(),
|
||||||
|
silero_vad_worker,
|
||||||
|
silero_model_epoch: crate::vad::silero_model_epoch(),
|
||||||
|
current_vad_backend: crate::VadBackend::Disabled,
|
||||||
|
fallback_warned_backend: None,
|
||||||
|
capture_frame_seq: 0,
|
||||||
|
vad_state: crate::voice_activity::VoiceActivityStateMachine::default(),
|
||||||
|
audio_processing_config,
|
||||||
|
pending_10ms: [0.0; crate::frame::FRAME_10MS_SAMPLES],
|
||||||
|
pending_10ms_len: 0,
|
||||||
|
mono_scratch: Vec::with_capacity(4096),
|
||||||
|
frame_scratch: Vec::with_capacity(FRAME_SAMPLES),
|
||||||
|
audio_processing_stats,
|
||||||
|
last_level_emit: std::time::Instant::now()
|
||||||
|
.checked_sub(LEVEL_METER_INTERVAL)
|
||||||
|
.unwrap_or_else(std::time::Instant::now),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(super) fn ingest<T: ToF32 + Copy>(&mut self, buf: &[T]) {
|
||||||
|
let in_channels = self.in_channels;
|
||||||
|
let mic_gain = self.mic_gain;
|
||||||
|
self.mono_scratch.clear();
|
||||||
|
let frame_count = buf.len() / in_channels.max(1);
|
||||||
|
self.mono_scratch.reserve(frame_count);
|
||||||
|
for frame in buf.chunks(in_channels) {
|
||||||
|
let sum: f32 = frame.iter().map(|s| s.to_f32_sample()).sum();
|
||||||
|
self.mono_scratch.push(sum / frame.len() as f32);
|
||||||
|
}
|
||||||
|
|
||||||
|
let now = std::time::Instant::now();
|
||||||
|
if now.duration_since(self.last_level_emit) >= LEVEL_METER_INTERVAL {
|
||||||
|
self.last_level_emit = now;
|
||||||
|
self.audio_processing_stats
|
||||||
|
.set_input_dbfs(crate::frame::dbfs(&self.mono_scratch));
|
||||||
|
}
|
||||||
|
|
||||||
|
if mic_gain != 1.0 {
|
||||||
|
for s in &mut self.mono_scratch {
|
||||||
|
*s *= mic_gain;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let vad_start_offset = self.pcm_accum.len();
|
||||||
|
if self.in_sample_rate == SAMPLE_RATE {
|
||||||
|
let (src, dst) = (&self.mono_scratch, &mut self.pcm_accum);
|
||||||
|
dst.extend_from_slice(src);
|
||||||
|
} else {
|
||||||
|
let mono = std::mem::take(&mut self.mono_scratch);
|
||||||
|
self.resample_into_accum(&mono);
|
||||||
|
self.mono_scratch = mono;
|
||||||
|
}
|
||||||
|
|
||||||
|
self.process_pending_vad_frames(vad_start_offset);
|
||||||
|
|
||||||
|
if !self.transmit_active.load(Ordering::Relaxed) {
|
||||||
|
self.pcm_accum.clear();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
while self.pcm_accum.len() >= FRAME_SAMPLES {
|
||||||
|
let frame = &mut self.frame_scratch;
|
||||||
|
frame.clear();
|
||||||
|
frame.extend(self.pcm_accum.drain(..FRAME_SAMPLES));
|
||||||
|
for s in frame.iter_mut() {
|
||||||
|
if *s > 1.0 {
|
||||||
|
*s = 1.0;
|
||||||
|
} else if *s < -1.0 {
|
||||||
|
*s = -1.0;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
match self
|
||||||
|
.encoder
|
||||||
|
.encode_float(&frame[..], &mut self.opus_out[..])
|
||||||
|
{
|
||||||
|
Ok(len) => {
|
||||||
|
crate::opus_voice::send_voip_frame(
|
||||||
|
&self.voice_out_tx,
|
||||||
|
&self.opus_out,
|
||||||
|
len,
|
||||||
|
|| {
|
||||||
|
warn!(
|
||||||
|
target: "chanora_audio",
|
||||||
|
"voice_out queue full; dropping frame"
|
||||||
|
);
|
||||||
|
},
|
||||||
|
|| {
|
||||||
|
warn!(target: "chanora_audio", "voice_out closed; stopping send");
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
error!(target: "chanora_audio", error = %e, "opus encode failed");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(super) fn process_pending_vad_frames(&mut self, start_offset: usize) {
|
||||||
|
let mut offset = start_offset.min(self.pcm_accum.len());
|
||||||
|
while offset < self.pcm_accum.len() {
|
||||||
|
let remaining = crate::frame::FRAME_10MS_SAMPLES - self.pending_10ms_len;
|
||||||
|
let take = remaining.min(self.pcm_accum.len() - offset);
|
||||||
|
self.pending_10ms[self.pending_10ms_len..self.pending_10ms_len + take]
|
||||||
|
.copy_from_slice(&self.pcm_accum[offset..offset + take]);
|
||||||
|
self.pending_10ms_len += take;
|
||||||
|
offset += take;
|
||||||
|
|
||||||
|
if self.pending_10ms_len == crate::frame::FRAME_10MS_SAMPLES {
|
||||||
|
let frame = self.pending_10ms;
|
||||||
|
self.process_10ms_capture_frame(&frame);
|
||||||
|
self.pending_10ms_len = 0;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn mark_vad_fallback_active(&mut self, failed_backend: crate::VadBackend) {
|
||||||
|
self.fallback_warned_backend = Some(failed_backend);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn sync_vad_backend(&mut self, voice_activity_mode: bool, vad_backend: crate::VadBackend) {
|
||||||
|
if !voice_activity_mode {
|
||||||
|
self.current_vad_backend = crate::VadBackend::Disabled;
|
||||||
|
self.fallback_warned_backend = None;
|
||||||
|
self.audio_processing_stats.set_vad_fallback_active(false);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
let silero_epoch = crate::vad::silero_model_epoch();
|
||||||
|
let silero_changed =
|
||||||
|
vad_backend == crate::VadBackend::SileroOnnx && silero_epoch != self.silero_model_epoch;
|
||||||
|
if vad_backend == self.current_vad_backend && !silero_changed {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
self.current_vad_backend = vad_backend;
|
||||||
|
self.silero_model_epoch = silero_epoch;
|
||||||
|
self.fallback_warned_backend = None;
|
||||||
|
self.vad_state.reset();
|
||||||
|
|
||||||
|
match vad_backend {
|
||||||
|
crate::VadBackend::SileroOnnx => {
|
||||||
|
let worker_available = self
|
||||||
|
.silero_vad_worker
|
||||||
|
.try_lock()
|
||||||
|
.map(|worker| worker.is_some())
|
||||||
|
.unwrap_or(false);
|
||||||
|
if worker_available {
|
||||||
|
self.audio_processing_stats.set_vad_fallback_active(false);
|
||||||
|
} else {
|
||||||
|
self.mark_vad_fallback_active(crate::VadBackend::SileroOnnx);
|
||||||
|
self.audio_processing_stats.set_vad_fallback_active(true);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
crate::VadBackend::WebrtcVad => {
|
||||||
|
self.audio_processing_stats.set_vad_fallback_active(false);
|
||||||
|
}
|
||||||
|
crate::VadBackend::EnergyDebug => {
|
||||||
|
self.audio_processing_stats.set_vad_fallback_active(true);
|
||||||
|
}
|
||||||
|
crate::VadBackend::Disabled => {
|
||||||
|
self.audio_processing_stats.set_vad_fallback_active(false);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(super) fn process_10ms_capture_frame(&mut self, frame: &[f32; crate::frame::FRAME_10MS_SAMPLES]) {
|
||||||
|
let input_dbfs = crate::frame::dbfs(frame);
|
||||||
|
let (vad_backend, vad_hangover) = self
|
||||||
|
.audio_processing_config
|
||||||
|
.try_lock()
|
||||||
|
.map(|cfg| (cfg.vad_backend, cfg.vad_hangover_ms))
|
||||||
|
.unwrap_or((
|
||||||
|
crate::VadBackend::WebrtcVad,
|
||||||
|
crate::voice_activity::VAD_HANGOVER_MS,
|
||||||
|
));
|
||||||
|
let voice_activity_mode = self
|
||||||
|
.voice_activity_selector
|
||||||
|
.as_ref()
|
||||||
|
.map(|selector| selector.mode() == crate::TransmitMode::VoiceActivity)
|
||||||
|
.unwrap_or(false);
|
||||||
|
|
||||||
|
if voice_activity_mode {
|
||||||
|
self.sync_vad_backend(true, vad_backend);
|
||||||
|
self.vad_state.configure(
|
||||||
|
crate::voice_activity::VAD_OPEN_AFTER_MS,
|
||||||
|
vad_hangover,
|
||||||
|
crate::voice_activity::VAD_MIN_TX_MS,
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
self.sync_vad_backend(false, vad_backend);
|
||||||
|
}
|
||||||
|
|
||||||
|
let (vad_probability, gate_open, used_fallback_vad) = if voice_activity_mode {
|
||||||
|
self.capture_frame_seq = self.capture_frame_seq.wrapping_add(1);
|
||||||
|
let capture_seq = self.capture_frame_seq;
|
||||||
|
let mut used_fallback_vad = false;
|
||||||
|
let vad = match vad_backend {
|
||||||
|
crate::VadBackend::Disabled => crate::vad::VadOutput {
|
||||||
|
probability: 1.0,
|
||||||
|
speech: true,
|
||||||
|
},
|
||||||
|
crate::VadBackend::SileroOnnx => {
|
||||||
|
let worker_output = {
|
||||||
|
let guard = self.silero_vad_worker.try_lock().ok();
|
||||||
|
guard.and_then(|guard| {
|
||||||
|
let worker = guard.as_ref()?;
|
||||||
|
if worker.try_send(capture_seq, frame) && !worker.is_stale(capture_seq)
|
||||||
|
{
|
||||||
|
let p = worker.latest_probability();
|
||||||
|
Some(crate::vad::VadOutput {
|
||||||
|
probability: p,
|
||||||
|
speech: p >= 0.5,
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
None
|
||||||
|
}
|
||||||
|
})
|
||||||
|
};
|
||||||
|
if let Some(output) = worker_output {
|
||||||
|
output
|
||||||
|
} else {
|
||||||
|
used_fallback_vad = true;
|
||||||
|
self.mark_vad_fallback_active(vad_backend);
|
||||||
|
crate::vad::VoiceActivityDetector::process_10ms(
|
||||||
|
&mut self.vad_detector,
|
||||||
|
frame,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
crate::VadBackend::WebrtcVad | crate::VadBackend::EnergyDebug => {
|
||||||
|
used_fallback_vad = vad_backend == crate::VadBackend::EnergyDebug;
|
||||||
|
crate::vad::VoiceActivityDetector::process_10ms(&mut self.vad_detector, frame)
|
||||||
|
}
|
||||||
|
};
|
||||||
|
(
|
||||||
|
vad.probability,
|
||||||
|
self.vad_state.update(vad.speech),
|
||||||
|
used_fallback_vad,
|
||||||
|
)
|
||||||
|
} else {
|
||||||
|
(0.0, false, false)
|
||||||
|
};
|
||||||
|
self.audio_processing_stats
|
||||||
|
.set_vad_fallback_active(used_fallback_vad);
|
||||||
|
|
||||||
|
let vad_active = voice_activity_mode && gate_open;
|
||||||
|
if let Some(selector) = &self.voice_activity_selector {
|
||||||
|
selector.set_voice_activity_open(vad_active);
|
||||||
|
}
|
||||||
|
self.audio_processing_stats.update_capture(
|
||||||
|
input_dbfs,
|
||||||
|
input_dbfs,
|
||||||
|
vad_probability,
|
||||||
|
vad_active,
|
||||||
|
self.transmit_active.load(Ordering::Relaxed),
|
||||||
|
);
|
||||||
|
self.audio_processing_stats
|
||||||
|
.record_capture_frame(frame.iter().all(|sample| sample.abs() <= 0.000_001));
|
||||||
|
}
|
||||||
|
|
||||||
|
fn resample_into_accum(&mut self, mono: &[f32]) {
|
||||||
|
if mono.is_empty() {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let ratio = self.in_sample_rate as f64 / SAMPLE_RATE as f64;
|
||||||
|
let mut pos = self.resample_pos;
|
||||||
|
while pos < mono.len() as f64 {
|
||||||
|
let i = pos.floor() as isize;
|
||||||
|
let frac = pos - i as f64;
|
||||||
|
let a = if i <= 0 {
|
||||||
|
self.resample_last
|
||||||
|
} else {
|
||||||
|
mono[(i - 1) as usize]
|
||||||
|
};
|
||||||
|
let b = if i < mono.len() as isize {
|
||||||
|
mono[i as usize]
|
||||||
|
} else {
|
||||||
|
a
|
||||||
|
};
|
||||||
|
self.pcm_accum
|
||||||
|
.push((a as f64 + frac * (b - a) as f64) as f32);
|
||||||
|
pos += ratio;
|
||||||
|
}
|
||||||
|
self.resample_pos = pos - mono.len() as f64;
|
||||||
|
self.resample_last = *mono.last().unwrap();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
trait ToF32 {
|
||||||
|
fn to_f32_sample(self) -> f32;
|
||||||
|
}
|
||||||
|
impl ToF32 for f32 {
|
||||||
|
fn to_f32_sample(self) -> f32 {
|
||||||
|
self
|
||||||
|
}
|
||||||
|
}
|
||||||
|
impl ToF32 for i16 {
|
||||||
|
fn to_f32_sample(self) -> f32 {
|
||||||
|
f32::from(self) / f32::from(i16::MAX)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
impl ToF32 for u16 {
|
||||||
|
fn to_f32_sample(self) -> f32 {
|
||||||
|
(f32::from(self) - f32::from(i16::MAX) - 1.0) / f32::from(i16::MAX)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn build_input_stream<T>(
|
||||||
|
device: &cpal::Device,
|
||||||
|
config: &cpal::StreamConfig,
|
||||||
|
state: Arc<Mutex<CaptureState>>,
|
||||||
|
) -> Result<cpal::Stream, AudioError>
|
||||||
|
where
|
||||||
|
T: SizedSample + ToF32 + Send + 'static,
|
||||||
|
{
|
||||||
|
let stream = device
|
||||||
|
.build_input_stream(
|
||||||
|
*config,
|
||||||
|
move |data: &[T], _: &cpal::InputCallbackInfo| {
|
||||||
|
let mut s = state.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
|
s.ingest(data);
|
||||||
|
},
|
||||||
|
move |e| {
|
||||||
|
error!(target: "chanora_audio", error = %e, "input stream error");
|
||||||
|
},
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.map_err(|e| AudioError::Backend(format!("build_input_stream: {e}")))?;
|
||||||
|
Ok(stream)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[doc(hidden)]
|
||||||
|
pub mod bench_seam {
|
||||||
|
use super::{Arc, AtomicBool, AtomicU32, CaptureState, OutPacket};
|
||||||
|
use tokio::sync::mpsc;
|
||||||
|
|
||||||
|
pub struct CaptureBenchHandle {
|
||||||
|
pub(crate) state: CaptureState,
|
||||||
|
_rx: mpsc::Receiver<OutPacket>,
|
||||||
|
transmit_active: Arc<AtomicBool>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl CaptureBenchHandle {
|
||||||
|
pub fn new(in_sample_rate: u32, in_channels: usize) -> Self {
|
||||||
|
let encoder =
|
||||||
|
crate::opus_voice::new_voip_encoder("cpal bench").expect("opus encoder init");
|
||||||
|
let (tx, rx) = mpsc::channel::<OutPacket>(64);
|
||||||
|
let transmit_active = Arc::new(AtomicBool::new(true));
|
||||||
|
let frames_sent = Arc::new(AtomicU32::new(0));
|
||||||
|
let voice_out_tx =
|
||||||
|
crate::opus_voice::start_out_packet_worker(tx, frames_sent, "cpal-bench")
|
||||||
|
.expect("start_out_packet_worker");
|
||||||
|
let state = CaptureState::new(
|
||||||
|
encoder,
|
||||||
|
in_sample_rate,
|
||||||
|
in_channels,
|
||||||
|
1.0,
|
||||||
|
voice_out_tx,
|
||||||
|
transmit_active.clone(),
|
||||||
|
None,
|
||||||
|
Arc::new(std::sync::Mutex::new(
|
||||||
|
crate::AudioProcessingConfig::default(),
|
||||||
|
)),
|
||||||
|
Arc::new(std::sync::Mutex::new(None)),
|
||||||
|
Arc::new(crate::SharedAudioProcessingStats::default()),
|
||||||
|
);
|
||||||
|
Self {
|
||||||
|
state,
|
||||||
|
_rx: rx,
|
||||||
|
transmit_active,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[inline]
|
||||||
|
pub fn ingest_f32(&mut self, buf: &[f32]) {
|
||||||
|
self.state.ingest(buf);
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn set_transmit_active(&self, active: bool) {
|
||||||
|
self.transmit_active
|
||||||
|
.store(active, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,176 @@
|
|||||||
|
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
|
||||||
|
use std::sync::{Arc, Mutex};
|
||||||
|
|
||||||
|
use cpal::traits::StreamTrait;
|
||||||
|
use cpal::SampleFormat;
|
||||||
|
use tracing::{error, warn};
|
||||||
|
use tsclientlib::audio::AudioHandler;
|
||||||
|
|
||||||
|
use crate::AudioError;
|
||||||
|
|
||||||
|
use super::SessionAudioId;
|
||||||
|
use super::SAMPLE_RATE;
|
||||||
|
|
||||||
|
pub(super) fn build_output_stream<T>(
|
||||||
|
device: &cpal::Device,
|
||||||
|
config: &cpal::StreamConfig,
|
||||||
|
handler: Arc<Mutex<AudioHandler<SessionAudioId>>>,
|
||||||
|
output_gain: Arc<AtomicU32>,
|
||||||
|
output_muted: Arc<AtomicBool>,
|
||||||
|
dev_sample_rate: u32,
|
||||||
|
dev_channels: usize,
|
||||||
|
) -> Result<cpal::Stream, AudioError>
|
||||||
|
where
|
||||||
|
T: cpal::SizedSample + FromF32 + Send + 'static,
|
||||||
|
{
|
||||||
|
let resample_ratio = SAMPLE_RATE as f64 / dev_sample_rate as f64;
|
||||||
|
let same_rate = dev_sample_rate == SAMPLE_RATE;
|
||||||
|
let resample_state: Arc<Mutex<PlaybackResampleState>> =
|
||||||
|
Arc::new(Mutex::new(PlaybackResampleState {
|
||||||
|
pos: 0.0,
|
||||||
|
last_l: 0.0,
|
||||||
|
last_r: 0.0,
|
||||||
|
}));
|
||||||
|
let mut scratch: Vec<f32> = Vec::with_capacity(8192);
|
||||||
|
let mut last_slow_warn = std::time::Instant::now()
|
||||||
|
.checked_sub(std::time::Duration::from_secs(2))
|
||||||
|
.unwrap_or_else(std::time::Instant::now);
|
||||||
|
let stream = device
|
||||||
|
.build_output_stream(
|
||||||
|
*config,
|
||||||
|
move |out: &mut [T], _: &cpal::OutputCallbackInfo| {
|
||||||
|
let cb_start = std::time::Instant::now();
|
||||||
|
let muted = output_muted.load(Ordering::Relaxed);
|
||||||
|
let dev_frames = out.len() / dev_channels.max(1);
|
||||||
|
let src_frames = if same_rate {
|
||||||
|
dev_frames
|
||||||
|
} else {
|
||||||
|
((dev_frames as f64 * resample_ratio).ceil() as usize) + 2
|
||||||
|
};
|
||||||
|
let needed = src_frames * 2;
|
||||||
|
if scratch.len() < needed {
|
||||||
|
scratch.resize(needed, 0.0);
|
||||||
|
}
|
||||||
|
scratch[..needed].fill(0.0);
|
||||||
|
{
|
||||||
|
let mut h = handler.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
|
h.fill_buffer(&mut scratch[..needed]);
|
||||||
|
}
|
||||||
|
|
||||||
|
if muted {
|
||||||
|
for dst in out.iter_mut() {
|
||||||
|
*dst = T::from_f32_sample(0.0);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
let gain = f32::from_bits(output_gain.load(Ordering::Relaxed));
|
||||||
|
if same_rate && dev_channels == 2 {
|
||||||
|
for (dst, s) in out.iter_mut().zip(scratch[..needed].iter().copied()) {
|
||||||
|
*dst = T::from_f32_sample(s * gain);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
let mut state = resample_state.lock().unwrap_or_else(|e| e.into_inner());
|
||||||
|
let mut pos = state.pos;
|
||||||
|
let mut last_l = state.last_l;
|
||||||
|
let mut last_r = state.last_r;
|
||||||
|
for frame_idx in 0..dev_frames {
|
||||||
|
let i = pos.floor() as isize;
|
||||||
|
let frac = pos - i as f64;
|
||||||
|
let (a_l, a_r) = if i <= 0 {
|
||||||
|
(last_l, last_r)
|
||||||
|
} else {
|
||||||
|
let idx = ((i - 1) as usize) * 2;
|
||||||
|
(scratch[idx], scratch[idx + 1])
|
||||||
|
};
|
||||||
|
let i_usize = i.max(0) as usize;
|
||||||
|
let (b_l, b_r) = if i_usize < src_frames {
|
||||||
|
let idx = i_usize * 2;
|
||||||
|
(scratch[idx], scratch[idx + 1])
|
||||||
|
} else {
|
||||||
|
(a_l, a_r)
|
||||||
|
};
|
||||||
|
let l = (a_l as f64 + frac * (b_l - a_l) as f64) as f32 * gain;
|
||||||
|
let r = (a_r as f64 + frac * (b_r - a_r) as f64) as f32 * gain;
|
||||||
|
let base = frame_idx * dev_channels;
|
||||||
|
if dev_channels == 1 {
|
||||||
|
out[base] = T::from_f32_sample((l + r) * 0.5);
|
||||||
|
} else {
|
||||||
|
out[base] = T::from_f32_sample(l);
|
||||||
|
if dev_channels >= 2 {
|
||||||
|
out[base + 1] = T::from_f32_sample(r);
|
||||||
|
}
|
||||||
|
for c in 2..dev_channels {
|
||||||
|
out[base + c] = T::from_f32_sample(0.0);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
pos += resample_ratio;
|
||||||
|
}
|
||||||
|
let consumed = pos.floor() as usize;
|
||||||
|
state.pos = pos - consumed as f64;
|
||||||
|
if consumed > 0 && consumed <= src_frames {
|
||||||
|
let idx = (consumed - 1) * 2;
|
||||||
|
last_l = scratch[idx];
|
||||||
|
last_r = scratch[idx + 1];
|
||||||
|
state.last_l = last_l;
|
||||||
|
state.last_r = last_r;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let elapsed = cb_start.elapsed();
|
||||||
|
let period_us = (dev_frames as u64 * 1_000_000) / dev_sample_rate as u64;
|
||||||
|
if elapsed.as_micros() as u64 > period_us / 2
|
||||||
|
&& last_slow_warn.elapsed() > std::time::Duration::from_secs(1)
|
||||||
|
{
|
||||||
|
last_slow_warn = std::time::Instant::now();
|
||||||
|
warn!(
|
||||||
|
target: "chanora_audio",
|
||||||
|
callback_us = elapsed.as_micros() as u64,
|
||||||
|
period_us,
|
||||||
|
dev_frames,
|
||||||
|
"output callback exceeded half the period budget — possible underrun cause"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
},
|
||||||
|
move |e| {
|
||||||
|
error!(target: "chanora_audio", error = %e, "output stream error");
|
||||||
|
},
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.map_err(|e| {
|
||||||
|
error!(
|
||||||
|
target: "chanora_audio",
|
||||||
|
error = %e,
|
||||||
|
requested_channels = config.channels,
|
||||||
|
requested_sample_rate = config.sample_rate,
|
||||||
|
"build_output_stream FAILED"
|
||||||
|
);
|
||||||
|
AudioError::Backend(format!("build_output_stream: {e}"))
|
||||||
|
})?;
|
||||||
|
Ok(stream)
|
||||||
|
}
|
||||||
|
|
||||||
|
struct PlaybackResampleState {
|
||||||
|
pos: f64,
|
||||||
|
last_l: f32,
|
||||||
|
last_r: f32,
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(super) trait FromF32 {
|
||||||
|
fn from_f32_sample(v: f32) -> Self;
|
||||||
|
}
|
||||||
|
impl FromF32 for f32 {
|
||||||
|
fn from_f32_sample(v: f32) -> Self {
|
||||||
|
v
|
||||||
|
}
|
||||||
|
}
|
||||||
|
impl FromF32 for i16 {
|
||||||
|
fn from_f32_sample(v: f32) -> Self {
|
||||||
|
(v.clamp(-1.0, 1.0) * f32::from(i16::MAX)) as i16
|
||||||
|
}
|
||||||
|
}
|
||||||
|
impl FromF32 for u16 {
|
||||||
|
fn from_f32_sample(v: f32) -> Self {
|
||||||
|
let s = (v.clamp(-1.0, 1.0) * f32::from(i16::MAX)) as i32;
|
||||||
|
(s + i32::from(i16::MAX) + 1) as u16
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user