From 3b71063909bce838d25b822ae6de1c55e4d3cfbd Mon Sep 17 00:00:00 2001 From: zhou-zhichao Date: Thu, 10 Sep 2026 02:02:07 +0200 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E8=BF=9C=E7=A8=8B=E5=BD=95?= =?UTF-8?q?=E9=9F=B3=E6=81=AF=E5=B1=8F=E4=B8=A2=E5=A4=B1=E5=B9=B6=E9=BB=98?= =?UTF-8?q?=E8=AE=A4=E4=BF=9D=E6=8C=81=E4=BA=AE=E5=B1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/remote-input-recording-recovery.md | 41 +++ .../app/crates/openless-core/src/api.rs | 52 ++++ .../app/crates/openless-core/src/domains.rs | 39 +++ .../openless-core/src/external_audio.rs | 191 +++++++++++- .../src/external_audio/archive.rs | 166 +++++++++++ .../app/crates/openless-core/src/lib.rs | 5 +- .../openless-core/src/remote_input_service.rs | 148 +++++++++- .../tests/remote_input_contract.rs | 226 +++++++++++++++ openless-all/app/linux-egui/src/backend.rs | 6 + .../app/linux-egui/src/remote_input.rs | 96 +++++- .../scripts/remote-input-audio-queue.test.mjs | 207 ++++++++++++- .../app/src-tauri/src/core_adapters.rs | 32 +- .../src-tauri/src/remote_server/assets/app.js | 274 +++++++++++++++++- .../src/remote_server/assets/index.html | 7 + .../src/remote_server/assets/style.css | 12 + .../app/src-tauri/src/remote_server/mod.rs | 79 ++++- 16 files changed, 1522 insertions(+), 59 deletions(-) create mode 100644 docs/remote-input-recording-recovery.md create mode 100644 openless-all/app/crates/openless-core/src/external_audio/archive.rs diff --git a/docs/remote-input-recording-recovery.md b/docs/remote-input-recording-recovery.md new file mode 100644 index 000000000..39163b092 --- /dev/null +++ b/docs/remote-input-recording-recovery.md @@ -0,0 +1,41 @@ +# 远程输入:息屏与录音恢复 + +手机录音页新增「录音时保持亮屏」,默认开启,并在当前浏览器记住选择。只在前台录音期间申请屏幕唤醒锁,结束、取消、断线或离开页面时释放。浏览器不支持,或系统因省电等原因拒绝时,页面会提示,录音仍可使用。 + +```text + [ 点击开始 / 点击结束 ] + 录音状态 + +录音时保持亮屏 [ 开 ] +息屏会结束本段录音,电脑继续处理已收到的部分。 + +电脑落字 [ 开 ] +``` + +## 息屏后会发生什么 + +- 自动息屏、手动锁屏、切到后台或麦克风被系统中断时,手机尽可能发送结束录音,电脑继续识别已经收到的音频。 +- 如果手机来不及通知,电脑在断线时收尾;连接仍在但连续 15 秒没有有效音频时,也会结束采集并开始识别。 +- 识别已经开始后,手机断线不会取消电脑上的识别与历史记录写入。主动取消、关闭远程输入、重置 PIN 仍会撤销会话。 +- 回到页面后重新认证并查询上次会话,找回结果或提示到电脑历史中重试。不会自动重新打开麦克风。 + +这不能让浏览器在锁屏后持续采集,也不能补回尚未传到电脑的音频。电脑退出、崩溃或磁盘不可写时,不能保证自动完成识别。 + +## 录音与结果保留 + +已收到的远程音频使用现有录音目录和保留设置写入 WAV。识别失败时,音频可供现有自动重试及桌面历史重新转录使用;识别成功后,除非开启保留成功录音,否则按现有策略删除音频。文字仍按历史记录设置保留。 + +手机只保存上次会话 ID 和单独生成的随机恢复凭据。恢复查询要求有效配对,不提供历史列表;仅凭公开会话 ID 无法读取结果。服务内最多保留最近 64 个会话的恢复权限,最长 24 小时,取消会话、重启服务或重置 PIN 后失效。此后请在电脑历史记录中查看。 + +## 真机验收(待补) + +记录手机型号、系统与浏览器版本、电脑平台、识别模式及测试 commit;分别检查 iPhone Safari 和 Android Chrome。 + +- 开关默认开启;关闭、刷新后仍关闭;录音中切换立即生效,结束后恢复系统自动息屏。 +- 点按和按住模式各录音 1–2 分钟:允许自动息屏、手动锁屏、切后台后,电脑保留已收到的部分并完成识别;回到手机能看到结果。 +- 识别过程中断开网络再重连;检查电脑历史以及手机找回的文字。 +- 模拟识别失败,检查 WAV 保留与历史重新转录;检查保留天数和录音数量上限。 +- 低电量或浏览器拒绝保持亮屏时,提示明确且能开始录音。 +- 识别中主动取消或重置 PIN,确认没有迟到落字,旧恢复凭据失效。 + +自动化测试覆盖模拟浏览器事件、两分钟 PCM、归档与失败重试、恢复鉴权和权限撤销;这些检查不能代替手机锁屏与省电行为的真机验收。 diff --git a/openless-all/app/crates/openless-core/src/api.rs b/openless-all/app/crates/openless-core/src/api.rs index 953d191a0..252de18c3 100644 --- a/openless-all/app/crates/openless-core/src/api.rs +++ b/openless-all/app/crates/openless-core/src/api.rs @@ -9346,6 +9346,58 @@ mod tests { let _ = std::fs::remove_dir_all(data_dir); } + #[tokio::test] + async fn external_audio_saves_failed_recordings_for_history_retry_and_prunes_successful_audio() + { + for fail in [false, true] { + let data_dir = std::env::temp_dir() + .join(format!("openless-remote-history-{}", uuid::Uuid::new_v4())); + let transcription = if fail { + crate::testing::FixtureTranscriptionEngine::failing(BackendError::new( + BackendErrorCode::Provider, + "fixture ASR failure", + )) + } else { + crate::testing::FixtureTranscriptionEngine::successful("received speech", 120_000) + }; + let engine = crate::PipelineDictationEngine::new( + Arc::new(crate::ExternalAudioRecorder::with_recordings_directory( + data_dir.join("recordings"), + )), + Arc::new(transcription.clone()), + Arc::new(crate::testing::FixtureTextPolisher::successful( + "complete transcription", + )), + ); + let backend = backend_with_dictation_engine(data_dir.clone(), Arc::new(engine)); + backend.start().await.unwrap(); + let session = backend.start_external_dictation().await.unwrap(); + let path = data_dir.join("recordings").join(format!("{session}.wav")); + for second in 0..120 { + backend + .feed_external_pcm(session, &vec![second; 32_000]) + .unwrap(); + } + let expected = transcription.pcm(); + assert_eq!(&std::fs::read(&path).unwrap()[44..], expected); + let result = backend.stop_dictation_session(session).await; + assert_eq!(result.is_err(), fail); + let history = backend.list_history().unwrap(); + assert_eq!(history.len(), 1); + assert_eq!(history[0].id, session.to_string()); + assert_eq!(history[0].has_audio_recording, Some(fail)); + if fail { + assert_eq!(history[0].error_code.as_deref(), Some("transcribeFailed")); + assert_eq!(&std::fs::read(path).unwrap()[44..], expected); + } else { + assert_eq!(history[0].final_text, "complete transcription"); + assert!(!path.exists()); + } + backend.shutdown().await.unwrap(); + let _ = std::fs::remove_dir_all(data_dir); + } + } + #[tokio::test] async fn dictation_freezes_channel_identity_protocol_and_model_for_the_session() { let data_dir = std::env::temp_dir().join(format!( diff --git a/openless-all/app/crates/openless-core/src/domains.rs b/openless-all/app/crates/openless-core/src/domains.rs index 1aefdd1e8..8423e13e8 100644 --- a/openless-all/app/crates/openless-core/src/domains.rs +++ b/openless-all/app/crates/openless-core/src/domains.rs @@ -1059,6 +1059,27 @@ pub trait RemoteInputRuntimeAdapter: Send + Sync { &self, session_id: SessionId, ) -> BoxFuture<'static, Result<(), BackendError>>; + /// 只读取指定会话;由 Core 校验手机持有的恢复凭据。 + fn read_audio_history( + &self, + _session_id: SessionId, + ) -> BoxFuture<'static, Result, BackendError>> { + Box::pin(async { Ok(None) }) + } +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "kind", rename_all = "camelCase")] +pub enum RemoteInputRecovery { + Pending, + Completed { + text: String, + }, + Failed { + #[serde(rename = "hasAudioRecording")] + has_audio_recording: bool, + }, + Unavailable, } pub trait RemoteInputApi: Send + Sync { @@ -1094,6 +1115,24 @@ pub trait RemoteInputApi: Send + Sync { ) -> BoxFuture<'static, Result<(), BackendError>> { unsupported("remote input") } + fn recover_stream( + &self, + _connection_id: SessionId, + _session_id: SessionId, + _recovery_key: crate::credentials::SecretValue, + ) -> BoxFuture<'static, Result> { + unsupported("remote input") + } + fn recovery_key( + &self, + _connection_id: SessionId, + _session_id: SessionId, + ) -> Result { + Err(BackendError::new( + BackendErrorCode::Unsupported, + "remote recovery is unavailable", + )) + } fn start_stream( &self, _connection_id: SessionId, diff --git a/openless-all/app/crates/openless-core/src/external_audio.rs b/openless-all/app/crates/openless-core/src/external_audio.rs index 793791fb5..24d0cbf8d 100644 --- a/openless-all/app/crates/openless-core/src/external_audio.rs +++ b/openless-all/app/crates/openless-core/src/external_audio.rs @@ -5,22 +5,30 @@ //! little-endian PCM contract and routes it to the active pipeline session. use std::collections::HashMap; +use std::path::PathBuf; use std::sync::{Arc, Mutex}; use futures_util::future::BoxFuture; use crate::dictation_context::{DictationAudioSource, DictationContext}; use crate::errors::{BackendError, BackendErrorCode}; -use crate::ports::{ActiveRecording, AudioConsumer, AudioRecorder, RecordingProgressSink}; +use crate::ports::{ + ActiveRecording, AudioConsumer, AudioRecorder, RecordingArchive, RecordingProgressSink, +}; use crate::types::SessionId; +mod archive; +use archive::ExternalRecordingArchive; + #[derive(Clone, Default)] pub struct ExternalAudioRecorder { sessions: Arc>>>, + recordings_dir: Option, } struct ExternalRecordingSession { state: Mutex, + archive: Option>, } struct ExternalRecordingState { @@ -37,6 +45,14 @@ struct ExternalActiveRecording { } impl ExternalAudioRecorder { + /// 与电脑历史记录共用 WAV 目录及保留策略,失败的远程录音也可重新转录。 + pub fn with_recordings_directory(directory: PathBuf) -> Self { + Self { + recordings_dir: Some(directory), + ..Self::default() + } + } + fn release( &self, session_id: SessionId, @@ -60,6 +76,9 @@ impl ExternalAudioRecorder { .lock() .expect("external audio state lock poisoned") .active = false; + if let Some(archive) = &expected.archive { + archive.finish(); + } sessions.remove(&session_id); Ok(()) } @@ -81,6 +100,29 @@ impl AudioRecorder for ExternalAudioRecorder { )) }); } + // 在归档创建/清理之前检查重复会话,避免误删仍在录音的文件。 + let mut sessions = self + .sessions + .lock() + .expect("external audio session lock poisoned"); + if sessions.contains_key(&session_id) { + return Box::pin(async { + Err(BackendError::new( + BackendErrorCode::Busy, + "external audio session already exists", + )) + }); + } + let archive = self + .recordings_dir + .as_ref() + .filter(|_| context.recording.archive_enabled) + .and_then(|directory| { + ExternalRecordingArchive::create(directory, session_id, &context.recording) + .map(Arc::new) + .map_err(|error| log::warn!("[remote-input] 录音归档不可用:{error}")) + .ok() + }); let session = Arc::new(ExternalRecordingSession { state: Mutex::new(ExternalRecordingState { active: true, @@ -88,22 +130,10 @@ impl AudioRecorder for ExternalAudioRecorder { consumer, progress, }), + archive, }); - { - let mut sessions = self - .sessions - .lock() - .expect("external audio session lock poisoned"); - if sessions.contains_key(&session_id) { - return Box::pin(async { - Err(BackendError::new( - BackendErrorCode::Busy, - "external audio session already exists", - )) - }); - } - sessions.insert(session_id, Arc::clone(&session)); - } + sessions.insert(session_id, Arc::clone(&session)); + drop(sessions); let recording = ExternalActiveRecording { recorder: self.clone(), session_id, @@ -141,6 +171,9 @@ impl AudioRecorder for ExternalAudioRecorder { "external audio session is not active", )); } + if let Some(archive) = &session.archive { + archive.append(pcm); + } state.consumer.consume_pcm_chunk(pcm); state.bytes_received = state.bytes_received.saturating_add(pcm.len() as u64); let elapsed_ms = state.bytes_received.saturating_mul(1_000) / 32_000; @@ -151,6 +184,13 @@ impl AudioRecorder for ExternalAudioRecorder { } impl ActiveRecording for ExternalActiveRecording { + fn archive(&self) -> Option> { + self.session + .archive + .as_ref() + .map(|archive| Arc::clone(archive) as Arc) + } + fn stop(self: Box) -> BoxFuture<'static, Result<(), BackendError>> { Box::pin(async move { self.recorder.release(self.session_id, &self.session) }) } @@ -253,6 +293,125 @@ mod tests { } } + #[tokio::test] + async fn remote_archive_keeps_two_minutes_of_pcm_and_survives_recording_stop() { + let directory = + std::env::temp_dir().join(format!("openless-remote-archive-{}", uuid::Uuid::new_v4())); + let recorder = ExternalAudioRecorder::with_recordings_directory(directory.clone()); + let id = SessionId::new(); + let mut context = DictationContext::default(); + context.audio_source = DictationAudioSource::External; + context.recording.archive_enabled = true; + let consumer = Arc::new(RecordingConsumer::default()); + let recording = recorder + .start( + id, + Arc::new(context.clone()), + consumer.clone(), + Arc::new(RecordingProgress::default()), + ) + .await + .unwrap(); + let archive = recording.archive().unwrap(); + assert!(!archive.is_available()); + for second in 0..120 { + recorder.feed_pcm(id, &vec![second; 32_000]).unwrap(); + } + let expected = consumer.0.lock().unwrap().clone(); + let path = directory.join(format!("{id}.wav")); + let wav = std::fs::read(&path).unwrap(); + assert_eq!(&wav[44..], expected); + assert_eq!( + u32::from_le_bytes(wav[40..44].try_into().unwrap()) as usize, + 120 * 32_000 + ); + assert!(recorder + .start( + id, + Arc::new(context), + consumer, + Arc::new(RecordingProgress::default()) + ) + .await + .is_err()); + assert_eq!(std::fs::read(&path).unwrap(), wav); + recording.stop().await.unwrap(); + assert!(archive.is_available()); + assert_eq!(archive.read_pcm().await.unwrap(), expected); + archive.discard().await.unwrap(); + assert!(!archive.is_available()); + assert!(!path.exists()); + let _ = std::fs::remove_dir_all(directory); + } + + #[tokio::test] + async fn disabled_or_unavailable_archive_does_not_block_remote_capture() { + let directory = std::env::temp_dir().join(format!( + "openless-remote-no-archive-{}", + uuid::Uuid::new_v4() + )); + for enabled in [false, true] { + if enabled { + std::fs::write(&directory, b"not a directory").unwrap(); + } + let recorder = ExternalAudioRecorder::with_recordings_directory(directory.clone()); + let id = SessionId::new(); + let mut context = DictationContext::default(); + context.audio_source = DictationAudioSource::External; + context.recording.archive_enabled = enabled; + let consumer = Arc::new(RecordingConsumer::default()); + let recording = recorder + .start( + id, + Arc::new(context), + consumer.clone(), + Arc::new(RecordingProgress::default()), + ) + .await + .unwrap(); + assert!(recording.archive().is_none()); + recorder.feed_pcm(id, &[1, 0, 2, 0]).unwrap(); + assert_eq!(*consumer.0.lock().unwrap(), vec![1, 0, 2, 0]); + recording.stop().await.unwrap(); + if !enabled { + assert!(!directory.exists()); + } + } + let _ = std::fs::remove_file(directory); + } + + #[tokio::test] + async fn remote_archive_applies_the_existing_count_limit_without_deleting_other_files() { + let directory = std::env::temp_dir().join(format!( + "openless-remote-retention-{}", + uuid::Uuid::new_v4() + )); + std::fs::create_dir_all(&directory).unwrap(); + std::fs::write(directory.join("user.wav"), b"keep").unwrap(); + let recorder = ExternalAudioRecorder::with_recordings_directory(directory.clone()); + let mut context = DictationContext::default(); + context.audio_source = DictationAudioSource::External; + context.recording.archive_enabled = true; + context.recording.max_entries = Some(2); + for _ in 0..4 { + let id = SessionId::new(); + let recording = recorder + .start( + id, + Arc::new(context.clone()), + Arc::new(RecordingConsumer::default()), + Arc::new(RecordingProgress::default()), + ) + .await + .unwrap(); + recorder.feed_pcm(id, &[1, 0]).unwrap(); + recording.stop().await.unwrap(); + } + assert_eq!(std::fs::read_dir(&directory).unwrap().count(), 3); + assert_eq!(std::fs::read(directory.join("user.wav")).unwrap(), b"keep"); + let _ = std::fs::remove_dir_all(directory); + } + #[tokio::test] async fn external_sessions_route_pcm_progress_and_reject_late_or_wrong_frames() { let microphone_starts = Arc::new(AtomicUsize::new(0)); diff --git a/openless-all/app/crates/openless-core/src/external_audio/archive.rs b/openless-all/app/crates/openless-core/src/external_audio/archive.rs new file mode 100644 index 000000000..14a786195 --- /dev/null +++ b/openless-all/app/crates/openless-core/src/external_audio/archive.rs @@ -0,0 +1,166 @@ +use std::io::{Seek, SeekFrom, Write}; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Mutex; + +use futures_util::future::BoxFuture; + +use crate::{BackendError, BackendErrorCode, RecordingArchive, RecordingPlan, SessionId}; + +pub(super) struct ExternalRecordingArchive { + path: PathBuf, + writer: Mutex<(Option, u32)>, + available: AtomicBool, +} + +impl ExternalRecordingArchive { + pub(super) fn create( + directory: &Path, + session_id: SessionId, + plan: &RecordingPlan, + ) -> std::io::Result { + std::fs::create_dir_all(directory)?; + // 只清理本应用生成的 UUID.wav,且为本次录音预留一个名额。 + let mut entries = Vec::new(); + for entry in std::fs::read_dir(directory)?.flatten() { + let path = entry.path(); + if path.extension().and_then(|ext| ext.to_str()) != Some("wav") + || path + .file_stem() + .and_then(|stem| stem.to_str()) + .and_then(|stem| uuid::Uuid::parse_str(stem).ok()) + .is_none() + { + continue; + } + let Ok(metadata) = entry.metadata() else { + continue; + }; + if !metadata.is_file() { + continue; + } + let modified = metadata.modified().unwrap_or(std::time::UNIX_EPOCH); + if plan.retention_days > 0 + && modified + .elapsed() + .is_ok_and(|age| age.as_secs() > u64::from(plan.retention_days) * 86_400) + { + let _ = std::fs::remove_file(path); + } else { + entries.push((path, modified)); + } + } + entries.sort_by(|left, right| right.1.cmp(&left.1)); + let cap = plan + .max_entries + .map(|count| (count as usize).clamp(1, crate::history::HISTORY_CAP)) + .unwrap_or(crate::history::HISTORY_CAP); + for (path, _) in entries.into_iter().skip(cap - 1) { + let _ = std::fs::remove_file(path); + } + let path = directory.join(format!("{session_id}.wav")); + let mut options = std::fs::OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + let mut file = options.open(&path)?; + file.write_all(&crate::audio::encode_dictation_wav(&[]).expect("empty PCM is valid"))?; + Ok(Self { + path, + writer: Mutex::new((Some(file), 0)), + available: AtomicBool::new(false), + }) + } + + pub(super) fn append(&self, pcm: &[u8]) { + let mut writer = self.writer.lock().expect("external archive lock poisoned"); + if writer.0.is_none() { + return; + } + let result = (|| -> std::io::Result<()> { + let size = u32::try_from(pcm.len()) + .ok() + .and_then(|size| writer.1.checked_add(size)) + .filter(|size| *size <= u32::MAX - 36) + .ok_or_else(|| std::io::Error::other("remote recording exceeds WAV size limit"))?; + let file = writer.0.as_mut().unwrap(); + file.write_all(pcm)?; + // 每帧修正标准头;无需等手机发送 stop,即可读取磁盘上已收到的音频。 + file.seek(SeekFrom::Start(4))?; + file.write_all(&(36 + size).to_le_bytes())?; + file.seek(SeekFrom::Start(40))?; + file.write_all(&size.to_le_bytes())?; + file.seek(SeekFrom::End(0))?; + writer.1 = size; + Ok(()) + })(); + if let Err(error) = result { + log::warn!("[remote-input] 写入录音归档失败:{error}"); + writer.0 = None; + self.available.store(false, Ordering::Release); + } else { + self.available.store(true, Ordering::Release); + } + } + + pub(super) fn finish(&self) { + let mut writer = self.writer.lock().expect("external archive lock poisoned"); + if let Some(file) = writer.0.take() { + if let Err(error) = file.sync_data() { + log::warn!("[remote-input] 同步录音归档失败:{error}"); + self.available.store(false, Ordering::Release); + } + } + if writer.1 == 0 { + let _ = std::fs::remove_file(&self.path); + } + } +} + +impl RecordingArchive for ExternalRecordingArchive { + fn is_available(&self) -> bool { + self.available.load(Ordering::Acquire) + } + + fn read_pcm(&self) -> BoxFuture<'static, Result, BackendError>> { + let path = self.path.clone(); + Box::pin(async move { + let wav = std::fs::read(path).map_err(|error| { + BackendError::new(BackendErrorCode::Persistence, error.to_string()) + })?; + if wav.len() <= 44 + || &wav[..4] != b"RIFF" + || &wav[8..12] != b"WAVE" + || !(wav.len() - 44).is_multiple_of(2) + { + return Err(BackendError::new( + BackendErrorCode::Persistence, + "remote recording WAV is empty or invalid", + )); + } + Ok(wav[44..].to_vec()) + }) + } + + fn discard(&self) -> BoxFuture<'static, Result<(), BackendError>> { + self.writer + .lock() + .expect("external archive lock poisoned") + .0 = None; + let result = match std::fs::remove_file(&self.path) { + Ok(()) => Ok(()), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(error) => Err(BackendError::new( + BackendErrorCode::Persistence, + error.to_string(), + )), + }; + if result.is_ok() { + self.available.store(false, Ordering::Release); + } + Box::pin(async move { result }) + } +} diff --git a/openless-all/app/crates/openless-core/src/lib.rs b/openless-all/app/crates/openless-core/src/lib.rs index 334294671..d6314fdf4 100644 --- a/openless-all/app/crates/openless-core/src/lib.rs +++ b/openless-all/app/crates/openless-core/src/lib.rs @@ -298,8 +298,9 @@ pub use providers::{ }; pub use qa_service::QaService; pub use remote_input_service::{ - constant_time_eq, validate_pairing_pin, RemoteFrameCodec, RemoteInputService, - RemoteStreamSequence, REMOTE_INPUT_MAX_PCM_FRAME_BYTES, REMOTE_INPUT_PAIRING_PIN_LEN, + constant_time_eq, finish_remote_input_connection, validate_pairing_pin, RemoteFrameCodec, + RemoteInputService, RemoteStreamSequence, REMOTE_INPUT_MAX_PCM_FRAME_BYTES, + REMOTE_INPUT_PAIRING_PIN_LEN, }; pub use selection_voice_intent::SelectionVoiceIntent; pub use settings::*; diff --git a/openless-all/app/crates/openless-core/src/remote_input_service.rs b/openless-all/app/crates/openless-core/src/remote_input_service.rs index b7e7257fc..ae77691bd 100644 --- a/openless-all/app/crates/openless-core/src/remote_input_service.rs +++ b/openless-all/app/crates/openless-core/src/remote_input_service.rs @@ -1,12 +1,12 @@ -use std::collections::HashMap; +use std::collections::{HashMap, VecDeque}; use std::sync::{Arc, Mutex}; use futures_util::future::BoxFuture; use crate::credentials::SecretValue; use crate::domains::{ - RemoteAuthResult, RemoteInputApi, RemoteInputConfig, RemoteInputRuntimeAdapter, - RemoteInputServerConfig, RemoteInputStatus, + RemoteAuthResult, RemoteInputApi, RemoteInputConfig, RemoteInputRecovery, + RemoteInputRuntimeAdapter, RemoteInputServerConfig, RemoteInputStatus, }; use crate::errors::{BackendError, BackendErrorCode}; use crate::events::{ @@ -24,6 +24,27 @@ const PIN_LOCK_SECS: u64 = 60; const PIN_FAILS_MAX_ENTRIES: usize = 256; const PIN_GLOBAL_MAX_FAILS: u32 = 20; const PIN_GLOBAL_WINDOW_SECS: u64 = 60; +const RECOVERY_SESSION_CAP: usize = 64; +const RECOVERY_TTL: std::time::Duration = std::time::Duration::from_secs(24 * 60 * 60); + +/// 意外断线后由宿主继续轮询原来的 stop,避免丢弃识别和历史记录收尾。 +/// 主动取消及撤销配对权限仍直接调用 disconnect,不经过此路径。 +pub async fn finish_remote_input_connection( + remote: &dyn RemoteInputApi, + connection_id: SessionId, + session_id: Option, + pending_stop: Option>>, +) -> Result<(), BackendError> { + let result = match pending_stop { + Some(stopping) => stopping.await, + None => match session_id { + Some(session_id) => remote.stop_stream(connection_id, session_id).await, + None => Ok(()), + }, + }; + let cleanup = remote.disconnect(connection_id).await; + result.and(cleanup) +} pub fn validate_pairing_pin(pin: &str) -> bool { pin.len() == REMOTE_INPUT_PAIRING_PIN_LEN && pin.bytes().all(|byte| byte.is_ascii_digit()) @@ -121,6 +142,8 @@ struct RemoteInputState { locale: String, pairing_pin: Option, connections: HashMap, + // 恢复凭据与会话 ID 分离,不能凭公开状态中的会话 ID 读取文字。 + recoverable_sessions: VecDeque<(SessionId, SecretValue, std::time::Instant)>, pin_fails: HashMap)>, global_pin_fails: (u32, std::time::Instant), } @@ -168,6 +191,7 @@ impl RemoteInputService { locale, pairing_pin: None, connections: HashMap::new(), + recoverable_sessions: VecDeque::new(), pin_fails: HashMap::new(), global_pin_fails: (0, std::time::Instant::now()), })), @@ -217,6 +241,7 @@ impl RemoteInputService { .filter_map(|connection| connection.stream.take().map(|stream| stream.session_id)) .collect::>(); state.connections.clear(); + state.recoverable_sessions.clear(); state.running = false; state.starting = false; state.urls.clear(); @@ -484,6 +509,18 @@ impl RemoteInputService { sequence: RemoteStreamSequence::new(session_id), finishing: false, }); + let now = std::time::Instant::now(); + state + .recoverable_sessions + .retain(|(_, _, created)| now.duration_since(*created) < RECOVERY_TTL); + state.recoverable_sessions.push_back(( + session_id, + SecretValue::new(uuid::Uuid::new_v4().to_string()), + now, + )); + while state.recoverable_sessions.len() > RECOVERY_SESSION_CAP { + state.recoverable_sessions.pop_front(); + } true } _ => false, @@ -538,6 +575,9 @@ impl RemoteInputService { let stream = ensure_remote_stream_mut(&mut state, connection_id, session_id)?; if cancel { state.connections.get_mut(&connection_id).unwrap().stream = None; + state + .recoverable_sessions + .retain(|(id, _, _)| *id != session_id); } else { if stream.finishing { return Err(BackendError::new( @@ -666,6 +706,108 @@ impl RemoteInputApi for RemoteInputService { Box::pin(async move { service.disconnect_inner(connection_id).await }) } + fn recover_stream( + &self, + connection_id: SessionId, + session_id: SessionId, + recovery_key: SecretValue, + ) -> BoxFuture<'static, Result> { + let service = self.clone(); + Box::pin(async move { + let allowed = |state: &RemoteInputState| { + state.running + && state.connections.contains_key(&connection_id) + && state.recoverable_sessions.iter().any(|(id, key, created)| { + *id == session_id + && created.elapsed() < RECOVERY_TTL + && constant_time_eq( + key.expose_secret().as_bytes(), + recovery_key.expose_secret().as_bytes(), + ) + }) + }; + let was_pending = { + let state = service + .state + .lock() + .expect("remote input state lock poisoned"); + if !allowed(&state) { + return Ok(RemoteInputRecovery::Unavailable); + } + state.connections.values().any(|connection| { + connection + .stream + .as_ref() + .is_some_and(|stream| stream.session_id == session_id) + }) + }; + let history = service.runtime.read_audio_history(session_id).await?; + let state = service + .state + .lock() + .expect("remote input state lock poisoned"); + // 查询过程中取消、关闭服务或重置 PIN,必须立即撤销恢复权限。 + if !allowed(&state) { + return Ok(RemoteInputRecovery::Unavailable); + } + if let Some(entry) = history { + let text = if entry.final_text.trim().is_empty() { + entry.raw_transcript + } else { + entry.final_text + }; + return Ok(if text.trim().is_empty() { + RemoteInputRecovery::Failed { + has_audio_recording: entry.has_audio_recording == Some(true), + } + } else { + RemoteInputRecovery::Completed { text } + }); + } + Ok( + if was_pending + || state.connections.values().any(|connection| { + connection + .stream + .as_ref() + .is_some_and(|stream| stream.session_id == session_id) + }) + { + RemoteInputRecovery::Pending + } else { + RemoteInputRecovery::Unavailable + }, + ) + }) + } + + fn recovery_key( + &self, + connection_id: SessionId, + session_id: SessionId, + ) -> Result { + let state = self.state.lock().expect("remote input state lock poisoned"); + let owns_stream = state.running + && state + .connections + .get(&connection_id) + .and_then(|connection| connection.stream.as_ref()) + .is_some_and(|stream| stream.session_id == session_id); + if owns_stream { + if let Some((_, key, _)) = state + .recoverable_sessions + .iter() + .find(|(id, _, _)| *id == session_id) + { + return Ok(key.clone()); + } + } + Err(BackendError::new( + BackendErrorCode::Cancelled, + "remote recovery is unavailable", + )) + } + fn start_stream( &self, connection_id: SessionId, diff --git a/openless-all/app/crates/openless-core/tests/remote_input_contract.rs b/openless-all/app/crates/openless-core/tests/remote_input_contract.rs index 104ebb5d1..62cf67c70 100644 --- a/openless-all/app/crates/openless-core/tests/remote_input_contract.rs +++ b/openless-all/app/crates/openless-core/tests/remote_input_contract.rs @@ -20,6 +20,8 @@ struct FixtureRemoteRuntime { audio_start_count: AtomicUsize, audio_stop_count: AtomicUsize, audio_cancel_count: AtomicUsize, + history_reads: AtomicUsize, + history: Mutex>, frames: Mutex)>>, insert_preferences: Mutex>, stop_started: Option>, @@ -136,6 +138,21 @@ impl RemoteInputRuntimeAdapter for FixtureRemoteRuntime { self.audio_cancel_count.fetch_add(1, Ordering::AcqRel); Box::pin(async { Ok(()) }) } + + fn read_audio_history( + &self, + session_id: SessionId, + ) -> BoxFuture<'static, Result, BackendError>> { + self.history_reads.fetch_add(1, Ordering::AcqRel); + let entry = self + .history + .lock() + .unwrap() + .iter() + .find(|entry| entry.id == session_id.to_string()) + .cloned(); + Box::pin(async move { Ok(entry) }) + } } fn backend(runtime: Arc) -> (OpenLessBackend, std::path::PathBuf) { @@ -474,6 +491,215 @@ async fn stream_association_validates_frames_and_rejects_duplicates_and_late_pcm let _ = std::fs::remove_dir_all(data_dir); } +#[tokio::test] +async fn interrupted_connection_finishes_two_minutes_of_received_pcm_without_cancelling() { + let runtime = Arc::new(FixtureRemoteRuntime::default()); + let (backend, data_dir) = backend(Arc::clone(&runtime)); + let remote = &backend.services().remote_input; + remote + .configure(RemoteInputConfig { + enabled: true, + port: 8443, + }) + .await + .unwrap(); + let connection = SessionId::new(); + authenticate(remote.as_ref(), connection).await; + let session = remote.start_stream(connection).await.unwrap(); + for second in 0..120 { + remote + .feed_pcm(connection, session, second, vec![second as u8; 32_000]) + .await + .unwrap(); + } + openless_core::finish_remote_input_connection(remote.as_ref(), connection, Some(session), None) + .await + .unwrap(); + assert_eq!(runtime.audio_stop_count.load(Ordering::Acquire), 1); + assert_eq!(runtime.audio_cancel_count.load(Ordering::Acquire), 0); + assert_eq!( + runtime + .frames + .lock() + .unwrap() + .iter() + .map(|(_, pcm)| pcm.len()) + .sum::(), + 120 * 32_000 + ); + assert_eq!(remote.status().unwrap().connection_count, 0); + assert!(remote + .feed_pcm(connection, session, 120, vec![1, 0]) + .await + .is_err()); + let _ = std::fs::remove_dir_all(data_dir); +} + +#[tokio::test] +async fn disconnect_keeps_an_inflight_stop_but_pin_rotation_can_still_revoke_it() { + for revoke in [false, true] { + let started = Arc::new(tokio::sync::Notify::new()); + let release = Arc::new(tokio::sync::Notify::new()); + let runtime = Arc::new(FixtureRemoteRuntime { + stop_started: Some(started), + release_stop: Some(Arc::clone(&release)), + ..Default::default() + }); + let (backend, data_dir) = backend(Arc::clone(&runtime)); + let remote = Arc::clone(&backend.services().remote_input); + remote + .configure(RemoteInputConfig { + enabled: true, + port: 8443, + }) + .await + .unwrap(); + let connection = SessionId::new(); + authenticate(remote.as_ref(), connection).await; + let session = remote.start_stream(connection).await.unwrap(); + let mut stopping = remote.stop_stream(connection, session); + assert!(futures_util::poll!(stopping.as_mut()).is_pending()); + let finishing = { + let remote = Arc::clone(&remote); + tokio::spawn(async move { + openless_core::finish_remote_input_connection( + remote.as_ref(), + connection, + Some(session), + Some(stopping), + ) + .await + }) + }; + tokio::task::yield_now().await; + assert_eq!(runtime.audio_stop_count.load(Ordering::Acquire), 1); + assert_eq!(runtime.audio_cancel_count.load(Ordering::Acquire), 0); + if revoke { + remote.regenerate_pairing_pin().await.unwrap(); + } + release.notify_one(); + let result = finishing.await.unwrap(); + if revoke { + assert_eq!(result.unwrap_err().code, BackendErrorCode::Cancelled); + } else { + result.unwrap(); + } + assert_eq!( + runtime.audio_cancel_count.load(Ordering::Acquire), + usize::from(revoke) + ); + assert_eq!(remote.status().unwrap().connection_count, 0); + let _ = std::fs::remove_dir_all(data_dir); + } +} + +#[tokio::test] +async fn recovery_requires_authentication_and_a_separate_key_and_obeys_revocation() { + use openless_core::RemoteInputRecovery; + let runtime = Arc::new(FixtureRemoteRuntime::default()); + let (backend, data_dir) = backend(Arc::clone(&runtime)); + let remote = &backend.services().remote_input; + remote + .configure(RemoteInputConfig { + enabled: true, + port: 8443, + }) + .await + .unwrap(); + let connection = SessionId::new(); + authenticate(remote.as_ref(), connection).await; + let session = remote.start_stream(connection).await.unwrap(); + let key = remote.recovery_key(connection, session).unwrap(); + let stranger = SessionId::new(); + assert_eq!( + remote + .recover_stream(stranger, session, key.clone()) + .await + .unwrap(), + RemoteInputRecovery::Unavailable + ); + authenticate(remote.as_ref(), stranger).await; + assert!(remote.recovery_key(stranger, session).is_err()); + assert_eq!( + remote + .recover_stream(stranger, session, SecretValue::new(session.to_string())) + .await + .unwrap(), + RemoteInputRecovery::Unavailable + ); + assert_eq!( + remote + .recover_stream(stranger, SessionId::new(), key.clone()) + .await + .unwrap(), + RemoteInputRecovery::Unavailable + ); + assert_eq!(runtime.history_reads.load(Ordering::Acquire), 0); + assert_eq!( + remote + .recover_stream(stranger, session, key.clone()) + .await + .unwrap(), + RemoteInputRecovery::Pending + ); + let public = format!( + "{} {}", + serde_json::to_string(&remote.status().unwrap()).unwrap(), + serde_json::to_string(&backend.replay_events_after(0)).unwrap() + ); + assert!(!public.contains(key.expose_secret())); + openless_core::finish_remote_input_connection(remote.as_ref(), connection, Some(session), None) + .await + .unwrap(); + runtime.history.lock().unwrap().push(serde_json::from_value(serde_json::json!({ + "id": session.to_string(), "createdAt":"2026-09-10T00:00:00Z", "rawTranscript":"received audio", + "finalText":"recovered text", "mode":"raw", "insertStatus":"notRequested", "hasAudioRecording":true, + "appName":"private desktop metadata" + })).unwrap()); + assert_eq!( + remote + .recover_stream(stranger, session, key.clone()) + .await + .unwrap(), + RemoteInputRecovery::Completed { + text: "recovered text".into() + } + ); + { + let mut history = runtime.history.lock().unwrap(); + history[0].raw_transcript.clear(); + history[0].final_text.clear(); + history[0].error_code = Some("transcribeFailed".into()); + } + assert_eq!( + remote + .recover_stream(stranger, session, key.clone()) + .await + .unwrap(), + RemoteInputRecovery::Failed { + has_audio_recording: true + } + ); + remote.regenerate_pairing_pin().await.unwrap(); + authenticate(remote.as_ref(), stranger).await; + assert_eq!( + remote.recover_stream(stranger, session, key).await.unwrap(), + RemoteInputRecovery::Unavailable + ); + let cancelled = remote.start_stream(stranger).await.unwrap(); + let cancelled_key = remote.recovery_key(stranger, cancelled).unwrap(); + remote.cancel_stream(stranger, cancelled).await.unwrap(); + assert_eq!( + remote + .recover_stream(stranger, cancelled, cancelled_key) + .await + .unwrap(), + RemoteInputRecovery::Unavailable + ); + remote.disconnect(stranger).await.unwrap(); + let _ = std::fs::remove_dir_all(data_dir); +} + #[tokio::test] async fn slow_finalization_remains_cancellable_and_cannot_clear_a_new_stream() { for disconnect in [false, true] { diff --git a/openless-all/app/linux-egui/src/backend.rs b/openless-all/app/linux-egui/src/backend.rs index 8775b8965..7c1ff6a7f 100644 --- a/openless-all/app/linux-egui/src/backend.rs +++ b/openless-all/app/linux-egui/src/backend.rs @@ -879,6 +879,12 @@ impl LinuxBackendBuilder { let recorder = self .recorder .unwrap_or_else(|| Arc::new(LinuxCpalRecorder::new(None)) as Arc); + let recorder: Arc = Arc::new(openless_core::AudioRecorderRouter::new( + recorder, + openless_core::ExternalAudioRecorder::with_recordings_directory( + self.config.data_dir.join("recordings"), + ), + )); let text_inserter = self .text_inserter .unwrap_or_else(|| Arc::new(Fcitx5TextInserter::new(true)) as Arc); diff --git a/openless-all/app/linux-egui/src/remote_input.rs b/openless-all/app/linux-egui/src/remote_input.rs index 070579d92..efad7eddd 100644 --- a/openless-all/app/linux-egui/src/remote_input.rs +++ b/openless-all/app/linux-egui/src/remote_input.rs @@ -172,6 +172,19 @@ impl RemoteInputRuntimeAdapter for LinuxRemoteInputRuntime { let backend = self.backend(); Box::pin(async move { backend?.cancel_dictation(Some(session_id)).await }) } + + fn read_audio_history( + &self, + session_id: SessionId, + ) -> BoxFuture<'static, Result, BackendError>> { + let backend = self.backend(); + Box::pin(async move { + Ok(backend? + .list_history()? + .into_iter() + .find(|entry| entry.id == session_id.to_string())) + }) + } } #[cfg(target_os = "linux")] @@ -208,6 +221,8 @@ mod assets { const KEEPALIVE_PING_SECS: u64 = 30; #[cfg(target_os = "linux")] const IDLE_TIMEOUT_SECS: u64 = 90; +#[cfg(target_os = "linux")] +const AUDIO_IDLE_TIMEOUT_SECS: u64 = 15; #[cfg(target_os = "linux")] struct LinuxRemoteServerHandle { @@ -498,18 +513,27 @@ async fn websocket_session(mut socket: WebSocket, state: Arc, peer: Ip let mut keepalive = tokio::time::interval(Duration::from_secs(KEEPALIVE_PING_SECS)); keepalive.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); let mut last_received = Instant::now(); + let mut last_audio = Instant::now(); + let mut audio_watchdog = tokio::time::interval(Duration::from_secs(2)); + audio_watchdog.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + let mut preserve_audio = true; + let mut receiving_audio = false; 'connection: loop { tokio::select! { incoming = socket.recv() => { last_received = Instant::now(); match incoming { Some(Ok(Message::Binary(frame))) => { + if !receiving_audio { continue; } if let Ok((session_id, sequence, pcm)) = openless_core::RemoteFrameCodec::decode(&frame) { - let _ = state.backend.services().remote_input - .feed_pcm(connection_id, session_id, sequence, pcm).await; + if state.backend.services().remote_input + .feed_pcm(connection_id, session_id, sequence, pcm).await.is_ok() { + last_audio = Instant::now(); + } } } Some(Ok(Message::Text(text))) => { + let previous_session = remote_session; if let Some(reply) = apply_control( &text, state.backend.services().remote_input.as_ref(), @@ -519,6 +543,11 @@ async fn websocket_session(mut socket: WebSocket, state: Arc, peer: Ip ).await { if socket.send(json_message(reply)).await.is_err() { break; } } + if remote_session != previous_session { + last_audio = Instant::now(); + receiving_audio = remote_session.is_some(); + } + if pending_stop.is_some() { receiving_audio = false; } } Some(Ok(Message::Close(_))) | Some(Err(_)) | None => break, _ => {} @@ -533,9 +562,7 @@ async fn websocket_session(mut socket: WebSocket, state: Arc, peer: Ip phase: openless_core::DictationPhase::Cancelled | openless_core::DictationPhase::Failed, .. })); for reply in remote_event(event.kind) { - // A failed outbound send means this socket no longer owns - // a usable transport. Leave the outer loop immediately so - // Core disconnect cancels any active external-audio lease. + // 下行失败也保留已收到的录音,由外层继续完成识别和持久化。 if socket.send(json_message(reply)).await.is_err() { break 'connection; } @@ -561,11 +588,37 @@ async fn websocket_session(mut socket: WebSocket, state: Arc, peer: Ip break; } } + _ = audio_watchdog.tick(), if receiving_audio && pending_stop.is_none() && remote_session.is_some() => { + if last_audio.elapsed() >= Duration::from_secs(AUDIO_IDLE_TIMEOUT_SECS) { + receiving_audio = false; + let session_id = remote_session.unwrap(); + pending_stop = Some((session_id, state.backend.services().remote_input.stop_stream(connection_id, session_id))); + } + } changed = shutdown.changed() => { - if changed.is_err() || *shutdown.borrow() { break; } + if changed.is_err() || *shutdown.borrow() { preserve_audio = false; break; } } } } + drop(socket); + if preserve_audio && !*shutdown.borrow() { + let finishing = openless_core::finish_remote_input_connection( + state.backend.services().remote_input.as_ref(), + connection_id, + remote_session, + pending_stop.take().map(|(_, future)| future), + ); + tokio::pin!(finishing); + tokio::select! { + result = &mut finishing => { + if let Err(error) = result { log::warn!("[remote-input] 断线录音收尾:{error}"); } + } + _ = shutdown.changed() => { + let _ = state.backend.services().remote_input.disconnect(connection_id).await; + } + } + return; + } let _ = state .backend .services() @@ -610,7 +663,10 @@ async fn apply_control( "start" => match remote.start_stream(connection_id).await { Ok(session_id) => { *remote_session = Some(session_id); - Some(serde_json::json!({"type":"started", "sessionId":session_id.to_string()})) + Some( + serde_json::json!({"type":"started", "sessionId":session_id.to_string(), + "recoveryKey":remote.recovery_key(connection_id, session_id).ok().map(|key| key.into_exposed())}), + ) } Err(error) => Some(serde_json::json!({"type":"busy", "reason":error.to_string()})), }, @@ -636,6 +692,32 @@ async fn apply_control( } None } + "recover" => { + let session_id = value + .get("sessionId") + .and_then(|value| value.as_str()) + .and_then(|value| uuid::Uuid::parse_str(value).ok()) + .map(SessionId::from_uuid); + let recovery = match session_id { + Some(session_id) => remote + .recover_stream( + connection_id, + session_id, + SecretValue::new( + value + .get("recoveryKey") + .and_then(|value| value.as_str()) + .unwrap_or_default(), + ), + ) + .await + .unwrap_or(openless_core::RemoteInputRecovery::Unavailable), + None => openless_core::RemoteInputRecovery::Unavailable, + }; + Some( + serde_json::json!({"type":"recovery", "sessionId":session_id, "recovery":recovery}), + ) + } "set_insert" => { let insert = value .get("value") diff --git a/openless-all/app/scripts/remote-input-audio-queue.test.mjs b/openless-all/app/scripts/remote-input-audio-queue.test.mjs index 86c2e88e3..180044abe 100644 --- a/openless-all/app/scripts/remote-input-audio-queue.test.mjs +++ b/openless-all/app/scripts/remote-input-audio-queue.test.mjs @@ -33,12 +33,23 @@ function fakeElement() { }; } -async function openRemotePage({ defaultMode, savedMode } = {}) { +const fixtureRecoveryKey = '40112233-4455-4677-8899-aabbccddeeff'; +async function openRemotePage({ defaultMode, savedMode, savedWakeLock, savedRecovery, wakeLockMode = 'supported' } = {}) { const elements = new Map(); const documentListeners = {}; const sent = []; let socket; let worklet; + let audioContext; + const windowListeners = {}; + const wakeLocks = []; + const tracks = []; + const wakeResolvers = []; + let wakeRequests = 0; + const timers = new Map(); + let now = 0; + let timerId = 0; + const flush = async () => { for (let i = 0; i < 12; i += 1) await Promise.resolve(); }; const element = (id) => { if (!elements.has(id)) elements.set(id, fakeElement()); @@ -73,11 +84,12 @@ async function openRemotePage({ defaultMode, savedMode } = {}) { class FakeAudioContext { constructor() { + audioContext = this; this.state = 'running'; this.sampleRate = 48_000; this.audioWorklet = { addModule: () => Promise.resolve() }; } - resume() { return Promise.resolve(); } + resume() { this.state = 'running'; return Promise.resolve(); } suspend() { this.state = 'suspended'; } createMediaStreamSource() { return { connect() {}, disconnect() {} }; } } @@ -105,24 +117,49 @@ async function openRemotePage({ defaultMode, savedMode } = {}) { AudioContext: FakeAudioContext, AudioWorkletNode: FakeAudioWorkletNode, WebSocket: FakeWebSocket, - clearTimeout, + clearTimeout: id => timers.delete(id), console, document, isNaN, localStorage: storage([ ['ol_remote_pin', '123456'], ...(savedMode === undefined ? [] : [['ol_remote_mode', savedMode]]), + ...(savedWakeLock === undefined ? [] : [['ol_remote_wake_lock', savedWakeLock]]), + ...(savedRecovery === undefined ? [] : [['ol_remote_recovery_session', JSON.stringify({ sessionId: savedRecovery, key: fixtureRecoveryKey })]]), ]), location: { host: 'localhost:8443', origin: 'https://localhost:8443', reload() {} }, navigator: { language: 'zh-CN', + wakeLock: wakeLockMode === 'unsupported' ? undefined : { + request: type => { + assert.equal(type, 'screen'); + wakeRequests++; + if (wakeLockMode === 'rejected') return Promise.reject(new Error('system denied')); + const sentinel = { + released: false, + releaseCount: 0, + listener: null, + addEventListener(type, listener) { assert.equal(type, 'release'); this.listener = listener; }, + release() { this.released = true; this.releaseCount++; this.listener?.(); return Promise.resolve(); }, + }; + wakeLocks.push(sentinel); + return wakeLockMode === 'deferred' + ? new Promise(resolve => wakeResolvers.push(() => resolve(sentinel))) + : Promise.resolve(sentinel); + }, + }, mediaDevices: { - getUserMedia: () => Promise.resolve({ getTracks: () => [{ stop() {} }] }), + getUserMedia: () => { + const track = { stopped: false, stop() { this.stopped = true; } }; + tracks.push(track); + return Promise.resolve({ getTracks: () => [track] }); + }, }, }, performance: { now: () => 100 }, sessionStorage: storage([['ol_reloaded_once', '1']]), - setTimeout, + setTimeout: (callback, delay) => { const id = ++timerId; timers.set(id, { callback, at: now + delay }); return id; }, + addEventListener: (type, listener) => { windowListeners[type] = listener; }, }; context.window = context; // Exercise the embedded HTML script too, so a missing template variable cannot @@ -141,11 +178,29 @@ async function openRemotePage({ defaultMode, savedMode } = {}) { documentListeners, element, sent, - socket, + get socket() { return socket; }, + get audioContext() { return audioContext; }, + get wakeRequests() { return wakeRequests; }, + wakeLocks, + tracks, + wakeResolvers, + windowListeners, + flush, + async advance(ms) { + now += ms; + for (const [id, timer] of [...timers]) { + if (timer.at <= now && timers.delete(id)) timer.callback(); + } + await flush(); + }, storage: context.localStorage, async start() { - element('btn-record').listeners.click(); - for (let i = 0; i < 8; i += 1) await Promise.resolve(); + if (element('btn-record').style.touchAction === 'none') { + element('btn-record').listeners.pointerdown({ preventDefault() {} }); + } else { + element('btn-record').listeners.click(); + } + await flush(); assert.ok(worklet?.port.onmessage, 'audio capture must be running'); }, pcm(bytes) { worklet.port.onmessage({ data: Uint8Array.from(bytes).buffer }); }, @@ -231,8 +286,9 @@ for (const terminal of ['cancel', 'busy', 'disconnect']) { await page.start(); page.pcm([6, 0]); if (terminal === 'cancel') { - page.document.hidden = true; - page.documentListeners.visibilitychange(); + page.element('mode-switch').listeners.click({ + target: { closest: () => ({ getAttribute: () => 'hold' }) }, + }); } else if (terminal === 'busy') { page.socket.onmessage({ data: JSON.stringify({ type: 'busy', reason: 'test' }) }); } else { @@ -259,4 +315,135 @@ for (const terminal of ['cancel', 'busy', 'disconnect']) { assert.equal(binaryFrames(page.sent).length, 0); } +const controls = (page, type) => page.sent.filter(value => typeof value === 'string') + .map(value => JSON.parse(value)).filter(value => value.type === type); +const recordingId = '30112233-4455-4677-8899-aabbccddeeff'; +const acknowledge = page => page.socket.onmessage({ data: JSON.stringify({ type: 'started', sessionId: recordingId, recoveryKey: fixtureRecoveryKey }) }); + +// 息屏结束两分钟录音,已发送的帧仍在 stop 之前;重复生命周期事件不会重复结束。 +for (const defaultMode of ['toggle', 'hold']) { + const page = await openRemotePage({ defaultMode }); + assert.equal(page.element('wake-lock-switch').checked, true); + assert.equal(page.storage.getItem('ol_remote_wake_lock'), null); + await page.start(); + acknowledge(page); + for (let second = 0; second < 120; second++) page.pcm(new Uint8Array(32_000).fill(second)); + page.document.hidden = true; + page.documentListeners.visibilitychange(); + page.windowListeners.pagehide(); + assert.equal(controls(page, 'stop').length, 1); + assert.equal(controls(page, 'cancel').length, 0); + assert.equal(binaryFrames(page.sent).reduce((bytes, frame) => bytes + frame.byteLength - 28, 0), 120 * 32_000); + assert.equal(JSON.parse(page.sent.at(-1)).type, 'stop'); + assert.equal(page.wakeLocks[0].released, true); + assert.deepEqual(JSON.parse(page.storage.getItem('ol_remote_recovery_session')), { sessionId: recordingId, key: fixtureRecoveryKey }); + page.document.hidden = false; + page.documentListeners.visibilitychange(); + assert.equal(controls(page, 'start').length, 1, '恢复前台不会自动打开麦克风'); + assert.equal(page.wakeRequests, 1); +} + +{ + const page = await openRemotePage(); + await page.start(); + page.pcm([7, 0, 8, 0]); + page.document.hidden = true; + page.documentListeners.visibilitychange(); + assert.equal(controls(page, 'cancel').length, 0); + acknowledge(page); + assert.deepEqual(payload(binaryFrames(page.sent)[0]), [7, 0, 8, 0]); + assert.equal(JSON.parse(page.sent.at(-1)).type, 'stop', 'ACK 迟到时保留首段音频并按序结束'); +} + +{ + const page = await openRemotePage({ savedWakeLock: '0' }); + assert.equal(page.element('wake-lock-switch').checked, false); + await page.start(); + assert.equal(page.wakeRequests, 0); + const toggle = page.element('wake-lock-switch'); + toggle.checked = true; + toggle.listeners.change(); + await page.flush(); + assert.equal(page.wakeRequests, 1); + assert.equal(page.storage.getItem('ol_remote_wake_lock'), '1'); + toggle.checked = false; + toggle.listeners.change(); + assert.equal(page.wakeLocks[0].releaseCount, 1); + assert.equal(page.storage.getItem('ol_remote_wake_lock'), '0'); +} + +for (const wakeLockMode of ['unsupported', 'rejected']) { + const page = await openRemotePage({ wakeLockMode }); + await page.start(); + assert.equal(controls(page, 'start').length, 1, '保持亮屏失败不阻止录音'); + assert.match(page.element('wake-lock-hint').textContent, /未允许/); + acknowledge(page); + page.audioContext.state = 'interrupted'; + page.audioContext.onstatechange(); + assert.equal(controls(page, 'stop').length, 1); + assert.equal(controls(page, 'cancel').length, 0); +} + +{ + const page = await openRemotePage(); + await page.start(); + acknowledge(page); + page.tracks[0].onended(); + assert.equal(controls(page, 'stop').length, 1); + assert.equal(page.tracks[0].stopped, true); + page.socket.onmessage({ data: JSON.stringify({ type: 'status', kind: 'done' }) }); + await page.start(); + assert.equal(page.tracks.length, 2, '系统中断后的下一次录音必须重新获取麦克风'); + assert.equal(page.tracks[1].stopped, false); +} + +{ + const page = await openRemotePage({ wakeLockMode: 'deferred' }); + await page.start(); + acknowledge(page); + page.element('btn-record').listeners.click(); + page.wakeResolvers[0](); + await page.flush(); + assert.equal(page.wakeLocks[0].releaseCount, 1, '录音结束后迟到的请求必须释放'); +} + +{ + const page = await openRemotePage(); + await page.start(); + await page.wakeLocks[0].release(); + assert.match(page.element('wake-lock-hint').textContent, /未允许/); + assert.equal(page.wakeRequests, 1, '系统释放后不能循环申请'); + acknowledge(page); + page.socket.readyState = 3; + page.socket.onclose(); + page.documentListeners.visibilitychange(); + page.socket.onopen(); + page.socket.onmessage({ data: JSON.stringify({ type: 'auth', ok: true }) }); + assert.equal(controls(page, 'recover').at(-1).sessionId, recordingId); + assert.equal(controls(page, 'recover').at(-1).recoveryKey, fixtureRecoveryKey); + page.socket.onmessage({ data: JSON.stringify({ type: 'recovery', sessionId: recordingId, recovery: { kind: 'pending' } }) }); + assert.equal(page.element('btn-record').disabled, true); + await page.advance(1500); + assert.equal(controls(page, 'recover').length, 2); + page.socket.onmessage({ data: JSON.stringify({ type: 'recovery', sessionId: recordingId, recovery: { kind: 'completed', text: '保留下来的两分钟录音' } }) }); + assert.equal(page.element('result-text').textContent, '保留下来的两分钟录音'); + assert.equal(page.element('btn-record').disabled, false); +} + +{ + const page = await openRemotePage({ savedRecovery: recordingId }); + const oldSocket = page.socket; + await page.advance(8000); + assert.notEqual(page.socket, oldSocket, '半开连接恢复超时后重新认证'); + assert.equal(oldSocket.readyState, 3); +} + +for (const recovery of [{ kind: 'failed', hasAudioRecording: true }, { kind: 'unavailable' }]) { + const page = await openRemotePage({ savedRecovery: recordingId }); + page.socket.onmessage({ data: JSON.stringify({ type: 'recovery', sessionId: recordingId, recovery }) }); + assert.match(page.element('status-text').textContent, /历史记录/); + assert.equal(page.element('btn-record').disabled, false); + if (recovery.kind === 'unavailable') assert.equal(page.storage.getItem('ol_remote_recovery_session'), null); +} + console.log('remote-input-audio-queue.test.mjs passed'); diff --git a/openless-all/app/src-tauri/src/core_adapters.rs b/openless-all/app/src-tauri/src/core_adapters.rs index b5e65b34b..a0b73e450 100644 --- a/openless-all/app/src-tauri/src/core_adapters.rs +++ b/openless-all/app/src-tauri/src/core_adapters.rs @@ -193,16 +193,27 @@ pub(crate) fn backend_dependencies( app: Arc::clone(&app), backend: Arc::clone(&backend), }); - let recorder = - AudioRecorderRouter::new(Arc::clone(&host_recorder), ExternalAudioRecorder::default()); + let external = match crate::persistence::data_dir() { + Ok(directory) => { + ExternalAudioRecorder::with_recordings_directory(directory.join("recordings")) + } + Err(error) => { + log::warn!("[remote-input] 无法定位录音目录:{error}"); + ExternalAudioRecorder::default() + } + }; + let recorder: Arc = Arc::new(AudioRecorderRouter::new( + Arc::clone(&host_recorder), + external, + )); let traditional = Arc::new(openless_core::PipelineDictationEngine::new( - Arc::new(recorder), + Arc::clone(&recorder), transcription, Arc::clone(&polisher), )); let dictation = Arc::new(openless_core::DictationEngineRouter::new(traditional)); let production_omni: Arc = Arc::new( - openless_core::SharedOmniDictationEngine::new(Arc::clone(&credential_store), host_recorder), + openless_core::SharedOmniDictationEngine::new(Arc::clone(&credential_store), recorder), ); for provider_type in openless_core::SHARED_OMNI_PROVIDER_TYPES { dictation @@ -1353,6 +1364,19 @@ impl openless_core::RemoteInputRuntimeAdapter for TauriRemoteInputRuntimeAdapter let backend = self.backend(); Box::pin(async move { backend?.cancel_dictation(Some(session_id)).await }) } + + fn read_audio_history( + &self, + session_id: SessionId, + ) -> BoxFuture<'static, Result, BackendError>> { + let backend = self.backend(); + Box::pin(async move { + Ok(backend? + .list_history()? + .into_iter() + .find(|entry| entry.id == session_id.to_string())) + }) + } } struct TauriPlatformApi { diff --git a/openless-all/app/src-tauri/src/remote_server/assets/app.js b/openless-all/app/src-tauri/src/remote_server/assets/app.js index e515a58cd..6c004f76a 100644 --- a/openless-all/app/src-tauri/src/remote_server/assets/app.js +++ b/openless-all/app/src-tauri/src/remote_server/assets/app.js @@ -14,6 +14,16 @@ // ============================================================ var I18N = { 'zh-CN': { + wakeLockLabel: "录音时保持亮屏", + wakeLockHint: "息屏会结束本段录音,电脑继续处理已收到的部分。", + wakeLockActive: "屏幕保持亮起,录音结束后允许自动息屏。", + wakeLockUnavailable: "浏览器或系统未允许保持亮屏;息屏后电脑会处理已收到的录音。", + interrupted: "录音已中断,电脑继续识别已收到的部分…", + offlineRecording: "连接已断开,电脑会继续处理已收到的录音。重连后可查看结果。", + recovering: "电脑正在处理上次录音…", + recovered: "已找回上次识别结果", + recoveryRetry: "识别未完成,录音已保存在电脑历史记录中,可重新转录。", + recoveryUnavailable: "暂未找到结果,请到电脑的历史记录中查看。", title: 'OpenLess 远程输入', brandTitle: 'OpenLess 远程输入', brandSub: '在手机上录音,实时输入到电脑', @@ -69,6 +79,16 @@ copied: '已复制 ✓', }, 'zh-TW': { + wakeLockLabel: "錄音時保持螢幕開啟", + wakeLockHint: "螢幕關閉會結束本段錄音,電腦繼續處理已收到的部分。", + wakeLockActive: "螢幕保持開啟,錄音結束後允許自動關閉螢幕。", + wakeLockUnavailable: "瀏覽器或系統未允許保持螢幕開啟;電腦會處理已收到的錄音。", + interrupted: "錄音已中斷,電腦繼續辨識已收到的部分…", + offlineRecording: "連線已中斷,電腦會繼續處理已收到的錄音。重新連線後可查看結果。", + recovering: "電腦正在處理上次錄音…", + recovered: "已找回上次辨識結果", + recoveryRetry: "辨識未完成,錄音已保存在電腦歷史記錄中,可重新轉錄。", + recoveryUnavailable: "暫未找到結果,請到電腦的歷史記錄中查看。", title: 'OpenLess 遠端輸入', brandTitle: 'OpenLess 遠端輸入', brandSub: '在手機上錄音,即時輸入到電腦', @@ -124,6 +144,16 @@ copied: '已複製 ✓', }, en: { + wakeLockLabel: "Keep screen awake while recording", + wakeLockHint: "Screen lock ends this recording. The computer processes the audio already received.", + wakeLockActive: "Screen stays awake until recording ends.", + wakeLockUnavailable: "The browser or system did not allow screen wake lock. Received audio will still be processed.", + interrupted: "Recording interrupted. The computer is processing the audio received…", + offlineRecording: "Disconnected. The computer continues processing received audio. Reconnect to see the result.", + recovering: "The computer is processing your last recording…", + recovered: "Last transcription recovered", + recoveryRetry: "Transcription failed. The recording is saved in computer history and can be retried.", + recoveryUnavailable: "Result unavailable. Please check history on the computer.", title: 'OpenLess Remote Input', brandTitle: 'OpenLess Remote Input', brandSub: 'Record on your phone, type to your computer in real time', @@ -179,6 +209,16 @@ copied: 'Copied ✓', }, ja: { + wakeLockLabel: "録音中は画面をオンにする", + wakeLockHint: "画面をロックすると録音を終了し、受信済みの音声をパソコンで処理します。", + wakeLockActive: "録音が終わるまで画面をオンに保ちます。", + wakeLockUnavailable: "ブラウザーまたはシステムが画面の維持を許可しませんでした。受信済みの音声は処理されます。", + interrupted: "録音が中断されました。受信済みの音声をパソコンで処理しています…", + offlineRecording: "接続が切れました。受信済みの音声の処理は続きます。再接続すると結果を確認できます。", + recovering: "前回の録音をパソコンで処理しています…", + recovered: "前回の文字起こし結果を復元しました", + recoveryRetry: "文字起こしが完了しませんでした。録音はパソコンの履歴に保存され、再試行できます。", + recoveryUnavailable: "結果が見つかりません。パソコンの履歴を確認してください。", title: 'OpenLess リモート入力', brandTitle: 'OpenLess リモート入力', brandSub: 'スマホで録音し、リアルタイムでパソコンに入力', @@ -234,6 +274,16 @@ copied: 'コピー済み ✓', }, ko: { + wakeLockLabel: "녹음 중 화면 켜짐 유지", + wakeLockHint: "화면을 잠그면 녹음이 끝나고 컴퓨터가 이미 받은 오디오를 처리합니다.", + wakeLockActive: "녹음이 끝날 때까지 화면을 켜진 상태로 유지합니다.", + wakeLockUnavailable: "브라우저 또는 시스템이 화면 켜짐 유지를 허용하지 않았습니다. 수신한 오디오는 계속 처리됩니다.", + interrupted: "녹음이 중단되었습니다. 컴퓨터가 받은 오디오를 처리하고 있습니다…", + offlineRecording: "연결이 끊겼습니다. 받은 오디오는 계속 처리됩니다. 다시 연결하면 결과를 볼 수 있습니다.", + recovering: "컴퓨터가 마지막 녹음을 처리하고 있습니다…", + recovered: "마지막 음성 인식 결과를 복구했습니다", + recoveryRetry: "음성 인식을 완료하지 못했습니다. 녹음은 컴퓨터 기록에 저장되며 다시 시도할 수 있습니다.", + recoveryUnavailable: "결과를 찾을 수 없습니다. 컴퓨터의 기록을 확인해 주세요.", title: 'OpenLess 원격 입력', brandTitle: 'OpenLess 원격 입력', brandSub: '휴대폰으로 녹음하여 실시간으로 컴퓨터에 입력', @@ -329,6 +379,8 @@ var MODE_KEY = 'ol_remote_mode'; // localStorage 键:录音方式 var PIN_KEY = 'ol_remote_pin'; // localStorage 键:上次成功的配对码 var INSERT_KEY = 'ol_remote_insert'; // localStorage 键:电脑落字开关(默认开) + var WAKE_LOCK_KEY = 'ol_remote_wake_lock'; + var RECOVERY_KEY = 'ol_remote_recovery_session'; var MIC_PREP_TIMEOUT_MS = 10000; // 麦克风准备超时:超过则判失败让用户重试,避免无限卡"准备中" var PCM_QUEUE_MAX_BYTES = 128 * 1024; @@ -355,6 +407,8 @@ var recTip = $('rec-tip'); var modeSwitch = $('mode-switch'); var insertSwitch = $('insert-switch'); + var wakeLockSwitch = $('wake-lock-switch'); + var wakeLockHint = $('wake-lock-hint'); var btnReconnect = $('btn-reconnect'); var offlineReason = $('offline-reason'); @@ -373,6 +427,14 @@ var finishAfterStarted = ''; // ACK 前松手/取消:'stop' | 'cancel' | '' var pendingPcm = []; var pendingPcmBytes = 0; + var awaitingResult = false; + var savedRecovery = readRecoverySession(); + var recoverySessionId = savedRecovery.sessionId; + var recoveryKey = savedRecovery.key; + var recoveryTimer = null; + var wakeLock = null; + var wakeLockGeneration = 0; + var wakeLockPending = null; // 音频相关 var audioCtx = null; @@ -403,6 +465,114 @@ } function clearPin() { try { localStorage.removeItem(PIN_KEY); } catch (e) {} + saveRecoverySession(''); + } + + // 恢复凭据仅用于本次随机会话;不请求电脑历史记录列表。 + function readRecoverySession() { + try { + var saved = JSON.parse(localStorage.getItem(RECOVERY_KEY) || 'null'); + if (saved && validSessionId(saved.sessionId) && validSessionId(saved.key)) return saved; + } catch (e) {} + return { sessionId: '', key: '' }; + } + function validSessionId(id) { + return typeof id === 'string' && /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i.test(id); + } + function saveRecoverySession(id, key) { + recoverySessionId = validSessionId(id) && validSessionId(key) ? id : ''; + recoveryKey = recoverySessionId ? key : ''; + try { + if (recoverySessionId) localStorage.setItem(RECOVERY_KEY, JSON.stringify({ sessionId: recoverySessionId, key: recoveryKey })); + else localStorage.removeItem(RECOVERY_KEY); + } catch (e) {} + } + function clearRecoveryTimer() { + if (recoveryTimer) { clearTimeout(recoveryTimer); recoveryTimer = null; } + } + function requestRecovery() { + clearRecoveryTimer(); + if (!authed || document.hidden || recording || startSent || !recoverySessionId) return; + wsSendJSON({ type: 'recover', sessionId: recoverySessionId, recoveryKey: recoveryKey }); + // 唤醒后的旧连接可能仍显示 OPEN,却再也收不到数据;超时重新认证。 + recoveryTimer = setTimeout(function () { + recoveryTimer = null; + if (!recording && authed && !document.hidden) { + var pin = readPin(); + if (pin) connect(pin); + } + }, 8000); + } + function handleRecovery(msg) { + if (msg.sessionId !== recoverySessionId || recording || startSent) return; + clearRecoveryTimer(); + clearWorkTimeout(); + var recovery = msg.recovery || {}; + awaitingResult = recovery.kind === 'pending'; + updateRecordBtnUI(); + if (recovery.kind === 'pending') { + setStatus(L.recovering, 'work'); + if (!document.hidden) recoveryTimer = setTimeout(requestRecovery, 1500); + } else if (recovery.kind === 'completed') { + showResult(recovery.text); + setStatus(L.recovered, 'ok'); + } else if (recovery.kind === 'failed') { + setStatus(recovery.hasAudioRecording ? L.recoveryRetry : L.recoveryUnavailable, 'error'); + } else { + saveRecoverySession(''); + setStatus(L.recoveryUnavailable, 'error'); + } + } + + function shouldKeepAwake() { + return recording && !document.hidden && wakeLockSwitch && wakeLockSwitch.checked; + } + function releaseWakeLock() { + wakeLockGeneration++; + wakeLockPending = null; + var previous = wakeLock; + wakeLock = null; + if (previous) previous.release().catch(function () {}); + } + function acquireWakeLock() { + if (!shouldKeepAwake() || wakeLock || wakeLockPending !== null) return; + if (!navigator.wakeLock || !navigator.wakeLock.request) { + if (wakeLockHint) wakeLockHint.textContent = L.wakeLockUnavailable; + return; + } + var generation = wakeLockGeneration; + wakeLockPending = generation; + navigator.wakeLock.request('screen').then(function (sentinel) { + if (wakeLockPending === generation) wakeLockPending = null; + if (generation !== wakeLockGeneration || !shouldKeepAwake()) { + sentinel.release().catch(function () {}); + return; + } + wakeLock = sentinel; + if (wakeLockHint) wakeLockHint.textContent = L.wakeLockActive; + sentinel.addEventListener('release', function () { + if (wakeLock !== sentinel) return; + wakeLock = null; + if (shouldKeepAwake() && wakeLockHint) wakeLockHint.textContent = L.wakeLockUnavailable; + }); + }).catch(function () { + if (wakeLockPending === generation) wakeLockPending = null; + if (generation === wakeLockGeneration && shouldKeepAwake() && wakeLockHint) { + wakeLockHint.textContent = L.wakeLockUnavailable; + } + }); + } + function initWakeLockSwitch() { + if (!wakeLockSwitch) return; + try { wakeLockSwitch.checked = localStorage.getItem(WAKE_LOCK_KEY) !== '0'; } + catch (e) { wakeLockSwitch.checked = true; } + if (wakeLockHint) wakeLockHint.textContent = L.wakeLockHint; + wakeLockSwitch.addEventListener('change', function () { + try { localStorage.setItem(WAKE_LOCK_KEY, wakeLockSwitch.checked ? '1' : '0'); } catch (e) {} + if (wakeLockHint) wakeLockHint.textContent = L.wakeLockHint; + if (wakeLockSwitch.checked) acquireWakeLock(); + else releaseWakeLock(); + }); } // ============================================================ @@ -548,7 +718,10 @@ workTimer = null; if (!recording && authed) { if (startSent) failRecording('❌ ' + L.errGeneric, true); + else if (recoverySessionId) requestRecovery(); else { + awaitingResult = false; + updateRecordBtnUI(); setStatus('❌ ' + L.errGeneric, 'error'); setLevel(0); } @@ -594,6 +767,7 @@ closeWS(); // 清理旧连接 authed = false; busy = false; + awaitingResult = false; var url = 'wss://' + location.host + '/ws'; try { @@ -626,16 +800,18 @@ clearConnectTimeout(); clearReadyTimer(); clearWorkTimeout(); + clearRecoveryTimer(); if (busyTimer) { clearTimeout(busyTimer); busyTimer = null; } var wasAuthed = authed; authed = false; recording = false; + awaitingResult = false; detachHoldEnd(); resetRemoteStreamState(); teardownAudio(); if (wasAuthed) { // 已进入录音屏后断开 → 断线屏 - offlineReason.textContent = L.offlineSub; + offlineReason.textContent = recoverySessionId ? L.offlineRecording : L.offlineSub; showScreen('offline'); } else { // 未认证就关闭(握手被拒/证书不受信任/网络中断)。无论当前是否在配对屏都给出 @@ -652,8 +828,10 @@ clearConnectTimeout(); clearReadyTimer(); clearWorkTimeout(); + clearRecoveryTimer(); if (busyTimer) { clearTimeout(busyTimer); busyTimer = null; } recording = false; + awaitingResult = false; detachHoldEnd(); resetRemoteStreamState(); teardownAudio(); @@ -675,6 +853,7 @@ clearConnectTimeout(); writePin(lastPin); // 配对成功 → 记住配对码,刷新后免重输 enterRecScreen(); + requestRecovery(); } else { authed = false; clearPin(); // 配对码失效(错误/锁定)→ 清除,避免下次自动重连又失败 @@ -691,7 +870,11 @@ break; case 'started': - handleStarted(msg.sessionId); + handleStarted(msg.sessionId, msg.recoveryKey); + break; + + case 'recovery': + handleRecovery(msg); break; case 'level': @@ -702,6 +885,7 @@ clearWorkTimeout(); busy = true; recording = false; + awaitingResult = false; resetRemoteStreamState(); teardownAudioCapture(); // 停止采集但保留 ctx updateRecordBtnUI(); @@ -719,6 +903,9 @@ case 'result': // 电脑落字完成后回传的最终文字,显示给手机用户看本次识别结果。 showResult(msg.text); + awaitingResult = false; + clearRecoveryTimer(); + updateRecordBtnUI(); break; } } @@ -735,6 +922,14 @@ setStatus(stripLeadingIcon(L.statusRecording), 'work'); break; case 'transcribing': + if (recording) { + recording = false; + detachHoldEnd(); + resetRemoteStreamState(); + teardownAudioCapture(); + } + awaitingResult = true; + updateRecordBtnUI(); setStatus(stripLeadingIcon(L.statusTranscribing), 'work'); if (statusDots) statusDots.hidden = false; // 识别中:三点加载动效 armWorkTimeout(); // 工作状态续上兜底超时,防止服务端中途无响应卡死 @@ -744,6 +939,8 @@ armWorkTimeout(); // 同上 break; case 'done': + awaitingResult = false; + updateRecordBtnUI(); clearWorkTimeout(); // 正常收尾,解除兜底超时 var n = (typeof msg.insertedChars === 'number') ? msg.insertedChars : 0; setStatus(stripLeadingIcon(fmt(L.statusDone, { n: n })), 'ok'); @@ -752,6 +949,8 @@ scheduleReady(); break; case 'error': + awaitingResult = false; + updateRecordBtnUI(); clearWorkTimeout(); // 服务端已明确报错,解除兜底超时 if (recording || startSent) failRecording('❌ ' + (msg.message || L.errGeneric), true); else { @@ -888,7 +1087,8 @@ // ============================================================ function updateRecordBtnUI() { recordBtn.classList.toggle('recording', recording); - recordBtn.classList.toggle('busy', busy && !recording); + recordBtn.classList.toggle('busy', (busy || awaitingResult) && !recording); + recordBtn.disabled = (busy || awaitingResult) && !recording; if (recording) { recordLabel.textContent = (mode === 'hold') ? L.labelHoldRec : L.labelToggleRec; } else { @@ -899,7 +1099,7 @@ // toggle 模式:click 切换 recordBtn.addEventListener('click', function () { if (mode !== 'toggle') return; - if (!authed || busy) return; + if (!authed || busy || awaitingResult) return; if (recording) stopRecording(); else startRecording(); }); @@ -928,7 +1128,7 @@ recordBtn.addEventListener('pointerdown', function (e) { if (mode !== 'hold') return; - if (!authed || busy) return; + if (!authed || busy || awaitingResult) return; e.preventDefault(); attachHoldEnd(); if (!recording) startRecording(); @@ -956,13 +1156,15 @@ } function startRecording() { - if (recording || startSent) return; + if (recording || startSent || awaitingResult) return; if (!ws || ws.readyState !== 1) { setStatus(L.connLost, 'error'); return; } // 先乐观置态,保证 iOS 在手势同步栈内 resume() recording = true; + clearRecoveryTimer(); + acquireWakeLock(); resetRemoteStreamState(); clearReadyTimer(); // 防止上一次 done 的回 ready 定时器迟到覆盖本次状态 clearWorkTimeout(); // 新一次录音开始,作废上一轮的识别兜底超时 @@ -1008,6 +1210,8 @@ setLevel(0); return; } + awaitingResult = true; + updateRecordBtnUI(); if (remoteSessionId) { wsSendJSON({ type: 'stop' }); resetRemoteStreamState(); @@ -1023,6 +1227,9 @@ function cancelRecording() { detachHoldEnd(); + awaitingResult = false; + clearRecoveryTimer(); + saveRecoverySession(''); if (!recording && !startSent) { teardownAudioCapture(); resetRemoteStreamState(); @@ -1079,6 +1286,11 @@ return Promise.reject(new Error('UNSUPPORTED:浏览器不支持录音,请升级或换浏览器')); } audioCtx = new AC(); + audioCtx.onstatechange = function () { + if (audioCtx && recording && startSent && audioCtx.state !== 'running') { + interruptRecording(); + } + }; } // 注意:iOS Safari 来电/Siri 后 ctx 处于私有的 'interrupted' 状态,只判 'suspended' @@ -1108,6 +1320,9 @@ return null; // 交给下一步判空直接放弃 } mediaStream = stream; + stream.getTracks().forEach(function (track) { + track.onended = function () { if (recording) interruptRecording(); }; + }); return stream; }); }) @@ -1304,6 +1519,7 @@ function failRecording(message, notifyBackend) { var waitingForAck = startSent && !remoteSessionId; recording = false; + awaitingResult = false; detachHoldEnd(); teardownAudioCapture(); if (notifyBackend && startSent) wsSendJSON({ type: 'cancel' }); @@ -1338,7 +1554,7 @@ return true; } - function handleStarted(sessionId) { + function handleStarted(sessionId, key) { if (!startSent) { clearPendingPcm(); return; @@ -1356,9 +1572,10 @@ } remoteSessionId = sessionId; + saveRecoverySession(sessionId, key); remoteSequence = 0; if (!flushPendingPcm()) { - failRecording(L.connLost, true); + interruptRecording(); return; } if (finishAfterStarted === 'stop') { @@ -1377,11 +1594,11 @@ function sendAudio(buf) { if (!recording || !buf || !buf.byteLength) return; if (!ws || ws.readyState !== 1) { - failRecording(L.connLost, false); + interruptRecording(); return; } if (remoteSessionId) { - if (!sendRemoteFrame(buf)) failRecording(L.connLost, true); + if (!sendRemoteFrame(buf)) interruptRecording(); else updateLocalLevel(buf); return; } @@ -1434,6 +1651,8 @@ // ============================================================ // 仅停止"采集/推流"(断开节点),保留 audioCtx & mediaStream 以便快速重启。 function teardownAudioCapture() { + releaseWakeLock(); + if (wakeLockHint) wakeLockHint.textContent = L.wakeLockHint; try { if (workletNode) { workletNode.port.onmessage = null; workletNode.disconnect(); } } catch (e) {} workletNode = null; @@ -1495,13 +1714,41 @@ } // ============================================================ - // 页面可见性:切后台时若在 hold 录音则取消,避免半截音频 + // 息屏和切后台结束本段录音,保留电脑已收到的部分。 // ============================================================ + function interruptRecording() { + var hadStarted = startSent; + if (recording) stopRecording(); + else if (startSent && remoteSessionId) { + wsSendJSON({ type: 'stop' }); + resetRemoteStreamState(); + awaitingResult = true; + teardownAudioCapture(); + updateRecordBtnUI(); + } + // 系统中断后释放旧轨道,下一次由用户开始录音时重新获取麦克风。 + teardownAudio(); + if (hadStarted) setStatus(L.interrupted, 'work'); + } document.addEventListener('visibilitychange', function () { - if (document.hidden && recording) { - cancelRecording(); + if (document.hidden) { + if (recording) interruptRecording(); + releaseWakeLock(); + clearRecoveryTimer(); + clearWorkTimeout(); + } else { + if (recording) acquireWakeLock(); + if (authed) requestRecovery(); + else if (!ws || ws.readyState > 1) { + var pin = readPin(); + if (pin) connect(pin); + } } }); + window.addEventListener('pagehide', function () { + if (recording) interruptRecording(); + releaseWakeLock(); + }); // ============================================================ // 初始化 @@ -1531,6 +1778,7 @@ applyStaticI18n(); syncModeUI(); initInsertSwitch(); + initWakeLockSwitch(); showScreen('pin'); showPinError(''); // 上次成功的配对码 → 自动填充并重连,刷新/重开页面免再输一次 diff --git a/openless-all/app/src-tauri/src/remote_server/assets/index.html b/openless-all/app/src-tauri/src/remote_server/assets/index.html index 14f6a12f1..c19ca850f 100644 --- a/openless-all/app/src-tauri/src/remote_server/assets/index.html +++ b/openless-all/app/src-tauri/src/remote_server/assets/index.html @@ -86,6 +86,13 @@

OpenLess 远程输入

点击大按钮开始录音,再次点击结束并识别。

+ +

+