Improved resampler
This commit is contained in:
+176
-135
@@ -19,20 +19,16 @@ use tauri::State;
|
||||
|
||||
use crate::commands::config::ConfigState;
|
||||
|
||||
const PACKET_SAMPLES: usize = 960; // 20ms @ 48kHz
|
||||
const TARGET_SAMPLE_RATE: u32 = 48_000;
|
||||
const PACKET_SAMPLES: usize = 960; // 20ms @ 48kHz mono
|
||||
const HEADER_SIZE: usize = 4;
|
||||
const MAX_PACKET_SIZE: usize = 4096;
|
||||
|
||||
const INITIAL_PACKET_CUSHION: usize = 3; // ~60ms cushion
|
||||
const INITIAL_PACKET_CUSHION: usize = 3;
|
||||
|
||||
// --- Voice Activity Detection (VAD) Settings ---
|
||||
const VAD_THRESHOLD: f32 = 0.01;
|
||||
const VAD_HANGOVER_FRAMES: usize = 10;
|
||||
|
||||
// ============================================================================
|
||||
// STATE & TYPES
|
||||
// ============================================================================
|
||||
|
||||
type AudioProducer = Arc<Mutex<CachingProd<Arc<SharedRb<Heap<f32>>>>>>;
|
||||
type AudioConsumer = Arc<Mutex<CachingCons<Arc<SharedRb<Heap<f32>>>>>>;
|
||||
|
||||
@@ -60,6 +56,39 @@ pub struct VoiceSession {
|
||||
pub consumer_out: AudioConsumer,
|
||||
}
|
||||
|
||||
// Simple Linear Resampler for real-time audio conversion
|
||||
struct LinearResampler {
|
||||
phase: f64,
|
||||
}
|
||||
|
||||
impl LinearResampler {
|
||||
fn new() -> Self {
|
||||
Self { phase: 0.0 }
|
||||
}
|
||||
|
||||
/// Resamples dynamic buffers from `src_rate` to `dst_rate`
|
||||
fn process(&mut self, input: &[f32], src_rate: u32, dst_rate: u32, output: &mut Vec<f32>) {
|
||||
if src_rate == dst_rate {
|
||||
output.extend_from_slice(input);
|
||||
return;
|
||||
}
|
||||
|
||||
let ratio = src_rate as f64 / dst_rate as f64;
|
||||
while self.phase < input.len() as f64 {
|
||||
let idx = self.phase as usize;
|
||||
let frac = (self.phase - idx as f64) as f32;
|
||||
let next_idx = (idx + 1).min(input.len() - 1);
|
||||
|
||||
let sample = input[idx] * (1.0 - frac) + input[next_idx] * frac;
|
||||
output.push(sample);
|
||||
|
||||
self.phase += ratio;
|
||||
}
|
||||
|
||||
self.phase -= input.len() as f64;
|
||||
}
|
||||
}
|
||||
|
||||
impl VoiceSession {
|
||||
pub fn update_input_device(
|
||||
&mut self,
|
||||
@@ -89,32 +118,7 @@ impl VoiceSession {
|
||||
.config();
|
||||
|
||||
let producer = Arc::clone(&self.producer_in);
|
||||
let inner_clone = Arc::clone(&state_inner);
|
||||
|
||||
let new_stream = device
|
||||
.build_input_stream(
|
||||
input_config,
|
||||
move |data: &[f32], _| {
|
||||
if let Ok(mut prod) = producer.lock() {
|
||||
let _ = prod.push_slice(data);
|
||||
}
|
||||
},
|
||||
move |err| {
|
||||
eprintln!("[vc] Input error: {err}. Attempting input stream recovery...");
|
||||
if let Ok(mut lock) = inner_clone.session.lock() {
|
||||
if let Some(session) = lock.as_mut() {
|
||||
let target_device = session.current_input_device.clone();
|
||||
if let Err(e) =
|
||||
session.update_input_device(target_device, Arc::clone(&inner_clone))
|
||||
{
|
||||
eprintln!("[vc] Input recovery failed: {e}");
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
None,
|
||||
)
|
||||
.map_err(|e| e.to_string())?;
|
||||
let new_stream = build_input_stream(&device, input_config, producer, state_inner)?;
|
||||
|
||||
new_stream
|
||||
.play()
|
||||
@@ -167,9 +171,136 @@ impl VoiceSession {
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// DEVICES
|
||||
// ============================================================================
|
||||
// Helper to build normalized Input Stream (Resampled & Downmixed to 48kHz Mono)
|
||||
fn build_input_stream(
|
||||
device: &cpal::Device,
|
||||
config: cpal::StreamConfig,
|
||||
producer: AudioProducer,
|
||||
state_inner: Arc<VoiceStateInner>,
|
||||
) -> Result<cpal::Stream, String> {
|
||||
let native_sample_rate = config.sample_rate;
|
||||
let channels = config.channels as usize;
|
||||
let mut resampler = LinearResampler::new();
|
||||
let mut mono_buffer = Vec::with_capacity(2048);
|
||||
let mut resampled_buffer = Vec::with_capacity(2048);
|
||||
|
||||
let inner_input_err = Arc::clone(&state_inner);
|
||||
|
||||
device
|
||||
.build_input_stream(
|
||||
config,
|
||||
move |data: &[f32], _| {
|
||||
mono_buffer.clear();
|
||||
resampled_buffer.clear();
|
||||
|
||||
// Downmix channels to mono
|
||||
for chunk in data.chunks_exact(channels) {
|
||||
let sum: f32 = chunk.iter().sum();
|
||||
mono_buffer.push(sum / channels as f32);
|
||||
}
|
||||
|
||||
// Resample to 48kHz standard target
|
||||
resampler.process(
|
||||
&mono_buffer,
|
||||
native_sample_rate,
|
||||
TARGET_SAMPLE_RATE,
|
||||
&mut resampled_buffer,
|
||||
);
|
||||
|
||||
if let Ok(mut prod) = producer.lock() {
|
||||
let _ = prod.push_slice(&resampled_buffer);
|
||||
}
|
||||
},
|
||||
move |err| {
|
||||
eprintln!("[vc] Input error: {err}. Attempting recovery...");
|
||||
if let Ok(mut lock) = inner_input_err.session.lock() {
|
||||
if let Some(session) = lock.as_mut() {
|
||||
let target_device = session.current_input_device.clone();
|
||||
let _ = session
|
||||
.update_input_device(target_device, Arc::clone(&inner_input_err));
|
||||
}
|
||||
}
|
||||
},
|
||||
None,
|
||||
)
|
||||
.map_err(|e| e.to_string())
|
||||
}
|
||||
|
||||
// Helper to build normalized Output Stream (48kHz Mono -> Device Native Channels & Rate)
|
||||
fn build_output_stream(
|
||||
device: &cpal::Device,
|
||||
config: cpal::StreamConfig,
|
||||
consumer: AudioConsumer,
|
||||
state_inner: Arc<VoiceStateInner>,
|
||||
) -> Result<cpal::Stream, String> {
|
||||
let native_sample_rate = config.sample_rate;
|
||||
let channels = config.channels as usize;
|
||||
let mut resampler = LinearResampler::new();
|
||||
let mut raw_mono_samples = Vec::with_capacity(2048);
|
||||
let mut resampled_mono = Vec::with_capacity(2048);
|
||||
let mut last_sample = 0.0f32;
|
||||
|
||||
let inner_output_err = Arc::clone(&state_inner);
|
||||
|
||||
device
|
||||
.build_output_stream(
|
||||
config,
|
||||
move |data: &mut [f32], _| {
|
||||
let required_mono_samples = (data.len() / channels) * TARGET_SAMPLE_RATE as usize
|
||||
/ native_sample_rate as usize;
|
||||
|
||||
raw_mono_samples.clear();
|
||||
resampled_mono.clear();
|
||||
|
||||
if let Ok(mut cons) = consumer.lock() {
|
||||
for _ in 0..required_mono_samples {
|
||||
if let Some(s) = cons.try_pop() {
|
||||
last_sample = s;
|
||||
raw_mono_samples.push(s);
|
||||
} else {
|
||||
// Exponential decay to prevent clicking when underflowing
|
||||
last_sample *= 0.92;
|
||||
raw_mono_samples.push(last_sample);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Resample from 48kHz mono to target native output rate
|
||||
resampler.process(
|
||||
&raw_mono_samples,
|
||||
TARGET_SAMPLE_RATE,
|
||||
native_sample_rate,
|
||||
&mut resampled_mono,
|
||||
);
|
||||
|
||||
// Interleave mono into hardware channels
|
||||
let mut res_idx = 0;
|
||||
let mut out_idx = 0;
|
||||
while out_idx < data.len() && res_idx < resampled_mono.len() {
|
||||
let mono_val = resampled_mono[res_idx];
|
||||
for ch in 0..channels {
|
||||
if out_idx + ch < data.len() {
|
||||
data[out_idx + ch] = mono_val;
|
||||
}
|
||||
}
|
||||
out_idx += channels;
|
||||
res_idx += 1;
|
||||
}
|
||||
},
|
||||
move |err| {
|
||||
eprintln!("[vc] Output error: {err}. Attempting recovery...");
|
||||
if let Ok(mut lock) = inner_output_err.session.lock() {
|
||||
if let Some(session) = lock.as_mut() {
|
||||
let target_device = session.current_output_device.clone();
|
||||
let _ = session
|
||||
.update_output_device(target_device, Arc::clone(&inner_output_err));
|
||||
}
|
||||
}
|
||||
},
|
||||
None,
|
||||
)
|
||||
.map_err(|e| e.to_string())
|
||||
}
|
||||
|
||||
#[tauri::command]
|
||||
pub fn list_input_devices() -> Result<Vec<String>, String> {
|
||||
@@ -191,10 +322,6 @@ pub fn list_output_devices() -> Result<Vec<String>, String> {
|
||||
.collect())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// DISCONNECT
|
||||
// ============================================================================
|
||||
|
||||
#[tauri::command]
|
||||
pub fn disconnect_from_vc(voice_state: State<'_, VoiceState>) -> Result<(), String> {
|
||||
let mut lock = voice_state
|
||||
@@ -205,7 +332,6 @@ pub fn disconnect_from_vc(voice_state: State<'_, VoiceState>) -> Result<(), Stri
|
||||
|
||||
if let Some(session) = lock.take() {
|
||||
session.shutdown.store(true, Ordering::SeqCst);
|
||||
|
||||
let _ = session.input_stream.pause();
|
||||
if let Ok(output) = session.output_stream.lock() {
|
||||
let _ = output.pause();
|
||||
@@ -215,10 +341,6 @@ pub fn disconnect_from_vc(voice_state: State<'_, VoiceState>) -> Result<(), Stri
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// CONNECT
|
||||
// ============================================================================
|
||||
|
||||
#[tauri::command]
|
||||
pub fn connect_to_vc(
|
||||
hostname: String,
|
||||
@@ -280,7 +402,7 @@ pub fn connect_to_vc(
|
||||
|
||||
let shutdown = Arc::new(AtomicBool::new(false));
|
||||
|
||||
// Output Ring Buffer setup
|
||||
// Ring Buffer Setup
|
||||
let rb_out = HeapRb::<f32>::new(19200);
|
||||
let (producer_out, consumer_out) = rb_out.split();
|
||||
let mut producer_out = producer_out;
|
||||
@@ -302,7 +424,6 @@ pub fn connect_to_vc(
|
||||
.map_err(|e| format!("Failed to start output: {e}"))?;
|
||||
let output_stream = Arc::new(Mutex::new(output_stream));
|
||||
|
||||
// Input Ring Buffer setup
|
||||
let rb_in = HeapRb::<f32>::new(19200);
|
||||
let (producer_in, mut consumer_in) = rb_in.split();
|
||||
let shared_producer_in = Arc::new(Mutex::new(producer_in));
|
||||
@@ -312,35 +433,14 @@ pub fn connect_to_vc(
|
||||
.map_err(|e| format!("Failed to get default input config: {e}"))?
|
||||
.config();
|
||||
|
||||
let cb_producer = Arc::clone(&shared_producer_in);
|
||||
let inner_input_err = Arc::clone(&state_inner);
|
||||
let input_stream = build_input_stream(
|
||||
&input_device,
|
||||
input_config,
|
||||
Arc::clone(&shared_producer_in),
|
||||
Arc::clone(&state_inner),
|
||||
)?;
|
||||
|
||||
let input_stream = input_device
|
||||
.build_input_stream(
|
||||
input_config,
|
||||
move |data: &[f32], _| {
|
||||
if let Ok(mut prod) = cb_producer.lock() {
|
||||
let _ = prod.push_slice(data);
|
||||
}
|
||||
},
|
||||
move |err| {
|
||||
eprintln!("[vc] Input error: {err}. Attempting input stream recovery...");
|
||||
if let Ok(mut lock) = inner_input_err.session.lock() {
|
||||
if let Some(session) = lock.as_mut() {
|
||||
let target_device = session.current_input_device.clone();
|
||||
if let Err(e) =
|
||||
session.update_input_device(target_device, Arc::clone(&inner_input_err))
|
||||
{
|
||||
eprintln!("[vc] Input recovery failed: {e}");
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
None,
|
||||
)
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
// UDP Sender Thread
|
||||
// Sender Thread
|
||||
{
|
||||
let input_socket = socket.clone();
|
||||
let shutdown = shutdown.clone();
|
||||
@@ -386,7 +486,7 @@ pub fn connect_to_vc(
|
||||
});
|
||||
}
|
||||
|
||||
// UDP Receiver Thread
|
||||
// Receiver Thread
|
||||
{
|
||||
let socket = socket.clone();
|
||||
let shutdown = shutdown.clone();
|
||||
@@ -477,65 +577,6 @@ pub fn connect_to_vc(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// AUDIO CALLBACK
|
||||
// ============================================================================
|
||||
|
||||
fn build_output_stream(
|
||||
device: &cpal::Device,
|
||||
config: cpal::StreamConfig,
|
||||
consumer: AudioConsumer,
|
||||
state_inner: Arc<VoiceStateInner>,
|
||||
) -> Result<cpal::Stream, String> {
|
||||
let channels = config.channels as usize;
|
||||
let mut last_sample = 0.0f32;
|
||||
let inner_output_err = Arc::clone(&state_inner);
|
||||
|
||||
device
|
||||
.build_output_stream(
|
||||
config,
|
||||
move |data: &mut [f32], _| {
|
||||
let mut idx = 0;
|
||||
let mut cons_guard = consumer.lock().ok();
|
||||
|
||||
while idx < data.len() {
|
||||
let sample = cons_guard.as_mut().and_then(|c| c.try_pop());
|
||||
if let Some(mono_sample) = sample {
|
||||
last_sample = mono_sample;
|
||||
for ch in 0..channels {
|
||||
data[idx + ch] = mono_sample;
|
||||
}
|
||||
idx += channels;
|
||||
} else {
|
||||
while idx < data.len() {
|
||||
last_sample *= 0.92;
|
||||
for ch in 0..channels {
|
||||
data[idx + ch] = last_sample;
|
||||
}
|
||||
idx += channels;
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
},
|
||||
move |err| {
|
||||
eprintln!("[vc] Output error: {err}. Attempting output stream recovery...");
|
||||
if let Ok(mut lock) = inner_output_err.session.lock() {
|
||||
if let Some(session) = lock.as_mut() {
|
||||
let target_device = session.current_output_device.clone();
|
||||
if let Err(e) = session
|
||||
.update_output_device(target_device, Arc::clone(&inner_output_err))
|
||||
{
|
||||
eprintln!("[vc] Output recovery failed: {e}");
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
None,
|
||||
)
|
||||
.map_err(|e| e.to_string())
|
||||
}
|
||||
|
||||
impl Drop for VoiceSession {
|
||||
fn drop(&mut self) {
|
||||
self.shutdown.store(true, Ordering::SeqCst);
|
||||
|
||||
Reference in New Issue
Block a user