Files
sundynix-agentix/sundynix-gateway/internal/voice/tts.go
T
Blizzard 6fe0a58f1b feat(voice): 火山双向流 TTS 客户端 + 下行接线(Phase 1 嘴巴)
回答 token 流 → 攒句器 → 火山双向 TTS → 音频帧回推客户端,端到端连续朗读打通。

- tts_frame.go: V3 事件族帧编解码(事件号+会话ID+gzip),与 ASR 简帧不同族;
  从官方参考实现核实 generate_header/parse_response 字节布局;3 解析单测
- tts.go: seed-tts-2.0 双向流客户端 StartTTS(握手ConnectionStarted/SessionStarted)
  /Speak(逐句TaskRequest)/Finish/Audio()/Close,PCM 24k;新版 API Key 鉴权
- voice_tts.go: speak() 先订阅token流再建TTS(core NATS无持久,握手期攒句入pending
  就绪补吐,不丢开头);音频泵首帧ServerSpeaking、收尾ServerTTSEnd;打断stopTTS
- voice.go: barge_in→stopTTS;会话结束连带停TTS;final转写→go speak(taskID)

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

190 lines
5.9 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 voice
import (
"context"
"encoding/json"
"fmt"
"net/http"
"sync"
"time"
"github.com/gorilla/websocket"
)
// 火山「双向流式 TTS V3」客户端(seed-tts-2.0)。边推文字边收音频,配 LLM token 流做连续朗读。
// 帧协议见 tts_frame.go;鉴权走**新版 API Key**Authorization: Bearer <APIKey>)。
//
// 一轮朗读的生命周期:StartTTS(连接+StartSession) → Speak(逐句推文字/事件200) → Finish(FinishSession)
// → 客户端 range Audio() 收音频直到 channel 关闭 → Close()。一个 TTSSession = 一次 Agent 回答。
const (
ttsEndpoint = "wss://openspeech.bytedance.com/api/v3/tts/bidirection"
TTSSampleRate = 24000 // 双向 TTS 回 PCM 24kHz 单声道(客户端按此播放)
)
// TTSSession 是一路双向 TTS 会话。Speak 推文字、Audio 出 PCM、Finish 收尾、Close 关闭。
type TTSSession struct {
conn *websocket.Conn
sessionID string
audio chan []byte
closed chan struct{}
closeOnce sync.Once
mu sync.Mutex
writeMu sync.Mutex // 串行化对火山连接的写(Speak 与 Finish 可能并发)
err error
}
// buildTTSStartSession 组 StartSession 的 req_params:音色 + 音频参数(PCM 24k)。
// namespace=BidirectionalTTS 是双向流式 TTS 的服务命名空间。
func buildTTSStartSession(cfg Config) []byte {
req := map[string]any{
"user": map[string]any{"uid": "sundynix"},
"namespace": "BidirectionalTTS",
"req_params": map[string]any{
"speaker": cfg.TTSVoiceType,
"audio_params": map[string]any{
"format": "pcm",
"sample_rate": TTSSampleRate,
},
},
}
b, _ := json.Marshal(req)
return b
}
// StartTTS 连火山双向 TTS,握手(StartConnection→StartSession)成功后返回会话。
func StartTTS(ctx context.Context, cfg Config) (*TTSSession, error) {
if !cfg.TTSEnabled() {
return nil, fmt.Errorf("TTS 未配置")
}
hdr := http.Header{}
hdr.Set("Authorization", "Bearer "+cfg.APIKey) // 新版 API Key 鉴权
hdr.Set("X-Api-Resource-Id", cfg.TTSResourceID)
hdr.Set("X-Api-Connect-Id", newConnectID())
dialer := websocket.Dialer{HandshakeTimeout: 10 * time.Second}
conn, _, err := dialer.DialContext(ctx, ttsEndpoint, hdr)
if err != nil {
return nil, fmt.Errorf("连接火山 TTS 失败: %w", err)
}
sid := newConnectID()
// StartConnection → 期待 ConnectionStarted。
if err := conn.WriteMessage(websocket.BinaryMessage, connEventFrame(evStartConnection, []byte("{}"))); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("发送 StartConnection 失败: %w", err)
}
if err := expectTTSEvent(conn, evConnectionStarted); err != nil {
_ = conn.Close()
return nil, err
}
// StartSession → 期待 SessionStarted。
if err := conn.WriteMessage(websocket.BinaryMessage, sessionEventFrame(evStartSession, sid, buildTTSStartSession(cfg))); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("发送 StartSession 失败: %w", err)
}
if err := expectTTSEvent(conn, evSessionStarted); err != nil {
_ = conn.Close()
return nil, err
}
s := &TTSSession{
conn: conn, sessionID: sid,
audio: make(chan []byte, 64), closed: make(chan struct{}),
}
go s.readLoop()
return s, nil
}
// expectTTSEvent 同步读一帧,校验是期望的事件(握手阶段用,此时还没起 readLoop)。
func expectTTSEvent(conn *websocket.Conn, want int32) error {
_ = conn.SetReadDeadline(time.Now().Add(10 * time.Second))
_, data, err := conn.ReadMessage()
conn.SetReadDeadline(time.Time{}) // 清除
if err != nil {
return fmt.Errorf("读 TTS 握手响应失败: %w", err)
}
r, perr := parseTTSFrame(data)
if perr != nil {
return fmt.Errorf("解析 TTS 握手响应失败: %w", perr)
}
if r.MsgType == ttsServerErr {
return fmt.Errorf("TTS 握手被拒 code=%d: %s", r.Code, string(r.Payload))
}
if r.Event != want && r.Event != 0 { // 0=未带事件号的容错
return fmt.Errorf("TTS 握手期待事件 %d,得 %d", want, r.Event)
}
return nil
}
// Speak 推一段文字给 TTSTaskRequest / 事件 200)。可多次调用逐句推。
func (s *TTSSession) Speak(text string) error {
payload, _ := json.Marshal(map[string]any{"text": text})
s.writeMu.Lock()
defer s.writeMu.Unlock()
return s.conn.WriteMessage(websocket.BinaryMessage, sessionEventFrame(evTaskRequest, s.sessionID, payload))
}
// Finish 发 FinishSession,告知本轮文字推完;服务端把剩余音频吐完后回 SessionFinishedAudio 随之关闭)。
func (s *TTSSession) Finish() error {
s.writeMu.Lock()
defer s.writeMu.Unlock()
return s.conn.WriteMessage(websocket.BinaryMessage, sessionEventFrame(evFinishSession, s.sessionID, []byte("{}")))
}
// Audio 返回下行音频流(PCM 24k);会话结束/出错时关闭。
func (s *TTSSession) Audio() <-chan []byte { return s.audio }
// Err 返回会话错误(Audio 关闭后读取)。
func (s *TTSSession) Err() error {
s.mu.Lock()
defer s.mu.Unlock()
return s.err
}
func (s *TTSSession) setErr(err error) {
s.mu.Lock()
if s.err == nil {
s.err = err
}
s.mu.Unlock()
}
// Close 关闭底层连接(读 goroutine 随之退出,Audio 关闭)。幂等。
func (s *TTSSession) Close() {
s.closeOnce.Do(func() {
close(s.closed)
_ = s.conn.Close()
})
}
func (s *TTSSession) readLoop() {
defer close(s.audio)
for {
_, data, err := s.conn.ReadMessage()
if err != nil {
return // 连接关闭 / 读错误
}
r, perr := parseTTSFrame(data)
if perr != nil {
continue
}
switch {
case r.MsgType == ttsServerErr:
s.setErr(fmt.Errorf("火山 TTS 错误 code=%d: %s", r.Code, string(r.Payload)))
return
case r.Event == evSessionFailed:
s.setErr(fmt.Errorf("火山 TTS 会话失败: %s", string(r.Payload)))
return
case r.Event == evSessionFinished || r.Event == evConnectionFinished:
return // 本轮朗读完毕
case r.IsAudio && len(r.Payload) > 0:
select {
case s.audio <- r.Payload:
case <-s.closed:
return
}
}
}
}