From cee32544d0ff9bf6dc31c47b9f2f67de37c52f33 Mon Sep 17 00:00:00 2001 From: Nahim El Atmani <2959826+Naam@users.noreply.github.com> Date: Fri, 7 Aug 2026 12:24:39 -0700 Subject: [PATCH] fix: recover when microphone disconnects --- src-tauri/src/audio_toolkit/audio/recorder.rs | 116 ++++++++++++++---- src-tauri/src/managers/audio.rs | 38 ++++-- src/stores/settingsStore.ts | 3 + 3 files changed, 120 insertions(+), 37 deletions(-) diff --git a/src-tauri/src/audio_toolkit/audio/recorder.rs b/src-tauri/src/audio_toolkit/audio/recorder.rs index 9bc77d05f0..280ecfc2ca 100644 --- a/src-tauri/src/audio_toolkit/audio/recorder.rs +++ b/src-tauri/src/audio_toolkit/audio/recorder.rs @@ -86,6 +86,8 @@ pub struct AudioRecorder { /// cleared whenever an open fails so a stale rate/format self-heals on the /// caller's retry. config_cache: Arc>>, + /// Set by cpal when the active input stream can no longer capture. + stream_error: Arc, } impl AudioRecorder { @@ -99,6 +101,7 @@ impl AudioRecorder { audio_cb: None, selected_channel: None, config_cache: Arc::new(Mutex::new(None)), + stream_error: Arc::new(AtomicBool::new(false)), }) } @@ -150,16 +153,15 @@ impl AudioRecorder { pub fn open(&mut self, device: Option) -> Result<(), Box> { if self.worker_handle.is_some() { - if !self.is_capture_worker_dead() { + if !self.needs_reopen() { return Ok(()); // already open } - // The worker exited on its own (see `is_capture_worker_dead`). Reap - // it so we rebuild the stream below instead of handing the caller - // back a recorder whose channels are already closed. - log::warn!("Capture worker exited; rebuilding microphone stream"); + log::warn!("Capture stream failed; rebuilding microphone stream"); let _ = self.close(); } + self.stream_error.store(false, Ordering::Relaxed); + let (sample_tx, sample_rx) = mpsc::channel::(); let (cmd_tx, cmd_rx) = mpsc::channel::(); let (init_tx, init_rx) = mpsc::sync_channel::>(1); @@ -180,6 +182,7 @@ impl AudioRecorder { let audio_cb = self.audio_cb.clone(); let selected_channel = self.selected_channel; let config_cache = Arc::clone(&self.config_cache); + let stream_error = Arc::clone(&self.stream_error); let worker = std::thread::spawn(move || { let stop_flag = Arc::new(AtomicBool::new(false)); @@ -235,6 +238,7 @@ impl AudioRecorder { channels, selected_channel, stop_flag_for_stream, + Arc::clone(&stream_error), ) .map_err(|e| format!("Failed to build input stream: {e}"))?, cpal::SampleFormat::I8 => AudioRecorder::build_stream::( @@ -244,6 +248,7 @@ impl AudioRecorder { channels, selected_channel, stop_flag_for_stream, + Arc::clone(&stream_error), ) .map_err(|e| format!("Failed to build input stream: {e}"))?, cpal::SampleFormat::I16 => AudioRecorder::build_stream::( @@ -253,6 +258,7 @@ impl AudioRecorder { channels, selected_channel, stop_flag_for_stream, + Arc::clone(&stream_error), ) .map_err(|e| format!("Failed to build input stream: {e}"))?, cpal::SampleFormat::I32 => AudioRecorder::build_stream::( @@ -262,6 +268,7 @@ impl AudioRecorder { channels, selected_channel, stop_flag_for_stream, + Arc::clone(&stream_error), ) .map_err(|e| format!("Failed to build input stream: {e}"))?, cpal::SampleFormat::F32 => AudioRecorder::build_stream::( @@ -271,6 +278,7 @@ impl AudioRecorder { channels, selected_channel, stop_flag_for_stream, + Arc::clone(&stream_error), ) .map_err(|e| format!("Failed to build input stream: {e}"))?, sample_format => { @@ -370,18 +378,16 @@ impl AudioRecorder { Ok(resp_rx.recv()?) // wait for the samples } - /// True once the capture worker has exited without anyone calling `close`. + /// True when the active capture stream must be rebuilt. /// - /// `run_consumer` is driven entirely by the sample channel, so when cpal - /// tears the stream down mid-session (device unplugged, USB/Bluetooth - /// dropout) `sample_rx.recv()` returns `Err`, the loop ends and the worker - /// thread finishes. `cmd_tx` and `worker_handle` are still populated at - /// that point, so the recorder looks open from the outside while every - /// command sent to it fails on a closed channel. - pub fn is_capture_worker_dead(&self) -> bool { - self.worker_handle - .as_ref() - .is_some_and(|handle| handle.is_finished()) + /// cpal may report a device disconnect asynchronously without closing its + /// callback channel, so also honor the error callback's explicit flag. + pub fn needs_reopen(&self) -> bool { + self.stream_error.load(Ordering::Relaxed) + || self + .worker_handle + .as_ref() + .is_some_and(|handle| handle.is_finished()) } pub fn close(&mut self) -> Result<(), Box> { @@ -402,6 +408,7 @@ impl AudioRecorder { channels: usize, selected_channel: Option, stop_flag: Arc, + stream_error: Arc, ) -> Result where T: Sample + SizedSample + Send + 'static, @@ -463,7 +470,10 @@ impl AudioRecorder { device.build_input_stream( &config.clone().into(), stream_cb, - |err| log::error!("Stream error: {}", err), + move |err| { + log::error!("Stream error: {}", err); + stream_error.store(true, Ordering::Relaxed); + }, None, ) } @@ -546,15 +556,60 @@ pub fn is_no_input_device_error(error_message: &str) -> bool { #[cfg(test)] mod tests { - use super::{is_microphone_access_denied, is_no_input_device_error, AudioRecorder}; + use super::{ + is_microphone_access_denied, is_no_input_device_error, run_consumer, AudioRecorder, Cmd, + }; + use std::{ + sync::{ + atomic::{AtomicBool, Ordering}, + mpsc, Arc, + }, + thread, + time::{Duration, Instant}, + }; #[test] - fn unopened_recorder_is_not_reported_dead() { + fn unopened_recorder_does_not_need_reopen() { // No worker has been spawned yet, so there is nothing to reap. Guards // against inverting the "no worker" case, which would make every first // open() take the rebuild path. let recorder = AudioRecorder::new().expect("recorder"); - assert!(!recorder.is_capture_worker_dead()); + assert!(!recorder.needs_reopen()); + } + + #[test] + fn stream_error_requires_reopen() { + let recorder = AudioRecorder::new().expect("recorder"); + recorder.stream_error.store(true, Ordering::Relaxed); + assert!(recorder.needs_reopen()); + } + + #[test] + fn shutdown_is_processed_without_audio_samples() { + let (sample_tx, sample_rx) = mpsc::channel(); + let (cmd_tx, cmd_rx) = mpsc::channel(); + let (done_tx, done_rx) = mpsc::channel(); + let worker = thread::spawn(move || { + run_consumer( + 48_000, + None, + sample_rx, + cmd_rx, + None, + None, + Arc::new(AtomicBool::new(false)), + Instant::now(), + ); + let _ = done_tx.send(()); + }); + + cmd_tx.send(Cmd::Shutdown).expect("send shutdown"); + let stopped = done_rx.recv_timeout(Duration::from_secs(1)); + + // Unblock the old implementation so a failing test still exits cleanly. + drop(sample_tx); + worker.join().expect("join consumer"); + assert!(stopped.is_ok(), "shutdown waited for an audio sample"); } #[test] @@ -680,19 +735,30 @@ fn run_consumer( } } - // Runs until the stream closes and `recv` returns `Err`. - while let Ok(chunk) = sample_rx.recv() { + // Poll commands even when a disconnected device stops producing samples + // without closing its CoreAudio stream. + loop { + let mut pending = match sample_rx.recv_timeout(Duration::from_millis(50)) { + Ok(chunk) => Some(chunk), + Err(mpsc::RecvTimeoutError::Timeout) => None, + Err(mpsc::RecvTimeoutError::Disconnected) => break, + }; + // Handle pending commands BEFORE the in-flight chunk so a Start // captures it. Commands used to be polled after processing, which // silently dropped one buffer period of audio (~10ms built-in, up to // ~100ms on Bluetooth) at every recording start. - let mut pending = Some(chunk); while let Ok(cmd) = cmd_rx.try_recv() { match cmd { Cmd::Start(policy, sent_at) => { log::debug!( - "Cmd::Start processed {:?} after send; capture begins with the in-flight chunk", - sent_at.elapsed() + "Cmd::Start processed {:?} after send; capture begins with {} chunk", + sent_at.elapsed(), + if pending.is_some() { + "the in-flight" + } else { + "the next available" + } ); awaiting_first_captured_chunk = Some(Instant::now()); stop_flag.store(false, Ordering::Relaxed); diff --git a/src-tauri/src/managers/audio.rs b/src-tauri/src/managers/audio.rs index 8cf50a2e5a..cd6276106c 100644 --- a/src-tauri/src/managers/audio.rs +++ b/src-tauri/src/managers/audio.rs @@ -8,14 +8,14 @@ use crate::audio_toolkit::{ }; use crate::helpers::clamshell; use crate::managers::transcription::StreamRouter; -use crate::settings::{get_settings, AppSettings}; +use crate::settings::{get_settings, write_settings, AppSettings}; use crate::utils; use log::{debug, error, info, trace, warn}; use std::path::Path; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -use tauri::Manager; +use tauri::{Emitter, Manager}; const STREAM_IDLE_TIMEOUT: Duration = Duration::from_secs(30); const VAD_THRESHOLD: f32 = 0.3; @@ -532,19 +532,17 @@ impl AudioRecordingManager { let mut open_flag = self.is_open.lock().unwrap(); if *open_flag { // `is_open` only records that we opened a stream at some point, not - // that one is still running. If the capture worker has since exited - // (mic unplugged mid-session, USB dropout), returning Ok here hands - // the caller a dead recorder: it captures nothing, then fails in - // stop() on the closed channel, and stays wedged until the - // on-demand close timeout eventually resets the manager. - let worker_dead = self + // that one is still running. If capture has since failed (mic + // unplugged mid-session, USB dropout), rebuild it before the next + // recording instead of handing the caller a stalled recorder. + let needs_reopen = self .recorder .lock() .unwrap() .as_ref() - .is_some_and(|rec| rec.is_capture_worker_dead()); + .is_some_and(|rec| rec.needs_reopen()); - if !worker_dead { + if !needs_reopen { // trace, not debug: with the aliveness check in // try_start_recording this now fires on every keypress in // always-on mode. @@ -564,12 +562,28 @@ impl AudioRecordingManager { } } if let Some(rec) = self.recorder.lock().unwrap().as_mut() { - // Skipping rec.stop() here: the worker is gone, so the command - // would only fail on the closed channel. let _ = rec.close(); } *self.is_recording.lock().unwrap() = false; *open_flag = false; + self.invalidate_device_cache(); + + // If the failed stream was the user's selected microphone, make + // the fallback explicit so the settings UI matches capture. + let mut settings = get_settings(&self.app_handle); + if settings.selected_microphone.is_some() + && self.desired_device_name(&settings) == settings.selected_microphone + { + settings.selected_microphone = None; + write_settings(&self.app_handle, settings); + let _ = self.app_handle.emit( + "settings-changed", + serde_json::json!({ + "setting": "selected_microphone", + "value": "Default" + }), + ); + } // Fall through and open a fresh stream. } diff --git a/src/stores/settingsStore.ts b/src/stores/settingsStore.ts index 5538908aed..bc110c87a0 100644 --- a/src/stores/settingsStore.ts +++ b/src/stores/settingsStore.ts @@ -610,6 +610,9 @@ export const useSettingsStore = create()( listen("model-state-changed", () => { get().refreshSettings(); }); + listen("settings-changed", () => { + get().refreshSettings(); + }); }, })), );