diff --git a/desktop/src-tauri/src/dictation.rs b/desktop/src-tauri/src/dictation.rs index 3ca48ef..cd36af6 100644 --- a/desktop/src-tauri/src/dictation.rs +++ b/desktop/src-tauri/src/dictation.rs @@ -127,11 +127,13 @@ pub fn stop(app: &AppHandle, canceled: bool) { let release_at = Instant::now(); tauri::async_runtime::spawn(async move { if !canceled { - // 事件驱动:收到 stop 后的尾部 final(或 usage 结算帧)即触发注入; - // 350ms 仅作超时兜底,避免 flush 慢于固定 sleep 截丢尾 final(18C)。 + // 事件驱动:收到 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(350)) => {} + _ = tokio::time::sleep(Duration::from_millis(800)) => {} } let buf = app.state::(); diff --git a/server/cmd/gummycheck/main.go b/server/cmd/gummycheck/main.go index ffcc788..c53db0d 100644 --- a/server/cmd/gummycheck/main.go +++ b/server/cmd/gummycheck/main.go @@ -20,6 +20,9 @@ import ( func main() { wavPath := flag.String("wav", "", "16kHz/16bit/mono wav 文件路径") + model := flag.String("model", "gummy-realtime-v1", "DashScope 实时识别模型") + paceMs := flag.Int("pace", 50, "每 100ms 帧的推流间隔 ms(50=2x 加速,100=实时)") + idleSec := flag.Int("idle", 0, "会话建立后先闲置 N 秒再推流(测预连接会话的存活时间)") flag.Parse() if *wavPath == "" { fmt.Fprintln(os.Stderr, "用法: rbw get dashscope-api-key | gummycheck -wav test.wav") @@ -41,6 +44,8 @@ func main() { fmt.Printf("音频 %d 字节 ≈ %.1fs\n", len(pcm), float64(len(pcm))/32000) p := asr.NewGummy(key) + p.Model = *model + fmt.Printf("模型 %s · 推流间隔 %dms/帧\n", *model, *paceMs) ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) defer cancel() @@ -51,6 +56,11 @@ func main() { os.Exit(1) } fmt.Printf("会话建立 %v\n", time.Since(start).Round(time.Millisecond)) + if *idleSec > 0 { + fmt.Printf("闲置 %ds...\n", *idleSec) + time.Sleep(time.Duration(*idleSec) * time.Second) + start = time.Now() // 重新计时,测闲置后的识别延迟 + } done := make(chan struct{}) var firstPartial time.Duration @@ -79,11 +89,11 @@ func main() { fmt.Fprintln(os.Stderr, "SendAudio 失败:", err) os.Exit(1) } - time.Sleep(50 * time.Millisecond) // 2x 实时推流 + time.Sleep(time.Duration(*paceMs) * time.Millisecond) } _ = sess.Close() <-done - fmt.Printf("首个结果延迟 %v(含 2x 推流速度因素)\n", firstPartial.Round(time.Millisecond)) + fmt.Printf("首个结果延迟 %v(推流 %dms/帧)\n", firstPartial.Round(time.Millisecond), *paceMs) } // readWavData 提取 wav 的 data chunk(仅支持 PCM)。 diff --git a/server/internal/asr/gummy.go b/server/internal/asr/gummy.go index 2fb7cdc..63b79be 100644 --- a/server/internal/asr/gummy.go +++ b/server/internal/asr/gummy.go @@ -22,8 +22,25 @@ const dashscopeWS = "wss://dashscope.aliyuncs.com/api-ws/v1/inference" type GummyProvider struct { APIKey string Model string // 默认 gummy-realtime-v1 + + // spare 预连接(12B/#24):后台常备一条已完成 run-task 握手的 DashScope + // 会话,start 直接取用,把 dial+task-started(实测暖路径 ~130ms、抖动可达 + // 2.5~3.7s)从 first_partial 关键路径上移除。协议固定 16kHz,spare 可互换。 + mu sync.Mutex + spare *spareEntry + dialing bool + maintainOne sync.Once } +type spareEntry struct { + sess *gummySession + born time.Time +} + +// spareMaxAge spare 的可用年龄上限:DashScope 空闲 60s 断连(实测 60s 存活、 +// 120s 报 Idle timeout),留余量在 45s 内取用、40s 定期换新。 +const spareMaxAge = 45 * time.Second + func NewGummy(apiKey string) *GummyProvider { return &GummyProvider{APIKey: apiKey, Model: "gummy-realtime-v1"} } @@ -31,6 +48,87 @@ func NewGummy(apiKey string) *GummyProvider { func (g *GummyProvider) Name() string { return "gummy" } func (g *GummyProvider) StartSession(ctx context.Context, cfg SessionConfig) (Session, error) { + g.maintainOne.Do(func() { go g.maintainSpare() }) + if cfg.SampleRate == 16000 { + if s := g.takeSpare(); s != nil { + go g.replenish() + return s, nil + } + defer func() { go g.replenish() }() + } + return g.dialSession(ctx, cfg) +} + +// takeSpare 取走当前 spare(若仍新鲜且存活);过期/已死的就地关闭丢弃。 +func (g *GummyProvider) takeSpare() *gummySession { + g.mu.Lock() + sp := g.spare + g.spare = nil + g.mu.Unlock() + if sp == nil { + return nil + } + if time.Since(sp.born) < spareMaxAge && sp.sess.alive() { + slog.Info("gummy: session from spare", "task", sp.sess.taskID, + "age_ms", time.Since(sp.born).Milliseconds()) + return sp.sess + } + _ = sp.sess.Close() + return nil +} + +// replenish 异步补位一条 spare;已有 spare 或正在补位时空转返回。 +func (g *GummyProvider) replenish() { + g.mu.Lock() + if g.dialing || g.spare != nil { + g.mu.Unlock() + return + } + g.dialing = true + g.mu.Unlock() + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + sess, err := g.dialSession(ctx, SessionConfig{SampleRate: 16000, SessionID: "spare"}) + cancel() + + g.mu.Lock() + g.dialing = false + if err == nil { + if g.spare == nil { + g.spare = &spareEntry{sess: sess.(*gummySession), born: time.Now()} + g.mu.Unlock() + return + } + g.mu.Unlock() + _ = sess.Close() // 竞态下已有他人补位:丢弃多余会话 + return + } + g.mu.Unlock() + slog.Warn("gummy: spare replenish failed", "err", err) +} + +// maintainSpare 每 40s 换新 spare,避免超过 DashScope 60s 空闲上限后 +// 取到已死会话、退化为冷启 dial。进程常驻,随 provider 生命周期运行。 +func (g *GummyProvider) maintainSpare() { + t := time.NewTicker(40 * time.Second) + defer t.Stop() + for range t.C { + g.mu.Lock() + sp := g.spare + stale := sp != nil && time.Since(sp.born) >= spareMaxAge-5*time.Second + if stale { + g.spare = nil + } + g.mu.Unlock() + if stale { + _ = sp.sess.Close() + } + g.replenish() + } +} + +// dialSession 建立一条全新的 DashScope 会话(原 StartSession 主体)。 +func (g *GummyProvider) dialSession(ctx context.Context, cfg SessionConfig) (Session, error) { header := map[string][]string{ "Authorization": {"bearer " + g.APIKey}, "X-DashScope-DataInspection": {"enable"}, @@ -109,6 +207,16 @@ type gummySession struct { lastText string // 当前句已下发文本,用于切分 partial/final } +// alive 会话底层连接是否仍然存活(readLoop 退出即 done 关闭)。 +func (s *gummySession) alive() bool { + select { + case <-s.done: + return false + default: + return true + } +} + func (s *gummySession) SendAudio(pcm []byte) error { s.writeMu.Lock() defer s.writeMu.Unlock() diff --git a/server/internal/gateway/gateway.go b/server/internal/gateway/gateway.go index 3979941..262705f 100644 --- a/server/internal/gateway/gateway.go +++ b/server/internal/gateway/gateway.go @@ -364,24 +364,26 @@ const finishWaitTimeout = 3 * time.Second // 注意 done/consumeMu 顺序(16D):先停 usageLoop 并 drain 在途 consumeDelta, // 再读 trialPart/balancePart 快照,避免漏记在途扣费。 func (s *session) finish(ctx context.Context, canceled bool) { + // 先冻结会话并停 usageLoop,再 flush provider(12B/#23):若 usageLoop 在 + // waitResults 期间(最长 3s)继续 tick,周期 usage 帧会抢在尾部 final 之前 + // 下发——客户端把"stop 后首个 final/usage"当提交信号,会以 partial 文本 + // 提前上屏、丢失 final 修正。冻结后 stop 之后唯一的 usage 是本函数末尾的 + // 结算帧,且必然在全部尾部 final 之后,客户端提交信号因此安全。 s.mu.Lock() if s.finished { s.mu.Unlock() return } + s.finished = true + close(s.done) s.mu.Unlock() _ = s.provider.Close() s.waitResults() // 尾部 final 全部下发后再收尾(带超时,16B) - // 先标记 finished 并停 usageLoop,再 drain 在途 consumeDelta(16D): - // 取 consumeMu 会等待任何已过 finished 检查、正等 Redis 返回的 consumeDelta - // 完成并把 trialPart/balancePart 计入,从而快照不漏账。 - s.mu.Lock() - s.finished = true - close(s.done) - s.mu.Unlock() - + // drain 在途 consumeDelta(16D):finished 已置位,取 consumeMu 会等待任何 + // 已过 finished 检查、正等 Redis 返回的 consumeDelta 完成并把 + // trialPart/balancePart 计入,从而快照不漏账。 s.consumeMu.Lock() // drain:等当前在途 consume(若有)完成 s.consumeMu.Unlock()