From 7ee50185d43a6b8992e47ec5fdf23c03e6e0573a Mon Sep 17 00:00:00 2001 From: Cassie Date: Fri, 2 Oct 2026 11:42:32 +0800 Subject: [PATCH 1/2] =?UTF-8?q?fix(audio):=20=E4=BC=98=E5=85=88=E9=BA=A6?= =?UTF-8?q?=E5=85=8B=E9=A3=8E=E5=8E=9F=E7=94=9F=E6=A0=BC=E5=BC=8F=E5=B9=B6?= =?UTF-8?q?=E6=8B=A6=E6=88=AA=E7=A9=BA=E5=BD=95=E9=9F=B3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../openless-core/src/dictation_engine.rs | 96 ++++++ .../app/scripts/test-recorder-module.py | 33 ++ openless-all/app/src-tauri/src/recorder.rs | 305 +++++++++++++----- 3 files changed, 362 insertions(+), 72 deletions(-) create mode 100644 openless-all/app/scripts/test-recorder-module.py diff --git a/openless-all/app/crates/openless-core/src/dictation_engine.rs b/openless-all/app/crates/openless-core/src/dictation_engine.rs index b25758c27..7056f36f9 100644 --- a/openless-all/app/crates/openless-core/src/dictation_engine.rs +++ b/openless-all/app/crates/openless-core/src/dictation_engine.rs @@ -1123,6 +1123,7 @@ struct BufferedTranscriptionInner { partials: Arc, progress: Arc, limit_notified: AtomicBool, + has_nonzero_pcm: AtomicBool, limit_threshold_bytes: usize, state: Mutex, } @@ -1164,6 +1165,7 @@ impl BufferedTranscriptionSession { partials, progress, limit_notified: AtomicBool::new(false), + has_nonzero_pcm: AtomicBool::new(false), limit_threshold_bytes, state: Mutex::new(BufferedTranscriptionState::Buffering(Vec::new())), }), @@ -1353,6 +1355,10 @@ impl Drop for NotifyOnDrop { impl AudioConsumer for BufferedTranscriptionSession { fn consume_pcm_chunk(&self, pcm: &[u8]) { + if !self.inner.has_nonzero_pcm.load(Ordering::Relaxed) && pcm.iter().any(|byte| *byte != 0) + { + self.inner.has_nonzero_pcm.store(true, Ordering::Release); + } let (downstream, buffer_limit_reached) = { let mut state = self .inner @@ -1407,6 +1413,24 @@ impl TranscriptionSession for BufferedTranscriptionSession { } fn finish(&self) -> BoxFuture<'static, Result> { + // Inspect after the caller stops capture. Cancellation remains independent + // of audio validity, and quiet non-zero samples do not trip a speech threshold. + let terminal = matches!( + &*self + .inner + .state + .lock() + .expect("buffered transcription lock poisoned"), + BufferedTranscriptionState::Failed(_) | BufferedTranscriptionState::Cancelled + ); + if !terminal && !self.inner.has_nonzero_pcm.load(Ordering::Acquire) { + return Box::pin(async { + Err(BackendError::new( + BackendErrorCode::InvalidArgument, + "未收到有效音频,请检查麦克风输入设备后重试", + )) + }); + } let attaching = self.attach(); Box::pin(async move { let downstream = attaching.await?; @@ -2197,6 +2221,78 @@ mod tests { assert_eq!(&*fixture.pcm.lock().unwrap(), &[1, 0, 2, 0]); } + #[tokio::test] + async fn empty_or_zero_pcm_never_finalizes_asr_and_can_still_be_cancelled() { + for (chunk, attached) in [ + (Vec::new(), false), + (vec![0; 640], false), + (Vec::new(), true), + (vec![0; 640], true), + ] { + let pcm = Arc::new(Mutex::new(Vec::new())); + let starts = Arc::new(AtomicUsize::new(0)); + let cancels = Arc::new(AtomicUsize::new(0)); + let transcriber = Arc::new(FixtureTranscriber { + session: Arc::new(FixtureTranscriptionSession { + pcm: pcm.clone(), + cancels: cancels.clone(), + finish_entered: None, + finish_release: None, + }), + starts: starts.clone(), + }); + let prepared = transcriber + .prepare(SessionId::new(), raw_dictation_context()) + .await + .unwrap(); + let buffered = BufferedTranscriptionSession::new( + prepared, + Arc::new(DiscardTextStream), + Arc::new(LimitRecordingProgress::default()), + ); + if attached { + buffered.attach().await.unwrap(); + } + buffered.consume_pcm_chunk(&chunk); + let error = buffered.finish().await.unwrap_err(); + assert_eq!(error.code, BackendErrorCode::InvalidArgument); + assert_eq!(starts.load(Ordering::Acquire), usize::from(attached)); + buffered.cancel().await.unwrap(); + assert_eq!(cancels.load(Ordering::Acquire), usize::from(attached)); + assert_eq!( + buffered.finish().await.unwrap_err().code, + BackendErrorCode::Cancelled + ); + } + } + + #[tokio::test] + async fn quiet_nonzero_pcm_is_not_rejected_as_silence() { + let pcm = Arc::new(Mutex::new(Vec::new())); + let starts = Arc::new(AtomicUsize::new(0)); + let transcriber = Arc::new(FixtureTranscriber { + session: Arc::new(FixtureTranscriptionSession { + pcm: pcm.clone(), + cancels: Arc::new(AtomicUsize::new(0)), + finish_entered: None, + finish_release: None, + }), + starts: starts.clone(), + }); + let prepared = transcriber + .prepare(SessionId::new(), raw_dictation_context()) + .await + .unwrap(); + let buffered = BufferedTranscriptionSession::new( + prepared, + Arc::new(DiscardTextStream), + Arc::new(LimitRecordingProgress::default()), + ); + buffered.consume_pcm_chunk(&[1, 0]); + assert_eq!(buffered.finish().await.unwrap().text, "raw text"); + assert_eq!(starts.load(Ordering::Acquire), 1); + } + #[tokio::test] async fn buffered_pcm_limit_requests_one_stop_and_preserves_the_tail() { let pcm = Arc::new(Mutex::new(Vec::new())); diff --git a/openless-all/app/scripts/test-recorder-module.py b/openless-all/app/scripts/test-recorder-module.py new file mode 100644 index 000000000..8932e43cc --- /dev/null +++ b/openless-all/app/scripts/test-recorder-module.py @@ -0,0 +1,33 @@ +#!/usr/bin/env python3 +"""Run the native recorder tests without building Tauri/MLX or fetching ASR submodules. + +The bridge trait below matches the only core interface used by recorder.rs. +This checks the production recorder module; it does not replace a desktop cargo check. +""" +import json +from pathlib import Path +import subprocess +import sys +import tempfile + +recorder = Path(__file__).resolve().parents[1] / "src-tauri/src/recorder.rs" +with tempfile.TemporaryDirectory(prefix="openless-recorder-tests-") as directory: + root = Path(directory) + (root / "src").mkdir() + (root / "Cargo.toml").write_text("""[package] +name = "openless-recorder-tests" +version = "0.1.0" +edition = "2021" +[dependencies] +cpal = "=0.15.3" +parking_lot = "0.12" +serde = { version = "1", features = ["derive"] } +thiserror = "1" +log = "0.4" +""") + (root / "src/lib.rs").write_text( + "extern crate self as openless_core;\n" + "pub trait AudioConsumer: Send + Sync { fn consume_pcm_chunk(&self, pcm: &[u8]); }\n" + + "#[path = " + json.dumps(str(recorder)) + "]\npub mod recorder;\n" + ) + subprocess.run(["cargo", "test", "--manifest-path", str(root / "Cargo.toml"), *sys.argv[1:]], check=True) diff --git a/openless-all/app/src-tauri/src/recorder.rs b/openless-all/app/src-tauri/src/recorder.rs index 401f04b62..0a81ce309 100644 --- a/openless-all/app/src-tauri/src/recorder.rs +++ b/openless-all/app/src-tauri/src/recorder.rs @@ -241,6 +241,7 @@ impl Recorder { } }; if let Err(error) = startup_result { + let _ = join_handle.join(); if let Some(archive) = archive_writer.as_ref() { archive.finish(); } @@ -332,15 +333,7 @@ fn run_audio_thread( } }; - if let Err(err) = stream.play() { - let _ = startup_tx.send(Err(RecorderError::EngineFailed(format!("play: {err}")))); - return; - } - - // Startup succeeded. - let _ = startup_tx.send(Ok(())); - - // Startup succeeded. + // build_input_stream returns only after native playback starts successfully. let _ = startup_tx.send(Ok(())); // Start the liveness watchdog: detect the capture callback silently stopping. @@ -467,7 +460,7 @@ fn run_audio_thread( } } -/// Select default input device + default config + build the Stream. +/// Try the selected microphone's native format before other advertised formats/devices. fn build_input_stream( microphone_device_name: Option, consumer: Arc, @@ -476,78 +469,162 @@ fn build_input_stream( runtime_error_tx: Sender, ) -> Result<(cpal::Stream, Arc), RecorderError> { let host = cpal::default_host(); - let device = select_input_device(&host, microphone_device_name.as_deref())?; - - let supported = device - .default_input_config() - .map_err(|e| classify_default_config_err(e.to_string()))?; - - let sample_format = supported.sample_format(); - let default_config: StreamConfig = supported.config(); - let config = stable_input_config_for_platform(&default_config); - let input_sr = config.sample_rate.0; - let channels = config.channels as usize; + let selected = select_input_device(&host, microphone_device_name.as_deref())?; + let selected_name = selected.name().ok(); + let start = |device: &cpal::Device| { + start_device_stream( + device, + &consumer, + &level_handler, + archiver.clone(), + &runtime_error_tx, + ) + }; + let initial_error = match start(&selected) { + Ok(stream) => return Ok(stream), + Err(RecorderError::PermissionDenied) => return Err(RecorderError::PermissionDenied), + Err(error) => error, + }; + log::warn!("[recorder] selected microphone failed: {initial_error}; trying other inputs"); + let default = host.default_input_device(); + let default_name = default.as_ref().and_then(|device| device.name().ok()); + let mut candidates = Vec::new(); + if let Some(default) = default { + if default_name != selected_name { + candidates.push(default); + } + } + if let Ok(devices) = host.input_devices() { + candidates.extend(devices.filter(|device| { + let name = device.name().ok(); + name.is_none() || (name != selected_name && name != default_name) + })); + } + try_input_candidates(candidates, initial_error, |device| start(&device)) +} - log::info!( - "[recorder] inputDevice={} inputFormat sampleRate={} channels={} fmt={:?}", - device.name().unwrap_or_else(|_| "".into()), - input_sr, - channels, - sample_format - ); +fn try_input_candidates( + candidates: impl IntoIterator, + mut last_error: RecorderError, + mut start: impl FnMut(T) -> Result, +) -> Result { + for candidate in candidates { + match start(candidate) { + Ok(stream) => return Ok(stream), + Err(RecorderError::PermissionDenied) => return Err(RecorderError::PermissionDenied), + Err(error) => last_error = error, + } + } + Err(last_error) +} - let state = Arc::new(StreamState::new()); - let stream = match build_stream_for_format( - &device, - &config, - sample_format, - Arc::clone(&consumer), - Arc::clone(&level_handler), - archiver.clone(), - Arc::clone(&state), - input_sr, - channels, - runtime_error_tx.clone(), - ) { - Ok(stream) => stream, - Err(err) if config != default_config => { - log::warn!( - "[recorder] stable input config failed; falling back to default config: {err}" - ); - build_stream_for_format( - &device, - &default_config, - sample_format, +fn start_device_stream( + device: &cpal::Device, + consumer: &Arc, + level_handler: &Arc, + archiver: Option>, + runtime_error_tx: &Sender, +) -> Result<(cpal::Stream, Arc), RecorderError> { + let mut default = None; + let initial_error = match device.default_input_config() { + Ok(config) => { + default = Some(config.clone()); + match start_native_config( + device, + config, consumer, level_handler, - archiver, - Arc::clone(&state), - default_config.sample_rate.0, - default_config.channels as usize, + archiver.clone(), runtime_error_tx, - )? + ) { + Ok(stream) => return Ok(stream), + Err(RecorderError::PermissionDenied) => { + return Err(RecorderError::PermissionDenied) + } + Err(error) => error, + } } - Err(err) => return Err(err), + Err(error) => classify_default_config_err(error.to_string()), }; - Ok((stream, state)) + if matches!(initial_error, RecorderError::PermissionDenied) { + return Err(initial_error); + } + // Query alternatives only after failure; never force mono or an ASR rate on hardware. + let ranges = match device.supported_input_configs() { + Ok(ranges) => ranges, + Err(_) => return Err(initial_error), + }; + let configs = native_input_alternatives(default.as_ref(), ranges); + try_input_candidates(configs, initial_error, |config| { + start_native_config( + device, + config, + consumer, + level_handler, + archiver.clone(), + runtime_error_tx, + ) + }) } -#[cfg(target_os = "android")] -fn stable_input_config_for_platform(default_config: &StreamConfig) -> StreamConfig { - let mut config = default_config.clone(); - if config.channels > 1 { - log::info!( - "[recorder] android forcing mono input channels: {} -> 1", - config.channels - ); - config.channels = 1; +fn native_input_alternatives( + default: Option<&cpal::SupportedStreamConfig>, + ranges: impl IntoIterator, +) -> Vec { + let native_rate = default.map(|config| config.sample_rate().0); + let mut candidates = Vec::new(); + for range in ranges { + let min = range.min_sample_rate().0; + let max = range.max_sample_rate().0; + for rate in [native_rate.unwrap_or(max).clamp(min, max), min, max] { + let candidate = range.with_sample_rate(cpal::SampleRate(rate)); + let same = |other: &cpal::SupportedStreamConfig| { + other.config() == candidate.config() + && other.sample_format() == candidate.sample_format() + }; + if !default.is_some_and(same) && !candidates.iter().any(same) { + candidates.push(candidate); + } + } } - config + candidates } -#[cfg(not(target_os = "android"))] -fn stable_input_config_for_platform(default_config: &StreamConfig) -> StreamConfig { - default_config.clone() +fn start_native_config( + device: &cpal::Device, + supported: cpal::SupportedStreamConfig, + consumer: &Arc, + level_handler: &Arc, + archiver: Option>, + runtime_error_tx: &Sender, +) -> Result<(cpal::Stream, Arc), RecorderError> { + let config = supported.config(); + let state = Arc::new(StreamState::new()); + state.accepting_audio.store(false, Ordering::Release); + let stream = build_stream_for_format( + device, + &config, + supported.sample_format(), + Arc::clone(consumer), + Arc::clone(level_handler), + archiver, + Arc::clone(&state), + config.sample_rate.0, + config.channels as usize, + runtime_error_tx.clone(), + )?; + stream + .play() + .map_err(|error| classify_default_config_err(format!("play: {error}")))?; + state.accepting_audio.store(true, Ordering::Release); + log::info!( + "[recorder] inputDevice={} inputFormat sampleRate={} channels={} fmt={:?}", + device.name().unwrap_or_else(|_| "".into()), + config.sample_rate.0, + config.channels, + supported.sample_format() + ); + Ok((stream, state)) } fn select_input_device( @@ -571,7 +648,12 @@ fn select_input_device( ); } - host.default_input_device() + if let Some(device) = host.default_input_device() { + return Ok(device); + } + host.input_devices() + .map_err(|error| RecorderError::EngineFailed(format!("input_devices: {error}")))? + .next() .ok_or(RecorderError::NoInputDevice) } @@ -645,7 +727,11 @@ fn build_stream_for_format( let archiver = archiver.clone(); let state = Arc::clone(&state); let runtime_error_tx = runtime_error_tx.clone(); + let error_state = Arc::clone(&state); let err_cb = move |err| { + if !error_state.accepting_audio.load(Ordering::Acquire) { + return; + } log::error!("[recorder] stream error: {err}"); let _ = runtime_error_tx.send(RecorderError::EngineFailed(format!("stream: {err}"))); @@ -677,6 +763,12 @@ fn build_stream_for_format( match sample_format { SampleFormat::F32 => make_stream!(f32, |s: f32| s), + SampleFormat::F64 => make_stream!(f64, |s: f64| s as f32), + SampleFormat::I64 => make_stream!(i64, |s: i64| (s as f64 / i64::MAX as f64) as f32), + SampleFormat::U64 => make_stream!(u64, |s: u64| ((s as f64 / u64::MAX as f64) * 2.0 - 1.0) + as f32), + SampleFormat::U32 => make_stream!(u32, |s: u32| (s as f64 / u32::MAX as f64 * 2.0 - 1.0) + as f32), SampleFormat::I16 => make_stream!(i16, |s: i16| s as f32 / i16::MAX as f32), SampleFormat::U16 => { make_stream!(u16, |s: u16| (s as f32 - 32768.0) / 32768.0) @@ -703,6 +795,7 @@ struct StreamState { /// next callback. last_sample: Mutex, callback_count: AtomicUsize, + accepting_audio: AtomicBool, peak_input_rms_milli: AtomicUsize, peak_output_rms_milli: AtomicUsize, /// Timestamp of the last successful consumer call (for liveness detection) @@ -715,6 +808,7 @@ impl StreamState { resample_phase: Mutex::new(0.0), last_sample: Mutex::new(0.0), callback_count: AtomicUsize::new(0), + accepting_audio: AtomicBool::new(true), peak_input_rms_milli: AtomicUsize::new(0), peak_output_rms_milli: AtomicUsize::new(0), // Start as None: timing begins only after the first callback, avoiding false @@ -734,7 +828,7 @@ fn process_callback( archiver: Option<&WavArchiveWriter>, state: &StreamState, ) { - if interleaved.is_empty() || channels == 0 { + if !state.accepting_audio.load(Ordering::Acquire) || interleaved.is_empty() || channels == 0 { return; } @@ -1006,6 +1100,73 @@ mod tests { .collect() } + #[test] + fn native_alternatives_preserve_channels_format_and_advertised_rates() { + let range = cpal::SupportedStreamConfigRange::new( + 2, + cpal::SampleRate(44_100), + cpal::SampleRate(48_000), + cpal::SupportedBufferSize::Unknown, + SampleFormat::F32, + ); + let native = range.with_sample_rate(cpal::SampleRate(48_000)); + let alternatives = native_input_alternatives(Some(&native), [range, range]); + assert_eq!(alternatives.len(), 1); + assert_eq!(alternatives[0].channels(), 2); + assert_eq!(alternatives[0].sample_format(), SampleFormat::F32); + assert_eq!(alternatives[0].sample_rate().0, 44_100); + } + + #[test] + fn unstarted_candidate_cannot_feed_asr_or_mark_audio_valid() { + let state = StreamState::new(); + state.accepting_audio.store(false, Ordering::Release); + let consumer = RecordingConsumer::default(); + process_callback( + &[0.5; 320], + 1, + TARGET_SAMPLE_RATE, + &consumer, + &|_| {}, + None, + &state, + ); + assert!(consumer.chunks.lock().unwrap().is_empty()); + assert!(state.last_callback_time.lock().is_none()); + } + + #[test] + fn startup_uses_candidates_in_order_and_stops_after_success() { + let mut attempted = Vec::new(); + let result = try_input_candidates( + ["native", "alternative", "other device"], + RecorderError::NoInputDevice, + |candidate| { + attempted.push(candidate); + if candidate == "alternative" { + Ok(candidate) + } else { + Err(RecorderError::EngineFailed(candidate.into())) + } + }, + ) + .unwrap(); + assert_eq!(result, "alternative"); + assert_eq!(attempted, ["native", "alternative"]); + } + + #[test] + fn permission_failure_never_probes_more_devices() { + let mut attempts = 0; + let result: Result<(), _> = + try_input_candidates([1, 2], RecorderError::NoInputDevice, |_| { + attempts += 1; + Err(RecorderError::PermissionDenied) + }); + assert!(matches!(result, Err(RecorderError::PermissionDenied))); + assert_eq!(attempts, 1); + } + #[test] fn downmix_to_mono_averages_complete_interleaved_frames() { let mono = downmix_to_mono(&[1.0, -1.0, 0.5, 0.25, 0.0], 2); From 8ada9e01cce3c7d09e6b271fed225b330ad9221a Mon Sep 17 00:00:00 2001 From: Cassie Date: Fri, 2 Oct 2026 11:49:36 +0800 Subject: [PATCH 2/2] =?UTF-8?q?test(audio):=20=E6=94=AF=E6=8C=81=20Unicode?= =?UTF-8?q?=20=E8=B7=AF=E5=BE=84=E7=9A=84=E5=BD=95=E9=9F=B3=E9=AA=8C?= =?UTF-8?q?=E8=AF=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- openless-all/app/scripts/test-recorder-module.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/openless-all/app/scripts/test-recorder-module.py b/openless-all/app/scripts/test-recorder-module.py index 8932e43cc..5903df059 100644 --- a/openless-all/app/scripts/test-recorder-module.py +++ b/openless-all/app/scripts/test-recorder-module.py @@ -28,6 +28,6 @@ (root / "src/lib.rs").write_text( "extern crate self as openless_core;\n" "pub trait AudioConsumer: Send + Sync { fn consume_pcm_chunk(&self, pcm: &[u8]); }\n" - + "#[path = " + json.dumps(str(recorder)) + "]\npub mod recorder;\n" + + "#[path = " + json.dumps(str(recorder), ensure_ascii=False) + "]\npub mod recorder;\n" ) subprocess.run(["cargo", "test", "--manifest-path", str(root / "Cargo.toml"), *sys.argv[1:]], check=True)