Self healing
This commit is contained in:
@@ -26,10 +26,7 @@ const MAX_PACKET_SIZE: usize = 4096;
|
|||||||
const INITIAL_PACKET_CUSHION: usize = 3; // ~60ms cushion
|
const INITIAL_PACKET_CUSHION: usize = 3; // ~60ms cushion
|
||||||
|
|
||||||
// --- Voice Activity Detection (VAD) Settings ---
|
// --- Voice Activity Detection (VAD) Settings ---
|
||||||
/// RMS energy threshold to classify audio as speech (Range: 0.0 to 1.0).
|
|
||||||
/// 0.01 ≈ -40 dBFS. Adjust upward if room noise is triggering transmission.
|
|
||||||
const VAD_THRESHOLD: f32 = 0.01;
|
const VAD_THRESHOLD: f32 = 0.01;
|
||||||
/// Number of 20ms frames to keep transmitting after falling below threshold (~200ms hangover)
|
|
||||||
const VAD_HANGOVER_FRAMES: usize = 10;
|
const VAD_HANGOVER_FRAMES: usize = 10;
|
||||||
|
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
@@ -39,7 +36,13 @@ const VAD_HANGOVER_FRAMES: usize = 10;
|
|||||||
type AudioProducer = Arc<Mutex<CachingProd<Arc<SharedRb<Heap<f32>>>>>>;
|
type AudioProducer = Arc<Mutex<CachingProd<Arc<SharedRb<Heap<f32>>>>>>;
|
||||||
type AudioConsumer = Arc<Mutex<CachingCons<Arc<SharedRb<Heap<f32>>>>>>;
|
type AudioConsumer = Arc<Mutex<CachingCons<Arc<SharedRb<Heap<f32>>>>>>;
|
||||||
|
|
||||||
|
#[derive(Clone, Default)]
|
||||||
pub struct VoiceState {
|
pub struct VoiceState {
|
||||||
|
pub inner: Arc<VoiceStateInner>,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Default)]
|
||||||
|
pub struct VoiceStateInner {
|
||||||
pub session: Mutex<Option<VoiceSession>>,
|
pub session: Mutex<Option<VoiceSession>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -50,18 +53,19 @@ pub struct VoiceSession {
|
|||||||
pub pin: u64,
|
pub pin: u64,
|
||||||
pub shutdown: Arc<AtomicBool>,
|
pub shutdown: Arc<AtomicBool>,
|
||||||
|
|
||||||
// Tracked settings for diffing config changes
|
|
||||||
pub current_input_device: Option<String>,
|
pub current_input_device: Option<String>,
|
||||||
pub current_output_device: Option<String>,
|
pub current_output_device: Option<String>,
|
||||||
|
|
||||||
// Shared ringbuffer handles for runtime hot-swapping
|
|
||||||
pub producer_in: AudioProducer,
|
pub producer_in: AudioProducer,
|
||||||
pub consumer_out: AudioConsumer,
|
pub consumer_out: AudioConsumer,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl VoiceSession {
|
impl VoiceSession {
|
||||||
/// Re-binds the input stream (microphone) without disconnecting the UDP thread
|
pub fn update_input_device(
|
||||||
pub fn update_input_device(&mut self, device_name: Option<String>) -> Result<(), String> {
|
&mut self,
|
||||||
|
device_name: Option<String>,
|
||||||
|
state_inner: Arc<VoiceStateInner>,
|
||||||
|
) -> Result<(), String> {
|
||||||
let host = cpal::default_host();
|
let host = cpal::default_host();
|
||||||
let device = match &device_name {
|
let device = match &device_name {
|
||||||
Some(name) => host
|
Some(name) => host
|
||||||
@@ -85,6 +89,8 @@ impl VoiceSession {
|
|||||||
.config();
|
.config();
|
||||||
|
|
||||||
let producer = Arc::clone(&self.producer_in);
|
let producer = Arc::clone(&self.producer_in);
|
||||||
|
let inner_clone = Arc::clone(&state_inner);
|
||||||
|
|
||||||
let new_stream = device
|
let new_stream = device
|
||||||
.build_input_stream(
|
.build_input_stream(
|
||||||
input_config,
|
input_config,
|
||||||
@@ -93,7 +99,19 @@ impl VoiceSession {
|
|||||||
let _ = prod.push_slice(data);
|
let _ = prod.push_slice(data);
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|err| eprintln!("[vc] input error: {err}"),
|
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,
|
None,
|
||||||
)
|
)
|
||||||
.map_err(|e| e.to_string())?;
|
.map_err(|e| e.to_string())?;
|
||||||
@@ -102,14 +120,16 @@ impl VoiceSession {
|
|||||||
.play()
|
.play()
|
||||||
.map_err(|e| format!("Failed to play input stream: {e}"))?;
|
.map_err(|e| format!("Failed to play input stream: {e}"))?;
|
||||||
|
|
||||||
// Dropping old stream stops capturing audio hardware
|
|
||||||
self.input_stream = new_stream;
|
self.input_stream = new_stream;
|
||||||
self.current_input_device = device_name;
|
self.current_input_device = device_name;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Re-binds the output stream (speakers/headphones) without interrupting UDP reception
|
pub fn update_output_device(
|
||||||
pub fn update_output_device(&mut self, device_name: Option<String>) -> Result<(), String> {
|
&mut self,
|
||||||
|
device_name: Option<String>,
|
||||||
|
state_inner: Arc<VoiceStateInner>,
|
||||||
|
) -> Result<(), String> {
|
||||||
let host = cpal::default_host();
|
let host = cpal::default_host();
|
||||||
let device = match &device_name {
|
let device = match &device_name {
|
||||||
Some(name) => host
|
Some(name) => host
|
||||||
@@ -133,7 +153,7 @@ impl VoiceSession {
|
|||||||
.config();
|
.config();
|
||||||
|
|
||||||
let consumer = Arc::clone(&self.consumer_out);
|
let consumer = Arc::clone(&self.consumer_out);
|
||||||
let new_stream = build_output_stream(&device, output_config, consumer)?;
|
let new_stream = build_output_stream(&device, output_config, consumer, state_inner)?;
|
||||||
new_stream
|
new_stream
|
||||||
.play()
|
.play()
|
||||||
.map_err(|e| format!("Failed to play output stream: {e}"))?;
|
.map_err(|e| format!("Failed to play output stream: {e}"))?;
|
||||||
@@ -147,14 +167,6 @@ impl VoiceSession {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for VoiceState {
|
|
||||||
fn default() -> Self {
|
|
||||||
Self {
|
|
||||||
session: Mutex::new(None),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
// DEVICES
|
// DEVICES
|
||||||
// ============================================================================
|
// ============================================================================
|
||||||
@@ -184,8 +196,12 @@ pub fn list_output_devices() -> Result<Vec<String>, String> {
|
|||||||
// ============================================================================
|
// ============================================================================
|
||||||
|
|
||||||
#[tauri::command]
|
#[tauri::command]
|
||||||
pub fn disconnect_from_vc(voice_state: State<VoiceState>) -> Result<(), String> {
|
pub fn disconnect_from_vc(voice_state: State<'_, VoiceState>) -> Result<(), String> {
|
||||||
let mut lock = voice_state.session.lock().map_err(|e| e.to_string())?;
|
let mut lock = voice_state
|
||||||
|
.inner
|
||||||
|
.session
|
||||||
|
.lock()
|
||||||
|
.map_err(|e| e.to_string())?;
|
||||||
|
|
||||||
if let Some(session) = lock.take() {
|
if let Some(session) = lock.take() {
|
||||||
session.shutdown.store(true, Ordering::SeqCst);
|
session.shutdown.store(true, Ordering::SeqCst);
|
||||||
@@ -207,12 +223,13 @@ pub fn disconnect_from_vc(voice_state: State<VoiceState>) -> Result<(), String>
|
|||||||
pub fn connect_to_vc(
|
pub fn connect_to_vc(
|
||||||
hostname: String,
|
hostname: String,
|
||||||
pin: u64,
|
pin: u64,
|
||||||
config_state: State<ConfigState>,
|
config_state: State<'_, ConfigState>,
|
||||||
voice_state: State<VoiceState>,
|
voice_state: State<'_, VoiceState>,
|
||||||
) -> Result<(), String> {
|
) -> Result<(), String> {
|
||||||
disconnect_from_vc(voice_state.clone())?;
|
disconnect_from_vc(voice_state.clone())?;
|
||||||
|
|
||||||
let config = config_state.0.lock().unwrap().clone();
|
let config = config_state.0.lock().unwrap().clone();
|
||||||
|
let state_inner = Arc::clone(&voice_state.inner);
|
||||||
|
|
||||||
let socket = Arc::new(
|
let socket = Arc::new(
|
||||||
UdpSocket::bind("0.0.0.0:0").map_err(|e| format!("Failed to bind UDP socket: {e}"))?,
|
UdpSocket::bind("0.0.0.0:0").map_err(|e| format!("Failed to bind UDP socket: {e}"))?,
|
||||||
@@ -266,7 +283,7 @@ pub fn connect_to_vc(
|
|||||||
// Output Ring Buffer setup
|
// Output Ring Buffer setup
|
||||||
let rb_out = HeapRb::<f32>::new(19200);
|
let rb_out = HeapRb::<f32>::new(19200);
|
||||||
let (producer_out, consumer_out) = rb_out.split();
|
let (producer_out, consumer_out) = rb_out.split();
|
||||||
let mut producer_out = producer_out; // Handed to receiver thread
|
let mut producer_out = producer_out;
|
||||||
let shared_consumer_out = Arc::new(Mutex::new(consumer_out));
|
let shared_consumer_out = Arc::new(Mutex::new(consumer_out));
|
||||||
|
|
||||||
let output_config = output_device
|
let output_config = output_device
|
||||||
@@ -278,6 +295,7 @@ pub fn connect_to_vc(
|
|||||||
&output_device,
|
&output_device,
|
||||||
output_config,
|
output_config,
|
||||||
Arc::clone(&shared_consumer_out),
|
Arc::clone(&shared_consumer_out),
|
||||||
|
Arc::clone(&state_inner),
|
||||||
)?;
|
)?;
|
||||||
output_stream
|
output_stream
|
||||||
.play()
|
.play()
|
||||||
@@ -295,6 +313,8 @@ pub fn connect_to_vc(
|
|||||||
.config();
|
.config();
|
||||||
|
|
||||||
let cb_producer = Arc::clone(&shared_producer_in);
|
let cb_producer = Arc::clone(&shared_producer_in);
|
||||||
|
let inner_input_err = Arc::clone(&state_inner);
|
||||||
|
|
||||||
let input_stream = input_device
|
let input_stream = input_device
|
||||||
.build_input_stream(
|
.build_input_stream(
|
||||||
input_config,
|
input_config,
|
||||||
@@ -303,12 +323,24 @@ pub fn connect_to_vc(
|
|||||||
let _ = prod.push_slice(data);
|
let _ = prod.push_slice(data);
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|err| eprintln!("[vc] input error: {err}"),
|
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,
|
None,
|
||||||
)
|
)
|
||||||
.map_err(|e| e.to_string())?;
|
.map_err(|e| e.to_string())?;
|
||||||
|
|
||||||
// UDP Sender Thread (With Voice Activity Detection)
|
// UDP Sender Thread
|
||||||
{
|
{
|
||||||
let input_socket = socket.clone();
|
let input_socket = socket.clone();
|
||||||
let shutdown = shutdown.clone();
|
let shutdown = shutdown.clone();
|
||||||
@@ -323,11 +355,9 @@ pub fn connect_to_vc(
|
|||||||
if consumer_in.occupied_len() >= PACKET_SAMPLES {
|
if consumer_in.occupied_len() >= PACKET_SAMPLES {
|
||||||
let _ = consumer_in.pop_slice(&mut frame_buf);
|
let _ = consumer_in.pop_slice(&mut frame_buf);
|
||||||
|
|
||||||
// 1. Calculate RMS energy of current audio frame
|
|
||||||
let sum_squares: f32 = frame_buf.iter().map(|&s| s * s).sum();
|
let sum_squares: f32 = frame_buf.iter().map(|&s| s * s).sum();
|
||||||
let rms = (sum_squares / PACKET_SAMPLES as f32).sqrt();
|
let rms = (sum_squares / PACKET_SAMPLES as f32).sqrt();
|
||||||
|
|
||||||
// 2. Check threshold and manage hangover counter
|
|
||||||
let is_speaking = if rms >= VAD_THRESHOLD {
|
let is_speaking = if rms >= VAD_THRESHOLD {
|
||||||
hangover_counter = VAD_HANGOVER_FRAMES;
|
hangover_counter = VAD_HANGOVER_FRAMES;
|
||||||
true
|
true
|
||||||
@@ -338,7 +368,6 @@ pub fn connect_to_vc(
|
|||||||
false
|
false
|
||||||
};
|
};
|
||||||
|
|
||||||
// 3. Only encode and transmit if VAD is active
|
|
||||||
if is_speaking {
|
if is_speaking {
|
||||||
net_packet[0..4].copy_from_slice(&sequence.to_be_bytes());
|
net_packet[0..4].copy_from_slice(&sequence.to_be_bytes());
|
||||||
for (i, sample) in frame_buf.iter().enumerate() {
|
for (i, sample) in frame_buf.iter().enumerate() {
|
||||||
@@ -433,7 +462,7 @@ pub fn connect_to_vc(
|
|||||||
.play()
|
.play()
|
||||||
.map_err(|e| format!("Failed to start input: {e}"))?;
|
.map_err(|e| format!("Failed to start input: {e}"))?;
|
||||||
|
|
||||||
*voice_state.session.lock().unwrap() = Some(VoiceSession {
|
*state_inner.session.lock().unwrap() = Some(VoiceSession {
|
||||||
input_stream,
|
input_stream,
|
||||||
output_stream,
|
output_stream,
|
||||||
socket,
|
socket,
|
||||||
@@ -456,9 +485,11 @@ fn build_output_stream(
|
|||||||
device: &cpal::Device,
|
device: &cpal::Device,
|
||||||
config: cpal::StreamConfig,
|
config: cpal::StreamConfig,
|
||||||
consumer: AudioConsumer,
|
consumer: AudioConsumer,
|
||||||
|
state_inner: Arc<VoiceStateInner>,
|
||||||
) -> Result<cpal::Stream, String> {
|
) -> Result<cpal::Stream, String> {
|
||||||
let channels = config.channels as usize;
|
let channels = config.channels as usize;
|
||||||
let mut last_sample = 0.0f32;
|
let mut last_sample = 0.0f32;
|
||||||
|
let inner_output_err = Arc::clone(&state_inner);
|
||||||
|
|
||||||
device
|
device
|
||||||
.build_output_stream(
|
.build_output_stream(
|
||||||
@@ -487,7 +518,19 @@ fn build_output_stream(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|err| eprintln!("[vc] output error: {err}"),
|
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,
|
None,
|
||||||
)
|
)
|
||||||
.map_err(|e| e.to_string())
|
.map_err(|e| e.to_string())
|
||||||
@@ -502,7 +545,6 @@ impl Drop for VoiceSession {
|
|||||||
let _ = output.pause();
|
let _ = output.pause();
|
||||||
}
|
}
|
||||||
|
|
||||||
// Output hardware stream handle clean-up
|
|
||||||
eprintln!("[vc] VoiceSession dropped and audio streams paused.");
|
eprintln!("[vc] VoiceSession dropped and audio streams paused.");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user