d22e02fa2e
#23 修复(服务端根治,五端客户端提交信号自动变安全): - finish 先冻结会话并停 usageLoop,再 flush provider——此前 usageLoop 在 waitResults(最长 3s)期间继续 tick,周期 usage 帧会抢在尾部 final 之前 下发,客户端以 partial 文本提前上屏、丢失 final 修正(12B 实测 8s 会话 4/4 复现;修复后 6/6 归零) - 桌面端收尾兜底超时 350ms → 800ms(实测 8s 会话 final flush 需 350~540ms, 350ms 会截丢尾 final) #24 预连接(12B 调优): - GummyProvider 常备一条已完成 run-task 握手的 spare 会话,start 直取, 后台异步补位 + 40s 定期换新(DashScope 空闲 60s 断连,实测 60s 存活/ 120s Idle timeout,留余量 45s) - 实测网关 start 处理 120~670ms → 3~4ms;消除 dial 抖动(实测 60ms~3.7s) 对 first_partial 尾部的放大 gummycheck 增加 -model / -pace / -idle 参数,支持模型对比与闲置存活实验 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
330 lines
12 KiB
Rust
330 lines
12 KiB
Rust
//! 听写编排(push-to-talk 核心):
|
||
//! 快捷键按下 → 浮层弹出(光标附近,不抢焦点)→ 采集推流;
|
||
//! 松开 → stop → 等尾部 final → 注入 committed 文本 → 浮层淡出。
|
||
|
||
use parking_lot::Mutex;
|
||
use serde_json::{json, Value};
|
||
use std::sync::atomic::{AtomicU64, Ordering};
|
||
use std::sync::mpsc;
|
||
use std::time::{Duration, Instant};
|
||
use tauri::{AppHandle, Emitter, EventTarget, Manager};
|
||
use tokio::sync::Notify;
|
||
|
||
/// 浮层窗口 label(仅 overlay 监听 asr/hotkey 事件,见 tauri.conf.json)。
|
||
const OVERLAY: &str = "overlay";
|
||
|
||
/// 向 overlay 窗口定向 emit(18G):避免广播给全部 webview。
|
||
fn emit_overlay(app: &AppHandle, event: &str, payload: Value) {
|
||
let _ = app.emit_to(EventTarget::webview_window(OVERLAY), event, payload);
|
||
}
|
||
|
||
#[derive(Default)]
|
||
pub struct DictationState {
|
||
inner: Mutex<Option<Session>>,
|
||
}
|
||
|
||
struct Session {
|
||
id: String,
|
||
/// 会话代际(18A/18B):每次 start 自增,on_server_msg / 收尾任务校验后才生效。
|
||
epoch: u64,
|
||
capture: Option<crate::audio::Capture>,
|
||
started: Instant,
|
||
got_first_partial: bool,
|
||
}
|
||
|
||
pub fn start(app: &AppHandle) {
|
||
let state = app.state::<DictationState>();
|
||
if state.inner.lock().is_some() {
|
||
return; // 已在录音
|
||
}
|
||
if crate::api::is_app_paused(app) {
|
||
return; // 已暂停使用
|
||
}
|
||
let session_id = uuid::Uuid::new_v4().to_string();
|
||
|
||
// 开新会话代际:递增 epoch,重置 buffer / 丢弃标志 / 收尾信号(18A/18B/18C)。
|
||
let buf = app.state::<CommitBuffer>();
|
||
let epoch = buf.begin_session();
|
||
|
||
set_tray_tooltip(app, "dudu — 录音中");
|
||
emit_overlay(app, "hotkey", json!({"state": "down"}));
|
||
show_overlay(app);
|
||
|
||
let ws = (*app.state::<crate::ws::WsHandle>()).clone();
|
||
let _ = ws.send(crate::ws::WsCmd::Start {
|
||
session_id: session_id.clone(),
|
||
});
|
||
|
||
// 音频帧经 std channel → 转发 tokio channel
|
||
let (tx, rx) = mpsc::channel::<Vec<u8>>();
|
||
let mic = app.state::<crate::settings::SettingsStore>().get().mic;
|
||
let t0 = Instant::now();
|
||
let capture = match crate::audio::start(&mic, tx) {
|
||
Ok(c) => {
|
||
let m = app.state::<crate::metrics::Metrics>();
|
||
m.record("audio.start_ms", json!({"ms": t0.elapsed().as_millis() as i64}));
|
||
Some(c)
|
||
}
|
||
Err(e) => {
|
||
log::error!("audio start failed: {e}");
|
||
emit_overlay(app, "asr", json!({"type":"error","code":"AUDIO","message": e}));
|
||
None
|
||
}
|
||
};
|
||
{
|
||
let ws = ws.clone();
|
||
std::thread::spawn(move || {
|
||
while let Ok(frame) = rx.recv() {
|
||
if ws.send(crate::ws::WsCmd::Audio(frame)).is_err() {
|
||
break;
|
||
}
|
||
}
|
||
});
|
||
}
|
||
|
||
*state.inner.lock() = Some(Session {
|
||
id: session_id,
|
||
epoch,
|
||
capture,
|
||
started: Instant::now(),
|
||
got_first_partial: false,
|
||
});
|
||
}
|
||
|
||
/// 当前是否有进行中的录音会话(18E:ws 重连后据此判断是否需中止会话)。
|
||
pub fn is_recording(app: &AppHandle) -> bool {
|
||
app.state::<DictationState>().inner.lock().is_some()
|
||
}
|
||
|
||
pub fn stop(app: &AppHandle, canceled: bool) {
|
||
let state = app.state::<DictationState>();
|
||
let Some(mut sess) = state.inner.lock().take() else {
|
||
return;
|
||
};
|
||
let epoch = sess.epoch;
|
||
emit_overlay(app, "hotkey", json!({"state": "up"}));
|
||
set_tray_tooltip(app, "dudu — 就绪");
|
||
sess.capture.take(); // 停止采集
|
||
|
||
let buf = app.state::<CommitBuffer>();
|
||
if canceled {
|
||
// 取消即丢弃:标记该 epoch,后续晚到 final 直接丢弃不入 buffer(18A)。
|
||
buf.discard(epoch);
|
||
}
|
||
|
||
let ws = (*app.state::<crate::ws::WsHandle>()).clone();
|
||
let cmd = if canceled {
|
||
crate::ws::WsCmd::Cancel { session_id: sess.id.clone() }
|
||
} else {
|
||
crate::ws::WsCmd::Stop { session_id: sess.id.clone() }
|
||
};
|
||
let _ = ws.send(cmd);
|
||
|
||
// 取尾部 final 完成信号(事件驱动注入,18C)。
|
||
let finalized = buf.finalize_signal();
|
||
|
||
let app = app.clone();
|
||
let release_at = Instant::now();
|
||
tauri::async_runtime::spawn(async move {
|
||
if !canceled {
|
||
// 事件驱动:收到 stop 后的尾部 final(或 usage 结算帧)即触发注入(18C)。
|
||
// 服务端已保证 stop 后不再下发周期 usage 帧(结算帧必在全部尾部 final
|
||
// 之后,#23),二者均可安全作为提交信号。800ms 仅作超时兜底:
|
||
// 12B 实测 8s 会话的 final flush 需 350~540ms,350ms 兜底会截丢尾 final。
|
||
tokio::select! {
|
||
_ = finalized.notified() => {}
|
||
_ = tokio::time::sleep(Duration::from_millis(800)) => {}
|
||
}
|
||
|
||
let buf = app.state::<CommitBuffer>();
|
||
// 会话代际校验:本任务只处理自己的 epoch,快速连说时旧任务醒来不串新会话(18B)。
|
||
let Some(committed) = buf.take_for_epoch(epoch) else {
|
||
// epoch 不匹配(已被新会话覆盖)或已被丢弃:不注入、不 hide(18A/18B)。
|
||
return;
|
||
};
|
||
if !committed.is_empty() {
|
||
if !crate::inject::accessibility_ok() {
|
||
// 无辅助功能权限(macOS):注入必失败 → 浮层引导去授权(10E)
|
||
emit_overlay(
|
||
&app,
|
||
"asr",
|
||
json!({"type":"error","code":"NO_ACCESSIBILITY","message":"需要辅助功能权限"}),
|
||
);
|
||
tokio::time::sleep(Duration::from_millis(2400)).await; // 让用户看清引导
|
||
} else {
|
||
let t = committed.clone();
|
||
let injected =
|
||
tauri::async_runtime::spawn_blocking(move || crate::inject::inject_text(&t)).await;
|
||
match injected {
|
||
Ok(Ok(())) => {
|
||
let m = app.state::<crate::metrics::Metrics>();
|
||
m.record(
|
||
"asr.release_to_commit_ms",
|
||
json!({"ms": release_at.elapsed().as_millis() as i64}),
|
||
);
|
||
}
|
||
_ => {
|
||
log::error!("inject failed: {injected:?}");
|
||
if !crate::inject::accessibility_ok() {
|
||
emit_overlay(
|
||
&app,
|
||
"asr",
|
||
json!({"type":"error","code":"NO_ACCESSIBILITY","message":"需要辅助功能权限"}),
|
||
);
|
||
tokio::time::sleep(Duration::from_millis(2400)).await;
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
} else {
|
||
// 取消:清空 buffer(仅当仍是本 epoch,避免误清新会话),不上屏(18A)。
|
||
let buf = app.state::<CommitBuffer>();
|
||
let _ = buf.take_for_epoch(epoch);
|
||
}
|
||
hide_overlay(&app);
|
||
});
|
||
}
|
||
|
||
/// CommitBuffer 跨 ws 任务与停止流程共享的"待注入文本"+ 会话代际/丢弃标志/收尾信号。
|
||
pub struct CommitBuffer {
|
||
inner: Mutex<BufferInner>,
|
||
/// 当前会话代际(18A/18B):start 自增,on_server_msg / 收尾任务据此校验归属。
|
||
epoch: AtomicU64,
|
||
/// 尾部 final 到达信号,驱动收尾注入(18C)。
|
||
finalized: std::sync::Arc<Notify>,
|
||
}
|
||
|
||
struct BufferInner {
|
||
text: (String, String), // (final 累计, 最新 partial)
|
||
/// text 归属的会话代际。
|
||
text_epoch: u64,
|
||
/// 被取消的会话代际:该 epoch 的后续 final 直接丢弃(18A)。
|
||
discarded_epoch: Option<u64>,
|
||
}
|
||
|
||
impl Default for CommitBuffer {
|
||
fn default() -> Self {
|
||
Self {
|
||
inner: Mutex::new(BufferInner {
|
||
text: (String::new(), String::new()),
|
||
text_epoch: 0,
|
||
discarded_epoch: None,
|
||
}),
|
||
epoch: AtomicU64::new(0),
|
||
finalized: std::sync::Arc::new(Notify::new()),
|
||
}
|
||
}
|
||
}
|
||
|
||
impl CommitBuffer {
|
||
/// 开新会话:epoch+1,清空文本与丢弃标志,返回新 epoch(18A/18B)。
|
||
fn begin_session(&self) -> u64 {
|
||
let epoch = self.epoch.fetch_add(1, Ordering::SeqCst) + 1;
|
||
let mut g = self.inner.lock();
|
||
g.text = (String::new(), String::new());
|
||
g.text_epoch = epoch;
|
||
g.discarded_epoch = None;
|
||
epoch
|
||
}
|
||
|
||
fn current_epoch(&self) -> u64 {
|
||
self.epoch.load(Ordering::SeqCst)
|
||
}
|
||
|
||
/// 标记某会话被取消,其后续 final 一律丢弃(18A)。
|
||
fn discard(&self, epoch: u64) {
|
||
self.inner.lock().discarded_epoch = Some(epoch);
|
||
}
|
||
|
||
/// 收尾信号克隆,供 stop 任务 await(18C)。
|
||
fn finalize_signal(&self) -> std::sync::Arc<Notify> {
|
||
self.finalized.clone()
|
||
}
|
||
|
||
/// 取出并清空文本——仅当 buffer 仍属指定 epoch 且未被丢弃时返回(18A/18B)。
|
||
/// epoch 不匹配(已被新会话覆盖)返回 None,调用方据此放弃注入/hide。
|
||
fn take_for_epoch(&self, epoch: u64) -> Option<String> {
|
||
let mut g = self.inner.lock();
|
||
if g.text_epoch != epoch {
|
||
return None; // 已被新会话覆盖
|
||
}
|
||
if g.discarded_epoch == Some(epoch) {
|
||
g.text = (String::new(), String::new());
|
||
return Some(String::new()); // 本会话已取消,返回空(不注入但允许 hide 自身浮层)
|
||
}
|
||
let committed = format!("{}{}", g.text.0, g.text.1);
|
||
g.text = (String::new(), String::new());
|
||
Some(committed)
|
||
}
|
||
}
|
||
|
||
/// ws 下行回调:维护文本缓冲 + 首字延迟打点(在 ws.rs 收包处调用)。
|
||
pub fn on_server_msg(app: &AppHandle, v: &Value) {
|
||
let buf = app.state::<CommitBuffer>();
|
||
let cur = buf.current_epoch();
|
||
match v.get("type").and_then(|t| t.as_str()) {
|
||
Some("partial") => {
|
||
if let Some(text) = v.get("text").and_then(|t| t.as_str()) {
|
||
let mut g = buf.inner.lock();
|
||
// 仅当属于当前会话且未被丢弃才写入(18A/18B)。
|
||
if g.text_epoch == cur && g.discarded_epoch != Some(cur) {
|
||
g.text.1 = text.to_string();
|
||
}
|
||
}
|
||
let state = app.state::<DictationState>();
|
||
let mut guard = state.inner.lock();
|
||
if let Some(sess) = guard.as_mut() {
|
||
if !sess.got_first_partial {
|
||
sess.got_first_partial = true;
|
||
let m = app.state::<crate::metrics::Metrics>();
|
||
m.record(
|
||
"asr.first_partial_ms",
|
||
json!({"ms": sess.started.elapsed().as_millis() as i64}),
|
||
);
|
||
}
|
||
}
|
||
}
|
||
Some("final") => {
|
||
if let Some(text) = v.get("text").and_then(|t| t.as_str()) {
|
||
let mut g = buf.inner.lock();
|
||
// 校验会话:不属当前会话或已被取消 → 丢弃,不写 buffer(18A/18B)。
|
||
if g.text_epoch == cur && g.discarded_epoch != Some(cur) {
|
||
g.text.0.push_str(text);
|
||
g.text.1.clear();
|
||
}
|
||
}
|
||
// 尾部 final 到达 → 唤醒收尾注入任务(18C)。
|
||
buf.finalized.notify_waiters();
|
||
}
|
||
// usage 结算帧(参考 iOS):作为收尾兜底信号,立即触发注入(18C)。
|
||
Some("usage") => {
|
||
buf.finalized.notify_waiters();
|
||
}
|
||
_ => {}
|
||
}
|
||
}
|
||
|
||
/// 托盘提示按录音状态切换(11B)。
|
||
fn set_tray_tooltip(app: &AppHandle, text: &str) {
|
||
if let Some(tray) = app.tray_by_id("main") {
|
||
let _ = tray.set_tooltip(Some(text));
|
||
}
|
||
}
|
||
|
||
fn show_overlay(app: &AppHandle) {
|
||
if let Some(w) = app.get_webview_window(OVERLAY) {
|
||
// 跟随光标:浮层出现在指针下方 24px,水平居中
|
||
if let Ok(pos) = app.cursor_position() {
|
||
let _ = w.set_position(tauri::PhysicalPosition::new(pos.x - 180.0, pos.y + 24.0));
|
||
}
|
||
let _ = w.show();
|
||
}
|
||
}
|
||
|
||
fn hide_overlay(app: &AppHandle) {
|
||
if let Some(w) = app.get_webview_window(OVERLAY) {
|
||
let _ = w.hide();
|
||
}
|
||
}
|