Files
sundynix-agentix/sundynix-gateway/internal/handler/voice.go
T
Blizzard 9ae06229b8 feat(voice): 显式 end 兜底提交 + 端到端模拟工具,全链路真机跑通
火山流式 ASR 只在 VAD 静音时发 Final;客户端点"停"(ClientEnd)不能干等——否则整段说完
却因没触发 VAD Final 而永不提交任务(实测复现)。改为:
- voice.go 累计 latestText(每条转写更新);ClientEnd 后 1.2s 用最新转写兜底提交;
  turnMu+submitted 保证 Final 与 ClientEnd 两路只提交一次(替换旧 lastFinal 去重)
- cmd/voicesim: 端到端模拟(免麦)——TTS 合成问话→灌网关语音WS→ASR转写→提交任务→
  大模型回答→TTS朗读回推,问答音频各存 wav,签发测试用户 JWT(auth.Issue)

真机验证:问"你是谁?你能做什么?"→ task_bab1c6ee → JARVIS 语音回答 18.6s(sim_answer.wav)。
麦克风音频→ASR→任务→大模型→TTS 全程走网关跑通。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-22 11:08:17 +08:00

236 lines
7.7 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package handler
import (
"context"
"encoding/json"
"log"
"net/http"
"strings"
"sync"
"time"
"github.com/gin-gonic/gin"
"github.com/gorilla/websocket"
"github.com/sundynix/sundynix-gateway/internal/voice"
)
// 语音交互 WebSocket 端点(JARVIS)。一条连接承载上行音频 + 下行转写 + 下行 TTS 音频,
// 协议见 voice/protocol.go。鉴权走 AuthFromHeaderOrQueryEventSource/WS 带不了 Bearer 头,
// 用 ?token=)。本文件是会话外壳 + 客户端↔网关协议循环;火山 ASR/TTS 客户端在下一步接入。
var voiceUpgrader = websocket.Upgrader{
ReadBufferSize: 4096,
WriteBufferSize: 4096,
// CheckOrigin 放行:鉴权已由 token 把关(跨源 WS 无法读响应,且我们不依赖 cookie)。
CheckOrigin: func(*http.Request) bool { return true },
}
const voiceWriteWait = 10 * time.Second
// VoiceStream: GET /api/v1/voice/stream —— 升级为 WebSocket 语音会话。
func (h *Handler) VoiceStream(c *gin.Context) {
uid := userID(c)
if uid == "" {
c.JSON(http.StatusUnauthorized, gin.H{"error": "需要登录"})
return
}
cfg := h.loadVoiceConfig(c.Request.Context())
if !cfg.ASREnabled() {
c.JSON(http.StatusServiceUnavailable, gin.H{"error": "语音服务未配置(缺 API Key / ASR resource-id"})
return
}
conn, err := voiceUpgrader.Upgrade(c.Writer, c.Request, nil)
if err != nil {
log.Printf("[voice] 升级 WS 失败 uid=%s: %v", uid, err)
return
}
defer conn.Close()
// 租户/会话在升级时(还握着 gin.Context)一并抓取,供 WS 读循环里提交任务复用共用关卡。
sess := &voiceSession{
conn: conn, uid: uid, cfg: cfg, h: h,
tenantID: tenantID(c), sessionID: sessionID(c),
}
sess.send(voice.ServerMsg{Type: voice.ServerReady})
sess.run()
sess.stopASR() // 连接结束,收掉在跑的识别会话
sess.stopTTS() // 连带停掉在朗读的下行 TTS
}
// voiceSession 是一次语音会话的外壳:持 WS 连接,跑协议循环。
// 上行 = 音频→ASR→转写→提交任务;下行(token流→攒句→TTS→音频)将在 TTS 步接上。
type voiceSession struct {
conn *websocket.Conn
h *Handler // 复用 preflightCore/launchCore 提交任务
uid string
tenantID string // 升级时抓取(读循环里无 gin.Context
sessionID string
cfg voice.Config
writeMu sync.Mutex // gorilla WS 不允许并发写:读循环与 ASR 结果 goroutine 都会 send,须串行化
asr *voice.ASRSession
asrCancel context.CancelFunc
ttsMu sync.Mutex // 护住当前下行 TTS 会话指针(打断/收尾从别的 goroutine 访问)
tts *voice.TTSSession
ttsCancel context.CancelFunc
pendingGraph string // 客户端 start 时带的画布编排图(语音触发既有编排),空则按转写现组
turnMu sync.Mutex // 护住一轮的转写累计 + 提交去重(ASR 结果 goroutine 与 ClientEnd 兜底 goroutine 都访问)
latestText string // 本轮最近一次转写(部分/最终);ClientEnd 时兜底用它提交
submitted bool // 本轮是否已提交——Final 与 ClientEnd 两条路径只落一次
}
// send 下发一条控制/事件消息(文本帧,JSON)。并发安全。
func (s *voiceSession) send(m voice.ServerMsg) {
b, _ := json.Marshal(m)
s.writeMu.Lock()
defer s.writeMu.Unlock()
_ = s.conn.SetWriteDeadline(time.Now().Add(voiceWriteWait))
if err := s.conn.WriteMessage(websocket.TextMessage, b); err != nil {
log.Printf("[voice] 写控制消息失败 uid=%s: %v", s.uid, err)
}
}
// sendAudio 下发一帧 TTS 音频(二进制帧)。并发安全。
func (s *voiceSession) sendAudio(pcm []byte) {
s.writeMu.Lock()
defer s.writeMu.Unlock()
_ = s.conn.SetWriteDeadline(time.Now().Add(voiceWriteWait))
if err := s.conn.WriteMessage(websocket.BinaryMessage, pcm); err != nil {
log.Printf("[voice] 写音频失败 uid=%s: %v", s.uid, err)
}
}
// run 是协议读循环:二进制帧=上行音频,文本帧=控制消息。
func (s *voiceSession) run() {
for {
mt, data, err := s.conn.ReadMessage()
if err != nil {
return // 客户端断开 / 读错误
}
switch mt {
case websocket.BinaryMessage:
s.onAudio(data)
case websocket.TextMessage:
var m voice.ClientMsg
if json.Unmarshal(data, &m) != nil {
continue
}
if s.onControl(m) {
return // bye
}
}
}
}
// onAudio 收到一帧上行音频 → 喂火山 ASR。
func (s *voiceSession) onAudio(pcm []byte) {
if s.asr == nil {
s.startASR() // 客户端没显式 start 就直接说话时,惰性开一路识别
}
if s.asr != nil {
if err := s.asr.PushAudio(pcm); err != nil {
log.Printf("[voice] 喂 ASR 音频失败 uid=%s: %v", s.uid, err)
}
}
}
// onControl 处理客户端控制消息,返回 true 表示会话应结束。
func (s *voiceSession) onControl(m voice.ClientMsg) (done bool) {
switch m.Type {
case voice.ClientBye:
return true
case voice.ClientStart:
s.pendingGraph = m.Graph // 客户端画布图(可空):本轮若有转写则语音触发它跑
s.turnMu.Lock()
s.latestText, s.submitted = "", false // 新一轮:清累计与提交标记
s.turnMu.Unlock()
s.stopASR()
s.startASR() // 新一轮:重开识别
case voice.ClientEnd:
if s.asr != nil {
_ = s.asr.Finish() // 告知火山本轮说完
}
// 火山流式 ASR 只在 VAD 静音时才发 Final;客户端显式 end(点停)时不能干等——
// 给点收尾时间让末尾部分结果到齐,再用"最新转写"兜底提交(trySubmit 去重,Final 先到就它先提交)。
go func() {
time.Sleep(1200 * time.Millisecond)
s.turnMu.Lock()
txt := s.latestText
s.turnMu.Unlock()
s.trySubmit(txt)
}()
case voice.ClientBargeIn:
s.stopTTS() // 打断:用户又开口,立刻掐掉正在朗读的 TTS
}
return false
}
// startASR 开一路火山流式识别,并起 goroutine 把转写实时回推客户端。
func (s *voiceSession) startASR() {
ctx, cancel := context.WithCancel(context.Background())
asr, err := voice.StartASR(ctx, s.cfg, s.uid)
if err != nil {
cancel()
log.Printf("[voice] 启动 ASR 失败 uid=%s: %v", s.uid, err)
s.send(voice.ServerMsg{Type: voice.ServerError, Msg: "语音识别启动失败"})
return
}
s.asr = asr
s.asrCancel = cancel
go func() {
for r := range asr.Results() {
if r.Err != nil {
return // 识别流结束/出错
}
s.send(voice.ServerMsg{Type: voice.ServerTranscript, Text: r.Text, Final: r.Final})
if t := strings.TrimSpace(r.Text); t != "" {
s.turnMu.Lock()
s.latestText = r.Text // 累计最新转写,供 ClientEnd 兜底提交
s.turnMu.Unlock()
}
if r.Final {
s.trySubmit(r.Text) // VAD 检出句末 → 直接提交(与 ClientEnd 兜底二选一,去重)
}
}
}()
}
// trySubmit 本轮提交一次任务:Final 与 ClientEnd 兜底两条路径抢先,submitted 保证只落一次。
func (s *voiceSession) trySubmit(text string) {
txt := strings.TrimSpace(text)
s.turnMu.Lock()
if txt == "" || s.submitted {
s.turnMu.Unlock()
return
}
s.submitted = true
s.turnMu.Unlock()
taskID, err := s.submitVoiceTask(txt, s.pendingGraph)
if err != nil {
s.send(voice.ServerMsg{Type: voice.ServerError, Msg: "任务提交失败:" + err.Error()})
return
}
s.pendingGraph = "" // 画布图一次性消费,避免后续转写重复触发同图
s.send(voice.ServerMsg{Type: voice.ServerTask, TaskID: taskID})
// 下行:订阅该任务 token 流 → 攒句 → TTS → 音频帧回推。独立 goroutine 跑,不堵 ASR 结果流。
go s.speak(taskID)
}
// stopASR 收掉当前识别会话(幂等)。
func (s *voiceSession) stopASR() {
if s.asr != nil {
s.asr.Close()
s.asr = nil
}
if s.asrCancel != nil {
s.asrCancel()
s.asrCancel = nil
}
}