Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
96 changes: 96 additions & 0 deletions openless-all/app/crates/openless-core/src/dictation_engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1123,6 +1123,7 @@ struct BufferedTranscriptionInner {
partials: Arc<dyn TextStreamSink>,
progress: Arc<dyn RecordingProgressSink>,
limit_notified: AtomicBool,
has_nonzero_pcm: AtomicBool,
limit_threshold_bytes: usize,
state: Mutex<BufferedTranscriptionState>,
}
Expand Down Expand Up @@ -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())),
}),
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -1407,6 +1413,24 @@ impl TranscriptionSession for BufferedTranscriptionSession {
}

fn finish(&self) -> BoxFuture<'static, Result<crate::ports::TranscriptOutput, BackendError>> {
// 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?;
Expand Down Expand Up @@ -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()));
Expand Down
33 changes: 33 additions & 0 deletions openless-all/app/scripts/test-recorder-module.py
Original file line number Diff line number Diff line change
@@ -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), ensure_ascii=False) + "]\npub mod recorder;\n"
)
subprocess.run(["cargo", "test", "--manifest-path", str(root / "Cargo.toml"), *sys.argv[1:]], check=True)
Loading