From ec747838955256b46dd4e15ef182aa0e25553423 Mon Sep 17 00:00:00 2001 From: a692570 Date: Fri, 11 Sep 2026 12:57:40 -0700 Subject: [PATCH 1/2] Add Telnyx STT provider --- README.md | 2 + README.zh-CN.md | 2 + config.toml.example | 8 +- docs/configuration.md | 6 +- docs/configuration.zh-CN.md | 6 +- docs/providers.md | 24 +++- docs/providers.zh-CN.md | 24 +++- internal/config/config.go | 14 +- internal/stt/stt.go | 7 +- internal/stt/telnyx.go | 255 ++++++++++++++++++++++++++++++++++++ internal/stt/telnyx_test.go | 144 ++++++++++++++++++++ 11 files changed, 481 insertions(+), 11 deletions(-) create mode 100644 internal/stt/telnyx.go create mode 100644 internal/stt/telnyx_test.go diff --git a/README.md b/README.md index f9f8166..5ead634 100644 --- a/README.md +++ b/README.md @@ -117,6 +117,8 @@ Providers: Deepgram, AssemblyAI, OpenAI, Cartesia, ElevenLabs, MiniMax, Speechif OpenAI STT supports `whisper-1`, `gpt-4o-transcribe`, and `gpt-4o-mini-transcribe` through the independent `openai.stt_model` setting. +Telnyx STT fronts a dozen engines behind one key (`telnyx.stt_engine`); the in-house `Telnyx` engine is finals-only, so barge-in and live captions need a hosted engine such as `Deepgram`. + ## Documentation | Page | What's in it | diff --git a/README.zh-CN.md b/README.zh-CN.md index 44ecf3c..109003d 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -117,6 +117,8 @@ StreamCore 位于「提示词 + 工具」类框架的下一层:媒体链路。 OpenAI STT 可通过独立的 `openai.stt_model` 配置选择 `whisper-1`、`gpt-4o-transcribe` 或 `gpt-4o-mini-transcribe`。 +Telnyx STT 用一个 key 前置十余种引擎(`telnyx.stt_engine`);自研 `Telnyx` 引擎只出最终结果,打断(barge-in)与实时字幕需要 hosted 引擎(如 `Deepgram`)。 + ## 文档 | 页面 | 内容 | diff --git a/config.toml.example b/config.toml.example index 4a97e8b..beeb3b7 100644 --- a/config.toml.example +++ b/config.toml.example @@ -39,7 +39,7 @@ turn_merge_ms = 350 # Debounce window for merging finals into one turn. provider = "" # Supported: grok (empty = classic pipeline) [stt] -provider = "deepgram" # Supported: aliyun, assemblyai, deepgram, openai, vibevoice, volcengine +provider = "deepgram" # Supported: aliyun, assemblyai, deepgram, openai, telnyx, vibevoice, volcengine [llm] provider = "openai" # Supported: openai, ollama, agent @@ -143,11 +143,15 @@ voice_id = "" # Optional; defaults to Geffen (geffen_32). Simba 3.2 us model = "" # Optional; defaults to simba-3.2 [telnyx] -api_key = "" # Required if tts.provider = "telnyx" +api_key = "" # Required if tts.provider = "telnyx" or stt.provider = "telnyx" voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # Any catalog voice from GET /v2/text-to-speech/voices # (Qwen3TTS voices use UUID ids). Availability varies by account; a voice your # key is not provisioned for fails the dial with HTTP 403 voice_speed = 1.0 # Playback-rate multiplier, clamped to 0.8-1.2 like the per-utterance delivery tags +stt_engine = "Telnyx" # STT. Transcription engine when stt.provider = "telnyx": the in-house recognizer + # or a hosted engine (Deepgram, AssemblyAI, Azure, Cohere, Google, Humain, + # Parakeet, Reson8, Soniox, Speechmatics, xAI). Case-sensitive. In-house is + # finals-only, so barge-in and live captions need a hosted engine [mimo] api_key = "" # Required if tts.provider = "mimo" diff --git a/docs/configuration.md b/docs/configuration.md index fa90135..4a82606 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -38,7 +38,7 @@ turn_merge_ms = 350 # Debounce window for merging finals into o provider = "" # "grok", or empty for the classic pipeline [stt] -provider = "deepgram" # aliyun | assemblyai | deepgram | openai | vibevoice | volcengine +provider = "deepgram" # aliyun | assemblyai | deepgram | openai | telnyx | vibevoice | volcengine [llm] provider = "openai" # openai | ollama | agent @@ -119,10 +119,12 @@ api_key = "" voice_id = "" model = "" -[telnyx] # Telnyx hosted synthesis, used when tts.provider = "telnyx" +[telnyx] # Telnyx hosted speech, used when tts.provider = "telnyx" or stt.provider = "telnyx"; one key covers both roles api_key = "" voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # Any catalog voice from GET /v2/text-to-speech/voices; availability varies by account voice_speed = 1.0 # Playback-rate multiplier, clamped to 0.8-1.2 +stt_engine = "Telnyx" # STT engine: the in-house recognizer or a hosted one (Deepgram, AssemblyAI, Azure, ...). Case-sensitive. + # In-house is finals-only — no barge-in or live captions; hosted engines stream interims and both work [minimax] api_key = "" diff --git a/docs/configuration.zh-CN.md b/docs/configuration.zh-CN.md index 7c42098..a2ceb25 100644 --- a/docs/configuration.zh-CN.md +++ b/docs/configuration.zh-CN.md @@ -32,7 +32,7 @@ turn_merge_ms = 350 # Debounce window for merging finals into o provider = "" # "grok", or empty for the classic pipeline [stt] -provider = "deepgram" # aliyun | assemblyai | deepgram | openai | vibevoice | volcengine +provider = "deepgram" # aliyun | assemblyai | deepgram | openai | telnyx | vibevoice | volcengine [llm] provider = "openai" # openai | ollama | agent @@ -113,10 +113,12 @@ api_key = "" voice_id = "" model = "" -[telnyx] # Telnyx 托管合成,当 tts.provider = "telnyx" 时使用 +[telnyx] # Telnyx 托管语音,当 tts.provider = "telnyx" 或 stt.provider = "telnyx" 时使用;一个 key 覆盖两个方向 api_key = "" voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # GET /v2/text-to-speech/voices 目录中的任意音色;可用性因账号而异 voice_speed = 1.0 # 播放速率倍数,限制在 0.8-1.2 +stt_engine = "Telnyx" # STT 引擎:自研识别器或托管引擎(Deepgram、AssemblyAI、Azure 等)。大小写敏感。 + # 自研引擎只出最终结果 —— 无打断(barge-in)、无实时字幕;托管引擎流式输出中间结果,两者均可用 [minimax] api_key = "" diff --git a/docs/providers.md b/docs/providers.md index 4532ba0..bb7c832 100644 --- a/docs/providers.md +++ b/docs/providers.md @@ -4,7 +4,7 @@ | Role | Providers | Required credentials | |------|-----------|----------------------| -| STT | `aliyun`, `assemblyai`, `deepgram`, `openai`, `vibevoice`, `volcengine` | Matching provider API key, or a local VibeVoice ASR server | +| STT | `aliyun`, `assemblyai`, `deepgram`, `openai`, `telnyx`, `vibevoice`, `volcengine` | Matching provider API key, or a local VibeVoice ASR server | | LLM | `openai`, `ollama`, `agent` | OpenAI API key, an Ollama instance you control, or your own HTTP agent endpoint | | TTS | `cartesia`, `deepgram`, `elevenlabs`, `mimo`, `minimax`, `speechify`, `telnyx`, `vibevoice` | Matching provider API key, or a local VibeVoice TTS server | | Speech-to-speech | `grok` | xAI API key — replaces STT, LLM, and TTS together | @@ -21,6 +21,7 @@ Notes: - `tts.provider = "mimo"` is Xiaomi's MiMo TTS, with Chinese and English voices and optional voice cloning on the paid models. - `stt.provider = "aliyun"` is Alibaba Cloud Model Studio (DashScope) streaming ASR; `vocabulary_id` biases it toward domain terms. - `stt.provider = "volcengine"` is Doubao streaming ASR — useful where Deepgram is slow to reach or its Mandarin is not good enough. The console gives a free hourly allowance. +- `stt.provider = "telnyx"` fronts Telnyx's in-house and a dozen hosted transcription engines over one WebSocket and one key. The in-house engine is finals-only, so barge-in and live captions need a hosted engine. See [Telnyx STT](#telnyx-stt). - `realtime.provider = "grok"` switches to speech-to-speech and ignores `[stt]`, `[llm]`, and `[tts]` entirely. Every key and knob lives in the [configuration reference](./configuration.md). @@ -148,6 +149,27 @@ Three things to know: Delivery tags map onto `voice_speed` (clamped to 0.8–1.2, the same conversational band as Cartesia), and `voice_speed` in config sets the baseline pace for untagged sentences. +## Telnyx STT + +Streaming transcription over the Telnyx speech-to-text WebSocket: raw linear16 binary frames in, JSON transcript frames out, at the pipeline's native 16 kHz mono so nothing resamples. The same `[telnyx]` section and API key as TTS cover both roles; `stt_engine` picks the recognizer. + +```toml +[stt] +provider = "telnyx" + +[telnyx] +api_key = "" +stt_engine = "Telnyx" # or a hosted engine — "Deepgram", "AssemblyAI", "Azure", ... Casing matters +``` + +Three things to know: + +- **The in-house `Telnyx` engine is finals-only.** It emits exactly one final frame after the caller stops speaking — no interims, no timestamps, no confidence. Barge-in and live captions depend on interim results (the pipeline gates interruption on partial text), so **neither works with this engine**: interruptions never fire, and the client transcript shows nothing until the final lands. Use it for final-transcript-only deployments. See the [design discussion](https://github.com/streamcoreai/streamcore-server/issues/75) for the trade-offs. +- **Hosted engines restore interims.** The same endpoint fronts AssemblyAI, Azure, Cohere, Deepgram, Google, Humain, Parakeet, Reson8, Soniox, Speechmatics, and xAI. Any engine other than `Telnyx` is requested with `interim_results=true`, partials stream as they do from any other provider, and barge-in and live captions work. `stt_engine = "Deepgram"` is the verified configuration. +- **The engine name is case-sensitive.** `telnyx` is rejected with a structured error frame that lists the supported engines; the valid values are exactly the ones in that list, first letter capitalised. + +Confidence arrives as `null` from the in-house engine and a 0–1 float from hosted ones; the pipeline treats `null` as unknown rather than low. + ## Local VibeVoice setup VibeVoice provides fully local STT and TTS with no API keys, using [VibeVoice-ASR](https://huggingface.co/mlx-community/VibeVoice-ASR-4bit) for recognition and [VibeVoice-Realtime-0.5B](https://huggingface.co/mlx-community/VibeVoice-Realtime-0.5B-6bit) for synthesis via two lightweight Python sidecars. On Apple Silicon they use [mlx-audio](https://github.com/Blaizzy/mlx-audio) (MLX); on Linux/Windows they fall back to PyTorch automatically. diff --git a/docs/providers.zh-CN.md b/docs/providers.zh-CN.md index 222ba37..19e2654 100644 --- a/docs/providers.zh-CN.md +++ b/docs/providers.zh-CN.md @@ -4,7 +4,7 @@ | 角色 | 服务商 | 所需凭据 | |------|-----------|----------------------| -| STT | `aliyun`、`assemblyai`、`deepgram`、`openai`、`vibevoice`、`volcengine` | 对应服务商的 API key,或一个本地 VibeVoice ASR 服务 | +| STT | `aliyun`、`assemblyai`、`deepgram`、`openai`、`telnyx`、`vibevoice`、`volcengine` | 对应服务商的 API key,或一个本地 VibeVoice ASR 服务 | | LLM | `openai`、`ollama`、`agent` | OpenAI API key、你自己掌控的 Ollama 实例,或你自己的 HTTP 智能体端点 | | TTS | `cartesia`、`deepgram`、`elevenlabs`、`mimo`、`minimax`、`speechify`、`telnyx`、`vibevoice` | 对应服务商的 API key,或一个本地 VibeVoice TTS 服务 | | 语音到语音 | `grok` | xAI API key —— 一并取代 STT、LLM 与 TTS | @@ -21,6 +21,7 @@ - `tts.provider = "mimo"` 是小米 MiMo TTS,中英文音色齐备,付费模型还支持声音克隆。 - `stt.provider = "aliyun"` 是阿里云百炼(DashScope)流式 ASR;`vocabulary_id` 可以把模型往你的领域词上带。 - `stt.provider = "volcengine"` 是豆包流式 ASR —— 适合 Deepgram 访问慢、或它的中文识别不够好的场景。控制台有免费时长可以先试。 +- `stt.provider = "telnyx"` 通过一条 WebSocket、一个 key 前置自研与十余种托管转写引擎。自研引擎只出最终结果,打断与实时字幕需要托管引擎。见 [Telnyx STT](#telnyx-stt)。 - `realtime.provider = "grok"` 切换到语音到语音模式,并完全忽略 `[stt]`、`[llm]` 与 `[tts]`。 所有 key 与可调项都在[配置参考](./configuration.zh-CN.md)里。 @@ -148,6 +149,27 @@ voice_speed = 1.0 表达标签映射到 `voice_speed`(限制在 0.8–1.2,与 Cartesia 相同的对话档位),配置里的 `voice_speed` 则是未打标签句子的基准语速。 +## Telnyx STT + +走 Telnyx 语音转文字 WebSocket 的流式识别:上行是裸 linear16 二进制帧,下行是 JSON 转写帧,均为流水线原生的 16 kHz 单声道,音频路径无需任何重采样。与 TTS 共用同一个 `[telnyx]` 配置段和 API key;`stt_engine` 选择识别引擎。 + +```toml +[stt] +provider = "telnyx" + +[telnyx] +api_key = "" +stt_engine = "Telnyx" # 或托管引擎 —— "Deepgram"、"AssemblyAI"、"Azure" 等。大小写敏感 +``` + +有三件事必须弄对: + +- **自研 `Telnyx` 引擎只出最终结果。** 它在来电者停止说话后才发出唯一一帧 final —— 没有中间结果、没有时间戳、没有置信度。打断(barge-in)与实时字幕依赖中间结果(流水线以部分文本来判定打断),因此**这两项在该引擎下不工作**:打断永远不会触发,客户端字幕也要等到 final 落地才有内容。只适合「仅需最终转写」的场景。取舍讨论见[设计讨论](https://github.com/streamcoreai/streamcore-server/issues/75)。 +- **托管引擎可以恢复中间结果。** 同一个端点前置 AssemblyAI、Azure、Cohere、Deepgram、Google、Humain、Parakeet、Reson8、Soniox、Speechmatics 与 xAI。除 `Telnyx` 以外的任何引擎都会带上 `interim_results=true` 请求,中间结果像其他服务商一样流式到达,打断与实时字幕正常工作。已验证的配置是 `stt_engine = "Deepgram"`。 +- **引擎名大小写敏感。** `telnyx` 会被一帧结构化错误拒绝,错误里列出支持的引擎;合法取值就是那个列表,首字母大写。 + +置信度:自研引擎返回 `null`,托管引擎返回 0–1 浮点数;流水线把 `null` 视为未知而不是低置信。 + ## 本地 VibeVoice 配置 VibeVoice 提供完全本地、无需 API key 的 STT 与 TTS:识别用 [VibeVoice-ASR](https://huggingface.co/mlx-community/VibeVoice-ASR-4bit),合成用 [VibeVoice-Realtime-0.5B](https://huggingface.co/mlx-community/VibeVoice-Realtime-0.5B-6bit),通过两个轻量 Python 边车进程运行。在 Apple Silicon 上使用 [mlx-audio](https://github.com/Blaizzy/mlx-audio)(MLX);在 Linux/Windows 上自动回退到 PyTorch。 diff --git a/internal/config/config.go b/internal/config/config.go index 699597b..606de99 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -293,8 +293,8 @@ type SpeechifyConfig struct { Model string `toml:"model"` } -// TelnyxConfig configures Telnyx hosted speech synthesis, used when -// tts.provider = "telnyx". +// TelnyxConfig configures Telnyx hosted speech, used when tts.provider = +// "telnyx" and when stt.provider = "telnyx". One API key covers both roles. type TelnyxConfig struct { APIKey string `toml:"api_key"` // Voice is a catalog voice from GET /v2/text-to-speech/voices. Empty @@ -305,6 +305,15 @@ type TelnyxConfig struct { // VoiceSpeed is a playback-rate multiplier (1.0 = normal), clamped to // the 0.8-1.2 conversational band. Zero defaults to 1.0. VoiceSpeed float64 `toml:"voice_speed"` + // STTEngine selects the transcription engine used when stt.provider = + // "telnyx": the in-house "Telnyx" recognizer or any of the hosted + // engines the endpoint fronts (AssemblyAI, Azure, Cohere, Deepgram, + // Google, Humain, Parakeet, Reson8, Soniox, Speechmatics, xAI). The + // value is case-sensitive and sent verbatim. The in-house engine emits + // one final per connection and no interims, so barge-in and live + // captions do not work with it; hosted engines stream interims and both + // work. Empty defaults to "Telnyx". + STTEngine string `toml:"stt_engine"` } type MiMoConfig struct { @@ -452,6 +461,7 @@ func Load(path string) (*Config, error) { if cfg.Telnyx.VoiceSpeed == 0 { cfg.Telnyx.VoiceSpeed = 1.0 } + setDefault(&cfg.Telnyx.STTEngine, "Telnyx") setDefault(&cfg.LLM.Provider, "openai") setDefault(&cfg.TTS.Provider, "cartesia") setDefault(&cfg.OpenAI.Model, "gpt-4o-mini") diff --git a/internal/stt/stt.go b/internal/stt/stt.go index 7dc41ba..0186e63 100644 --- a/internal/stt/stt.go +++ b/internal/stt/stt.go @@ -51,9 +51,14 @@ func NewClient(ctx context.Context, cfg *config.Config, onResult func(Transcript return nil, fmt.Errorf("stt provider %q requires [volcengine] api_key to be set", cfg.STT.Provider) } return NewVolcengineClient(ctx, cfg.Volcengine, onResult) + case "telnyx": + if cfg.Telnyx.APIKey == "" { + return nil, fmt.Errorf("stt provider %q requires [telnyx] api_key to be set", cfg.STT.Provider) + } + return NewTelnyxClient(ctx, cfg.Telnyx, onResult) case "vibevoice": return NewVibeVoiceClient(ctx, cfg.VibeVoice.ASRURL, onResult) default: - return nil, fmt.Errorf("unknown stt provider %q (supported: aliyun, assemblyai, deepgram, openai, vibevoice, volcengine)", cfg.STT.Provider) + return nil, fmt.Errorf("unknown stt provider %q (supported: aliyun, assemblyai, deepgram, openai, telnyx, vibevoice, volcengine)", cfg.STT.Provider) } } diff --git a/internal/stt/telnyx.go b/internal/stt/telnyx.go new file mode 100644 index 0000000..93c2c29 --- /dev/null +++ b/internal/stt/telnyx.go @@ -0,0 +1,255 @@ +package stt + +import ( + "context" + "encoding/json" + "fmt" + "io" + "log" + "net/http" + "net/url" + "strings" + "sync" + "time" + + "github.com/gorilla/websocket" + + "github.com/streamcoreai/streamcore-server/internal/config" +) + +const ( + // telnyxSTTURL is the streaming transcription WebSocket. One endpoint + // fronts every engine; the engine rides the transcription_engine query + // parameter. + telnyxSTTURL = "wss://api.telnyx.com/v2/speech-to-text/transcription" + // telnyxInHouseEngine is Telnyx's own recognizer, reached only by + // spelling the name exactly: the parameter is case-sensitive, and a + // lowercase "telnyx" is rejected with a structured error frame listing + // the supported engines. + telnyxInHouseEngine = "Telnyx" + // telnyxFinalSpentReason is what SendAudio reports once the in-house + // engine's single final has been delivered. Saying "connection closed" + // here would read as a network failure when it is in fact the documented + // one-final-per-connection contract. + telnyxFinalSpentReason = "the in-house Telnyx engine emits one final per connection and this connection's final has been delivered; set stt_engine to a streaming engine (e.g. \"Deepgram\") for continuous transcription" +) + +// TelnyxClient implements the Client interface against Telnyx's streaming +// speech-to-text WebSocket. +// +// The wire protocol is raw linear16 PCM upstream (binary frames at the +// pipeline's native 16 kHz mono, so nothing resamples) and JSON transcript +// frames downstream. One endpoint fronts a dozen engines — Telnyx's in-house +// recognizer plus hosted Deepgram, AssemblyAI, Azure, and others — selected +// by the transcription_engine query parameter, so one [telnyx] section and +// one API key cover both STT and TTS. +// +// The engines differ in what they stream, and the difference is not cosmetic: +// +// - The in-house engine emits exactly one final frame after the caller +// stops speaking — no interims, no timestamps, confidence null — and +// then holds the socket open indefinitely; the server never closes it. +// The client therefore tears the connection down itself once that final +// is delivered, and the next SendAudio fails with the reason above +// rather than pretending to accept audio. Barge-in and live captions +// hard-depend on interim results (internal/pipeline/inbound.go gates +// interruption on partial text), so this engine suits +// final-transcript-only deployments; see the design discussion in +// issue #75. +// - Every other engine streams partials when asked, so the client appends +// interim_results=true for them and barge-in works as with any other +// provider. On the in-house engine the parameter is silently ignored — +// the final still arrives — so it is simply not sent there. +type TelnyxClient struct { + conn *websocket.Conn + cancel context.CancelFunc + + // closeOnFinal tears the connection down after the first final, the + // in-house engine's one-shot contract above. + closeOnFinal bool + + // writeMu serialises writes: gorilla/websocket allows only one concurrent + // writer, and Close races SendAudio when a call ends mid-utterance. + writeMu sync.Mutex + closed bool + closeReason string +} + +// NewTelnyxClient dials the Telnyx transcription endpoint and starts a +// goroutine that forwards transcripts to onResult. Empty stt_engine defaults +// to the in-house Telnyx engine. +func NewTelnyxClient(ctx context.Context, cfg config.TelnyxConfig, onResult func(TranscriptResult)) (Client, error) { + sttCtx, cancel := context.WithCancel(ctx) + + engine := strings.TrimSpace(cfg.STTEngine) + if engine == "" { + engine = telnyxInHouseEngine + } + + header := http.Header{} + header.Set("Authorization", "Bearer "+cfg.APIKey) + + dialer := websocket.Dialer{HandshakeTimeout: 10 * time.Second} + conn, resp, err := dialer.DialContext(sttCtx, telnyxSTTEndpoint(engine), header) + if err != nil { + cancel() + if resp != nil { + // The body carries the actual reason — a rejected key reads as a + // bare 401 — which the status code alone does not convey. + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + resp.Body.Close() + return nil, fmt.Errorf("telnyx stt dial: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) + } + return nil, fmt.Errorf("telnyx stt dial: %w", err) + } + + c := &TelnyxClient{ + conn: conn, + cancel: cancel, + closeOnFinal: engine == telnyxInHouseEngine, + } + + log.Printf("[stt] connected to Telnyx STT (engine %s)", engine) + go c.readLoop(sttCtx, onResult) + + return c, nil +} + +// telnyxSTTEndpoint builds the dial URL. The engine value is URL-encoded +// because it reaches the server as a query parameter, and interim_results is +// appended for every engine except the in-house one: only the hosted engines +// emit partials, and the pipeline's barge-in hard-depends on them. +func telnyxSTTEndpoint(engine string) string { + endpoint := fmt.Sprintf("%s?transcription_engine=%s&input_format=linear16&sample_rate=16000", + telnyxSTTURL, url.QueryEscape(engine)) + if engine != telnyxInHouseEngine { + endpoint += "&interim_results=true" + } + return endpoint +} + +// telnyxSTTFrame is one transcript message. The in-house engine reports +// confidence as null; hosted engines report a 0-1 float. +type telnyxSTTFrame struct { + Transcript string `json:"transcript"` + // Confidence is a pointer so null — the in-house engine — stays + // distinguishable from a genuine 0.0. Null maps to 0, which + // TranscriptResult defines as "unknown", not "low". + Confidence *float64 `json:"confidence"` + IsFinal bool `json:"is_final"` + // Errors carries a rejected request (an unsupported engine, say) instead + // of a transcript; the server closes the connection right after. + Errors []telnyxSTTError `json:"errors"` +} + +type telnyxSTTError struct { + Code string `json:"code"` + Title string `json:"title"` + Source struct { + Parameter string `json:"parameter"` + } `json:"source"` + Detail string `json:"detail"` +} + +// routeTelnyxSTTFrame classifies one server frame: emit reports whether the +// frame carries a transcript the pipeline should see, and err is set for +// error frames and undecodable messages. Frames with no transcript text — +// utterance-end markers, keepalives — come back with emit false so the +// caller skips them; emitting one would surface an empty final turn. +func routeTelnyxSTTFrame(msg []byte) (result TranscriptResult, emit bool, err error) { + var f telnyxSTTFrame + if jerr := json.Unmarshal(msg, &f); jerr != nil { + return TranscriptResult{}, false, fmt.Errorf("telnyx stt frame decode: %w", jerr) + } + if len(f.Errors) > 0 { + // The detail names the actual problem ("Unsupported + // transcription_engine 'telnyx'. Supported engines: …"), which beats + // the code and title for an operator reading the log. It is the + // first entry: the service reports one error per frame. + e := f.Errors[0] + detail := e.Detail + if detail == "" { + detail = e.Title + } + return TranscriptResult{}, false, fmt.Errorf("telnyx stt: %s", detail) + } + text := strings.TrimSpace(f.Transcript) + if text == "" { + return TranscriptResult{}, false, nil + } + var confidence float64 + if f.Confidence != nil { + confidence = *f.Confidence + } + return TranscriptResult{Text: text, IsFinal: f.IsFinal, Confidence: confidence}, true, nil +} + +// readLoop forwards transcript frames to onResult until the connection ends. +func (c *TelnyxClient) readLoop(ctx context.Context, onResult func(TranscriptResult)) { + for { + _, msg, err := c.conn.ReadMessage() + if err != nil { + if ctx.Err() == nil { + log.Printf("[stt] telnyx read: %v", err) + } + return + } + + result, emit, rerr := routeTelnyxSTTFrame(msg) + if rerr != nil { + // The server closes the connection right after an error frame; + // there is nothing left to read. + log.Printf("[stt] telnyx: %v", rerr) + return + } + if !emit { + continue + } + if result.IsFinal { + log.Printf("[stt] final: %q", result.Text) + } + onResult(result) + + // The in-house engine's contract is one final per connection, and + // the server holds the socket open afterwards — it never closes. + // Tear it down from this side so the spent session does not linger + // until call teardown, and let SendAudio report why. + if result.IsFinal && c.closeOnFinal { + c.shutdown(telnyxFinalSpentReason) + return + } + } +} + +// SendAudio forwards a chunk of PCM. The pipeline calls this every 20ms. +func (c *TelnyxClient) SendAudio(data []byte) error { + c.writeMu.Lock() + defer c.writeMu.Unlock() + if c.closed { + return fmt.Errorf("telnyx stt: %s", c.closeReason) + } + return c.conn.WriteMessage(websocket.BinaryMessage, data) +} + +// Close tears down the connection. There is no flush to wait for: the +// service settles an utterance on its own silence detection, so a +// client-side close mid-utterance simply drops the stream. +func (c *TelnyxClient) Close() { + c.shutdown("connection closed") +} + +// shutdown closes the connection once, recording why so SendAudio can tell a +// spent finals-only session apart from a normal teardown. +func (c *TelnyxClient) shutdown(reason string) { + c.writeMu.Lock() + if c.closed { + c.writeMu.Unlock() + return + } + c.closed = true + c.closeReason = reason + c.writeMu.Unlock() + + _ = c.conn.Close() + c.cancel() +} diff --git a/internal/stt/telnyx_test.go b/internal/stt/telnyx_test.go new file mode 100644 index 0000000..0e67b32 --- /dev/null +++ b/internal/stt/telnyx_test.go @@ -0,0 +1,144 @@ +package stt + +import ( + "strings" + "testing" +) + +// The in-house engine's final frame carries confidence null. A nil pointer +// must map to 0, which TranscriptResult defines as "unknown" — a naive +// dereference would panic, and a misread null would look like a failed +// recognition. +func TestTelnyxSTTFinalFrameMapsNullConfidence(t *testing.T) { + frame := `{"transcript":" The quick brown fox jumps over the lazy dog, testing one, two, three.","confidence":null,"is_final":true}` + + result, emit, err := routeTelnyxSTTFrame([]byte(frame)) + if err != nil { + t.Fatalf("final frame rejected: %v", err) + } + if !emit { + t.Fatal("final frame not emitted") + } + if !result.IsFinal { + t.Error("final frame reported partial, want final") + } + if result.Text != "The quick brown fox jumps over the lazy dog, testing one, two, three." { + t.Errorf("text = %q, want the leading space trimmed", result.Text) + } + if result.Confidence != 0 { + t.Errorf("confidence = %v, want 0 (null maps to unknown)", result.Confidence) + } +} + +// Hosted engines stream partials with a real confidence. Dropping the partial +// path would kill barge-in, which gates interruption on partial text. +func TestTelnyxSTTRoutesPartialFrame(t *testing.T) { + frame := `{"transcript":"The quick brown","confidence":0.8442383,"is_final":false,"speech_final":false}` + + result, emit, err := routeTelnyxSTTFrame([]byte(frame)) + if err != nil { + t.Fatalf("partial frame rejected: %v", err) + } + if !emit { + t.Fatal("partial frame not emitted") + } + if result.IsFinal { + t.Error("partial frame reported final, want partial") + } + if result.Text != "The quick brown" { + t.Errorf("text = %q, want %q", result.Text, "The quick brown") + } + if result.Confidence != 0.8442383 { + t.Errorf("confidence = %v, want 0.8442383", result.Confidence) + } +} + +// End-of-utterance markers arrive with an empty transcript. Emitting one +// would surface an empty final turn and re-answer the sentence before it. +func TestTelnyxSTTSkipsEmptyTranscript(t *testing.T) { + _, emit, err := routeTelnyxSTTFrame([]byte(`{"transcript":"","is_final":true,"utterance_end":true}`)) + if err != nil { + t.Fatalf("utterance-end frame rejected: %v", err) + } + if emit { + t.Error("utterance-end frame emitted, want it skipped") + } +} + +// An unsupported engine is reported in a structured errors array, and the +// detail names the actual problem including the supported list. Surfacing +// only the code or title would leave an operator with "Invalid Parameter" +// and no idea what to fix. +func TestTelnyxSTTSurfacesErrorDetail(t *testing.T) { + frame := `{"errors":[{"code":"40007","title":"Invalid Parameter","source":{"parameter":"transcription_engine"},"detail":"Unsupported transcription_engine 'telnyx'. Supported engines: AssemblyAI, Azure, Cohere, Deepgram, Google, Humain, Parakeet, Reson8, Soniox, Speechmatics, Telnyx, xAI"}]}` + + _, _, err := routeTelnyxSTTFrame([]byte(frame)) + if err == nil { + t.Fatal("error frame not surfaced") + } + if !strings.Contains(err.Error(), "Unsupported transcription_engine 'telnyx'") { + t.Errorf("error = %v, want it to carry the server detail", err) + } + if !strings.Contains(err.Error(), "Telnyx") { + t.Errorf("error = %v, want the supported-engine list to survive", err) + } +} + +// A frame that is not JSON means the connection is misbehaving; continuing +// would silently drop whatever transcript followed. +func TestTelnyxSTTRejectsMalformedJSON(t *testing.T) { + if _, _, err := routeTelnyxSTTFrame([]byte(`{not json`)); err == nil { + t.Fatal("malformed frame accepted, want a decode error") + } +} + +// The engine rides the dial URL as a query parameter, so it must be encoded; +// and interim_results decides whether barge-in works at all. The in-house +// engine ignores the parameter, so it is not sent there; every hosted engine +// streams partials only when it is present. +func TestTelnyxSTTEndpointBuildsQuery(t *testing.T) { + cases := []struct { + name string + engine string + want string + }{ + { + "in-house engine, no interims", + "Telnyx", + telnyxSTTURL + "?transcription_engine=Telnyx&input_format=linear16&sample_rate=16000", + }, + { + "hosted engine gets interims", + "Deepgram", + telnyxSTTURL + "?transcription_engine=Deepgram&input_format=linear16&sample_rate=16000&interim_results=true", + }, + { + "engine value is URL-encoded", + "Two Words", + telnyxSTTURL + "?transcription_engine=Two+Words&input_format=linear16&sample_rate=16000&interim_results=true", + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + if got := telnyxSTTEndpoint(tc.engine); got != tc.want { + t.Errorf("endpoint = %q, want %q", got, tc.want) + } + }) + } +} + +// After the in-house engine's single final, SendAudio must refuse rather +// than pretend to accept audio, and the message must point at the fix — +// "connection closed" would read as a network failure when the connection +// is behaving exactly as documented. +func TestTelnyxSTTSendAudioReportsSpentSession(t *testing.T) { + c := &TelnyxClient{closed: true, closeReason: telnyxFinalSpentReason} + + err := c.SendAudio([]byte{0x00, 0x01}) + if err == nil { + t.Fatal("send after close accepted, want the spent-session reason") + } + if !strings.Contains(err.Error(), "stt_engine") { + t.Errorf("error = %v, want it to point at the stt_engine fix", err) + } +} From 2975ed0e162e64bef93b4ab549bd60a4da0de2b4 Mon Sep 17 00:00:00 2001 From: a692570 Date: Fri, 11 Sep 2026 13:05:51 -0700 Subject: [PATCH 2/2] Reshape Telnyx STT per review: transcription_engine, Deepgram default, per-utterance lifecycle --- README.md | 2 +- README.zh-CN.md | 2 +- config.toml.example | 9 +- docs/configuration.md | 5 +- docs/configuration.zh-CN.md | 4 +- docs/providers.md | 16 +- docs/providers.zh-CN.md | 16 +- internal/config/config.go | 20 +- internal/stt/telnyx.go | 525 +++++++++++++++++++++++++----------- internal/stt/telnyx_test.go | 357 ++++++++++++++++++++++-- 10 files changed, 750 insertions(+), 206 deletions(-) diff --git a/README.md b/README.md index 5ead634..bf30358 100644 --- a/README.md +++ b/README.md @@ -117,7 +117,7 @@ Providers: Deepgram, AssemblyAI, OpenAI, Cartesia, ElevenLabs, MiniMax, Speechif OpenAI STT supports `whisper-1`, `gpt-4o-transcribe`, and `gpt-4o-mini-transcribe` through the independent `openai.stt_model` setting. -Telnyx STT fronts a dozen engines behind one key (`telnyx.stt_engine`); the in-house `Telnyx` engine is finals-only, so barge-in and live captions need a hosted engine such as `Deepgram`. +Telnyx STT fronts a dozen engines behind one key (`telnyx.transcription_engine`, default `Deepgram` so barge-in and live captions work out of the box); the in-house `Telnyx` engine is finals-only, so both are off when it is selected, and a startup log line says so. ## Documentation diff --git a/README.zh-CN.md b/README.zh-CN.md index 109003d..10458fd 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -117,7 +117,7 @@ StreamCore 位于「提示词 + 工具」类框架的下一层:媒体链路。 OpenAI STT 可通过独立的 `openai.stt_model` 配置选择 `whisper-1`、`gpt-4o-transcribe` 或 `gpt-4o-mini-transcribe`。 -Telnyx STT 用一个 key 前置十余种引擎(`telnyx.stt_engine`);自研 `Telnyx` 引擎只出最终结果,打断(barge-in)与实时字幕需要 hosted 引擎(如 `Deepgram`)。 +Telnyx STT 用一个 key 前置十余种引擎(`telnyx.transcription_engine`,默认 `Deepgram`,打断与实时字幕开箱即用);自研 `Telnyx` 引擎只出最终结果,选中它时两者关闭,启动日志会写明。 ## 文档 diff --git a/config.toml.example b/config.toml.example index beeb3b7..2376c05 100644 --- a/config.toml.example +++ b/config.toml.example @@ -148,10 +148,11 @@ voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # Any catalog vo # (Qwen3TTS voices use UUID ids). Availability varies by account; a voice your # key is not provisioned for fails the dial with HTTP 403 voice_speed = 1.0 # Playback-rate multiplier, clamped to 0.8-1.2 like the per-utterance delivery tags -stt_engine = "Telnyx" # STT. Transcription engine when stt.provider = "telnyx": the in-house recognizer - # or a hosted engine (Deepgram, AssemblyAI, Azure, Cohere, Google, Humain, - # Parakeet, Reson8, Soniox, Speechmatics, xAI). Case-sensitive. In-house is - # finals-only, so barge-in and live captions need a hosted engine +transcription_engine = "Deepgram" # STT. Transcription engine when stt.provider = "telnyx". + # Verified values: "Deepgram" (streams partials, so barge-in and live captions + # work) or "Telnyx" (in-house, finals-only, so barge-in and live captions are + # off and a startup log line says so). Case-sensitive, sent verbatim; other + # hosted engines the endpoint fronts pass through untested [mimo] api_key = "" # Required if tts.provider = "mimo" diff --git a/docs/configuration.md b/docs/configuration.md index 4a82606..f41c7aa 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -123,8 +123,9 @@ model = "" api_key = "" voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # Any catalog voice from GET /v2/text-to-speech/voices; availability varies by account voice_speed = 1.0 # Playback-rate multiplier, clamped to 0.8-1.2 -stt_engine = "Telnyx" # STT engine: the in-house recognizer or a hosted one (Deepgram, AssemblyAI, Azure, ...). Case-sensitive. - # In-house is finals-only — no barge-in or live captions; hosted engines stream interims and both work +transcription_engine = "Deepgram" # STT engine, verified values: "Deepgram" (partials, barge-in works; the default) or + # "Telnyx" (in-house, finals-only: barge-in and live captions are off, and a startup + # log line says so). Case-sensitive, sent verbatim [minimax] api_key = "" diff --git a/docs/configuration.zh-CN.md b/docs/configuration.zh-CN.md index a2ceb25..500f0c3 100644 --- a/docs/configuration.zh-CN.md +++ b/docs/configuration.zh-CN.md @@ -117,8 +117,8 @@ model = "" api_key = "" voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # GET /v2/text-to-speech/voices 目录中的任意音色;可用性因账号而异 voice_speed = 1.0 # 播放速率倍数,限制在 0.8-1.2 -stt_engine = "Telnyx" # STT 引擎:自研识别器或托管引擎(Deepgram、AssemblyAI、Azure 等)。大小写敏感。 - # 自研引擎只出最终结果 —— 无打断(barge-in)、无实时字幕;托管引擎流式输出中间结果,两者均可用 +transcription_engine = "Deepgram" # STT 引擎,已验证取值:"Deepgram"(有中间结果,打断可用;默认值)或 + # "Telnyx"(自研,只出最终结果:打断与实时字幕关闭,启动日志会写明)。大小写敏感,原样透传 [minimax] api_key = "" diff --git a/docs/providers.md b/docs/providers.md index bb7c832..8845470 100644 --- a/docs/providers.md +++ b/docs/providers.md @@ -12,7 +12,8 @@ Notes: -- `stt.provider = "openai"` uses batch final transcription instead of streaming partials; choose `whisper-1`, `gpt-4o-transcribe`, or `gpt-4o-mini-transcribe` with `openai.stt_model`. +- `stt.provider = "openai"` uses batch final transcription instead of streaming partials, so barge-in and live captions do not work, since both depend on partials; choose `whisper-1`, `gpt-4o-transcribe`, or `gpt-4o-mini-transcribe` with `openai.stt_model`. +- `stt.provider = "telnyx"` fronts Telnyx's in-house and a dozen hosted transcription engines over one WebSocket and one key. `transcription_engine` defaults to `Deepgram`, so barge-in and live captions work out of the box; the in-house `Telnyx` engine is finals-only, so both are off when it is selected. See [Telnyx STT](#telnyx-stt). - `llm.provider = "ollama"` targets any Ollama-compatible endpoint via `base_url` — local or on your own infrastructure. - `llm.provider = "agent"` POSTs each turn to an HTTP endpoint you host; your agent owns memory, prompting, and tools, and replies stream back as SSE, chunked text, or JSON. See [Bring your own agent](./bring-your-own-agent.md). - `stt.provider = "vibevoice"` and `tts.provider = "vibevoice"` use local models; start the Python sidecars first. @@ -21,7 +22,6 @@ Notes: - `tts.provider = "mimo"` is Xiaomi's MiMo TTS, with Chinese and English voices and optional voice cloning on the paid models. - `stt.provider = "aliyun"` is Alibaba Cloud Model Studio (DashScope) streaming ASR; `vocabulary_id` biases it toward domain terms. - `stt.provider = "volcengine"` is Doubao streaming ASR — useful where Deepgram is slow to reach or its Mandarin is not good enough. The console gives a free hourly allowance. -- `stt.provider = "telnyx"` fronts Telnyx's in-house and a dozen hosted transcription engines over one WebSocket and one key. The in-house engine is finals-only, so barge-in and live captions need a hosted engine. See [Telnyx STT](#telnyx-stt). - `realtime.provider = "grok"` switches to speech-to-speech and ignores `[stt]`, `[llm]`, and `[tts]` entirely. Every key and knob lives in the [configuration reference](./configuration.md). @@ -151,7 +151,7 @@ Delivery tags map onto `voice_speed` (clamped to 0.8–1.2, the same conversatio ## Telnyx STT -Streaming transcription over the Telnyx speech-to-text WebSocket: raw linear16 binary frames in, JSON transcript frames out, at the pipeline's native 16 kHz mono so nothing resamples. The same `[telnyx]` section and API key as TTS cover both roles; `stt_engine` picks the recognizer. +Streaming transcription over the Telnyx speech-to-text WebSocket: raw linear16 binary frames in, JSON transcript frames out, at the pipeline's native 16 kHz mono so nothing resamples. The same `[telnyx]` section and API key as TTS cover both roles; `transcription_engine` picks the recognizer. ```toml [stt] @@ -159,16 +159,16 @@ provider = "telnyx" [telnyx] api_key = "" -stt_engine = "Telnyx" # or a hosted engine — "Deepgram", "AssemblyAI", "Azure", ... Casing matters +transcription_engine = "Deepgram" # verified: "Deepgram" (partial results, barge-in works) or "Telnyx" (in-house, finals-only). Case matters ``` Three things to know: -- **The in-house `Telnyx` engine is finals-only.** It emits exactly one final frame after the caller stops speaking — no interims, no timestamps, no confidence. Barge-in and live captions depend on interim results (the pipeline gates interruption on partial text), so **neither works with this engine**: interruptions never fire, and the client transcript shows nothing until the final lands. Use it for final-transcript-only deployments. See the [design discussion](https://github.com/streamcoreai/streamcore-server/issues/75) for the trade-offs. -- **Hosted engines restore interims.** The same endpoint fronts AssemblyAI, Azure, Cohere, Deepgram, Google, Humain, Parakeet, Reson8, Soniox, Speechmatics, and xAI. Any engine other than `Telnyx` is requested with `interim_results=true`, partials stream as they do from any other provider, and barge-in and live captions work. `stt_engine = "Deepgram"` is the verified configuration. -- **The engine name is case-sensitive.** `telnyx` is rejected with a structured error frame that lists the supported engines; the valid values are exactly the ones in that list, first letter capitalised. +- **`Deepgram` is the default.** It streams interim results exactly like the built-in Deepgram provider, so barge-in and live captions work out of the box. A finals-only default would silently switch both off for anyone who just sets the provider and starts talking. +- **The in-house `Telnyx` engine is finals-only.** It emits exactly one final per utterance, only after the caller stops speaking: no interims, no timestamps, confidence `null`. Barge-in and live captions depend on interim results (the pipeline gates interruption on partial text), so **neither works with this engine**: interruptions never fire, and the client transcript shows nothing until the final lands. Because the engine answers only after audio stops, the client runs its own endpointing (`internal/vad`, the same detector the pipeline uses, with the silence timeout `openai.go` uses) and opens one socket per utterance, closing it once the final lands, since the server keeps it open. A startup log line says that barge-in and live captions are off for the session, so the operator learns it from the log rather than from a caller talking over the agent with nothing happening. See the [design discussion](https://github.com/streamcoreai/streamcore-server/issues/75) for the trade-offs. +- **The engine name is case-sensitive and sent verbatim.** `telnyx` is rejected with a structured error frame that lists the supported engines. Only `Deepgram` and `Telnyx` are verified here; the other hosted engines the endpoint fronts (AssemblyAI, Azure, and the rest of that list) pass through untested. -Confidence arrives as `null` from the in-house engine and a 0–1 float from hosted ones; the pipeline treats `null` as unknown rather than low. +Confidence arrives as `null` from the in-house engine and a 0-1 float from hosted ones; the pipeline treats `null` as unknown rather than low. ## Local VibeVoice setup diff --git a/docs/providers.zh-CN.md b/docs/providers.zh-CN.md index 19e2654..756da9f 100644 --- a/docs/providers.zh-CN.md +++ b/docs/providers.zh-CN.md @@ -12,7 +12,8 @@ 注意: -- `stt.provider = "openai"` 使用批量最终转写而不是流式中间结果;可通过 `openai.stt_model` 选择 `whisper-1`、`gpt-4o-transcribe` 或 `gpt-4o-mini-transcribe`。 +- `stt.provider = "openai"` 使用批量最终转写而不是流式中间结果,打断(barge-in)与实时字幕都依赖中间结果,因此都不工作;可通过 `openai.stt_model` 选择 `whisper-1`、`gpt-4o-transcribe` 或 `gpt-4o-mini-transcribe`。 +- `stt.provider = "telnyx"` 通过一条 WebSocket、一个 key 前置自研与十余种托管转写引擎。`transcription_engine` 默认 `Deepgram`,打断与实时字幕开箱即用;自研 `Telnyx` 引擎只出最终结果,选中它时两者关闭。见 [Telnyx STT](#telnyx-stt)。 - `llm.provider = "ollama"` 通过 `base_url` 指向任何兼容 Ollama 的端点 —— 本地或你自己的基础设施均可。 - `llm.provider = "agent"` 把每一轮对话 POST 到你托管的 HTTP 端点;记忆、提示词与工具都由你的智能体掌控,回复以 SSE、分块文本或 JSON 流式返回。见[接入你自己的智能体](./bring-your-own-agent.zh-CN.md)。 - `stt.provider = "vibevoice"` 与 `tts.provider = "vibevoice"` 使用本地模型;请先启动 Python 边车进程。 @@ -21,7 +22,6 @@ - `tts.provider = "mimo"` 是小米 MiMo TTS,中英文音色齐备,付费模型还支持声音克隆。 - `stt.provider = "aliyun"` 是阿里云百炼(DashScope)流式 ASR;`vocabulary_id` 可以把模型往你的领域词上带。 - `stt.provider = "volcengine"` 是豆包流式 ASR —— 适合 Deepgram 访问慢、或它的中文识别不够好的场景。控制台有免费时长可以先试。 -- `stt.provider = "telnyx"` 通过一条 WebSocket、一个 key 前置自研与十余种托管转写引擎。自研引擎只出最终结果,打断与实时字幕需要托管引擎。见 [Telnyx STT](#telnyx-stt)。 - `realtime.provider = "grok"` 切换到语音到语音模式,并完全忽略 `[stt]`、`[llm]` 与 `[tts]`。 所有 key 与可调项都在[配置参考](./configuration.zh-CN.md)里。 @@ -151,7 +151,7 @@ voice_speed = 1.0 ## Telnyx STT -走 Telnyx 语音转文字 WebSocket 的流式识别:上行是裸 linear16 二进制帧,下行是 JSON 转写帧,均为流水线原生的 16 kHz 单声道,音频路径无需任何重采样。与 TTS 共用同一个 `[telnyx]` 配置段和 API key;`stt_engine` 选择识别引擎。 +走 Telnyx 语音转文字 WebSocket 的流式识别:上行是裸 linear16 二进制帧,下行是 JSON 转写帧,均为流水线原生的 16 kHz 单声道,音频路径无需任何重采样。与 TTS 共用同一个 `[telnyx]` 配置段和 API key;`transcription_engine` 选择识别引擎。 ```toml [stt] @@ -159,16 +159,16 @@ provider = "telnyx" [telnyx] api_key = "" -stt_engine = "Telnyx" # 或托管引擎 —— "Deepgram"、"AssemblyAI"、"Azure" 等。大小写敏感 +transcription_engine = "Deepgram" # 已验证取值:"Deepgram"(有中间结果,打断可用)或 "Telnyx"(自研,只出最终结果)。大小写敏感 ``` 有三件事必须弄对: -- **自研 `Telnyx` 引擎只出最终结果。** 它在来电者停止说话后才发出唯一一帧 final —— 没有中间结果、没有时间戳、没有置信度。打断(barge-in)与实时字幕依赖中间结果(流水线以部分文本来判定打断),因此**这两项在该引擎下不工作**:打断永远不会触发,客户端字幕也要等到 final 落地才有内容。只适合「仅需最终转写」的场景。取舍讨论见[设计讨论](https://github.com/streamcoreai/streamcore-server/issues/75)。 -- **托管引擎可以恢复中间结果。** 同一个端点前置 AssemblyAI、Azure、Cohere、Deepgram、Google、Humain、Parakeet、Reson8、Soniox、Speechmatics 与 xAI。除 `Telnyx` 以外的任何引擎都会带上 `interim_results=true` 请求,中间结果像其他服务商一样流式到达,打断与实时字幕正常工作。已验证的配置是 `stt_engine = "Deepgram"`。 -- **引擎名大小写敏感。** `telnyx` 会被一帧结构化错误拒绝,错误里列出支持的引擎;合法取值就是那个列表,首字母大写。 +- **默认是 `Deepgram`。** 它像内置的 Deepgram 服务商一样流式输出中间结果,打断(barge-in)与实时字幕开箱即用。一个只出最终结果的默认引擎,会让任何只改了 provider 就开始说话的人悄无声息地失去这两项能力。 +- **自研 `Telnyx` 引擎只出最终结果。** 它在来电者停止说话后才发出唯一一帧 final —— 没有中间结果、没有时间戳、置信度为 `null`。打断与实时字幕依赖中间结果(流水线以部分文本来判定打断),因此**这两项在该引擎下不工作**:打断永远不会触发,客户端字幕也要等到 final 落地才有内容。因为该引擎只在整个话语结束后才应答,客户端自己做端点检测(`internal/vad`,即流水线在用的同一个检测器,静音窗口与 `openai.go` 相同),每个话语开一条新连接,final 落地后由客户端关闭(服务端会一直握着连接不放)。启动时日志里会写明本会话的打断与实时字幕已关闭,运维从日志就能知道,而不用等到来电者对着智能体说话却毫无反应。取舍讨论见[设计讨论](https://github.com/streamcoreai/streamcore-server/issues/75)。 +- **引擎名大小写敏感,原样透传。** `telnyx` 会被一帧结构化错误拒绝,错误里列出支持的引擎。此处只验证了 `Deepgram` 与 `Telnyx` 两个取值;该端点前置的其他托管引擎(AssemblyAI、Azure 及列表中的其余引擎)可透传但未经测试。 -置信度:自研引擎返回 `null`,托管引擎返回 0–1 浮点数;流水线把 `null` 视为未知而不是低置信。 +置信度:自研引擎返回 `null`,托管引擎返回 0-1 浮点数;流水线把 `null` 视为未知而不是低置信。 ## 本地 VibeVoice 配置 diff --git a/internal/config/config.go b/internal/config/config.go index 606de99..ac9aa9e 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -305,15 +305,15 @@ type TelnyxConfig struct { // VoiceSpeed is a playback-rate multiplier (1.0 = normal), clamped to // the 0.8-1.2 conversational band. Zero defaults to 1.0. VoiceSpeed float64 `toml:"voice_speed"` - // STTEngine selects the transcription engine used when stt.provider = - // "telnyx": the in-house "Telnyx" recognizer or any of the hosted - // engines the endpoint fronts (AssemblyAI, Azure, Cohere, Deepgram, - // Google, Humain, Parakeet, Reson8, Soniox, Speechmatics, xAI). The - // value is case-sensitive and sent verbatim. The in-house engine emits - // one final per connection and no interims, so barge-in and live - // captions do not work with it; hosted engines stream interims and both - // work. Empty defaults to "Telnyx". - STTEngine string `toml:"stt_engine"` + // TranscriptionEngine selects the recognizer used when stt.provider = + // "telnyx". The value is case-sensitive and sent to the endpoint + // verbatim. "Deepgram" (the default) streams partials, so barge-in and + // live captions work out of the box; the in-house "Telnyx" engine + // emits one final per utterance and no interims, so barge-in and live + // captions do not work with it and a startup log line says so. Other + // engines the endpoint fronts (AssemblyAI, Azure, ...) pass through + // untested. + TranscriptionEngine string `toml:"transcription_engine"` } type MiMoConfig struct { @@ -461,7 +461,7 @@ func Load(path string) (*Config, error) { if cfg.Telnyx.VoiceSpeed == 0 { cfg.Telnyx.VoiceSpeed = 1.0 } - setDefault(&cfg.Telnyx.STTEngine, "Telnyx") + setDefault(&cfg.Telnyx.TranscriptionEngine, "Deepgram") setDefault(&cfg.LLM.Provider, "openai") setDefault(&cfg.TTS.Provider, "cartesia") setDefault(&cfg.OpenAI.Model, "gpt-4o-mini") diff --git a/internal/stt/telnyx.go b/internal/stt/telnyx.go index 93c2c29..79ada34 100644 --- a/internal/stt/telnyx.go +++ b/internal/stt/telnyx.go @@ -14,111 +14,377 @@ import ( "github.com/gorilla/websocket" + "github.com/streamcoreai/streamcore-server/internal/audio" "github.com/streamcoreai/streamcore-server/internal/config" + "github.com/streamcoreai/streamcore-server/internal/vad" ) +// telnyxSTTURL is the streaming transcription WebSocket. One endpoint +// fronts every engine; the engine rides the transcription_engine query +// parameter. It is a variable only so hermetic tests can point the client +// at a fake server, the same override pattern as grokDialURL. +var telnyxSTTURL = "wss://api.telnyx.com/v2/speech-to-text/transcription" + const ( - // telnyxSTTURL is the streaming transcription WebSocket. One endpoint - // fronts every engine; the engine rides the transcription_engine query - // parameter. - telnyxSTTURL = "wss://api.telnyx.com/v2/speech-to-text/transcription" // telnyxInHouseEngine is Telnyx's own recognizer, reached only by // spelling the name exactly: the parameter is case-sensitive, and a // lowercase "telnyx" is rejected with a structured error frame listing // the supported engines. telnyxInHouseEngine = "Telnyx" - // telnyxFinalSpentReason is what SendAudio reports once the in-house - // engine's single final has been delivered. Saying "connection closed" - // here would read as a network failure when it is in fact the documented - // one-final-per-connection contract. - telnyxFinalSpentReason = "the in-house Telnyx engine emits one final per connection and this connection's final has been delivered; set stt_engine to a streaming engine (e.g. \"Deepgram\") for continuous transcription" + // telnyxDefaultEngine is what an unset transcription_engine falls back + // to. It streams partials, so barge-in and live captions work without + // any extra configuration; a finals-only default would silently switch + // both off for anyone who just sets the provider and starts talking. + telnyxDefaultEngine = "Deepgram" + + // Client-side endpointing for the in-house engine, which transcribes + // only a finished utterance. Windows are counted in pipeline frames: + // the inbound loop delivers 20ms of audio per SendAudio call. + telnyxVADThreshold = 1200.0 // the same base the pipeline's VADs start from + telnyxVADSpeechFrames = 10 // 200ms of speech to open an utterance + telnyxVADSilentFrames = 30 // 600ms of silence to close one, the silence timeout openai.go uses for the same job + + // One frame of linear16 at the pipeline's native rate: 320 samples, + // two bytes each. + telnyxFrameBytes = audio.FrameSize * 2 + + // Idle audio retained before onset, so the frames that lifted the + // detector over its start threshold are not lost: onset fires on the + // 10th loud frame, and 1s of lookback covers that with margin. + telnyxOnsetLookbackBytes = 50 * telnyxFrameBytes + + // An utterance is capped like openai.go caps its buffer: past 30s the + // audio is sent as-is and a fresh utterance opens, so a caller who + // never pauses still gets finals. + telnyxMaxUtteranceBytes = 1500 * telnyxFrameBytes + + // The utterance upload is paced at the fastest rate verified against + // the live endpoint: 200ms of audio per write, a 30ms pause between + // writes. Sending slower only delays the final; sending faster is + // untested. + telnyxSendBatchFrames = 10 + telnyxSendPace = 30 * time.Millisecond + + // Safety net for a socket that never delivers its final: the server + // normally answers within a few seconds of audio stopping. + telnyxFinalWait = 30 * time.Second ) -// TelnyxClient implements the Client interface against Telnyx's streaming -// speech-to-text WebSocket. -// -// The wire protocol is raw linear16 PCM upstream (binary frames at the -// pipeline's native 16 kHz mono, so nothing resamples) and JSON transcript -// frames downstream. One endpoint fronts a dozen engines — Telnyx's in-house -// recognizer plus hosted Deepgram, AssemblyAI, Azure, and others — selected -// by the transcription_engine query parameter, so one [telnyx] section and -// one API key cover both STT and TTS. +// NewTelnyxClient returns the right client for the configured engine. The +// two engines have opposite socket lifecycles, because they stream +// different things: // -// The engines differ in what they stream, and the difference is not cosmetic: +// - Every hosted engine (Deepgram, AssemblyAI, ...) streams interims, +// so it gets one long-lived socket per session: partials flow into +// onResult with IsFinal false, the same shape as deepgram.go, and +// barge-in and live captions work. +// - The in-house engine emits exactly one final per utterance, only +// after audio stops, and holds the socket open afterwards. A +// long-lived socket would be spent after the first sentence, so it +// gets one socket per utterance instead: internal/vad.Detector +// endpoints the speech client-side, the utterance is uploaded, the +// single final is awaited, and the client closes the socket. This is +// the same endpointing problem openai.go solves for its batch +// transcription, reused here per utterance. // -// - The in-house engine emits exactly one final frame after the caller -// stops speaking — no interims, no timestamps, confidence null — and -// then holds the socket open indefinitely; the server never closes it. -// The client therefore tears the connection down itself once that final -// is delivered, and the next SendAudio fails with the reason above -// rather than pretending to accept audio. Barge-in and live captions -// hard-depend on interim results (internal/pipeline/inbound.go gates -// interruption on partial text), so this engine suits -// final-transcript-only deployments; see the design discussion in -// issue #75. -// - Every other engine streams partials when asked, so the client appends -// interim_results=true for them and barge-in works as with any other -// provider. On the in-house engine the parameter is silently ignored — -// the final still arrives — so it is simply not sent there. -type TelnyxClient struct { +// The engine string is sent to the endpoint verbatim: it is +// case-sensitive on the server side, so it is never trimmed or +// normalised. Only "Deepgram" and "Telnyx" are verified; other engines +// pass through untested. +func NewTelnyxClient(ctx context.Context, cfg config.TelnyxConfig, onResult func(TranscriptResult)) (Client, error) { + engine := cfg.TranscriptionEngine + if engine == "" { + engine = telnyxDefaultEngine + } + + if engine == telnyxInHouseEngine { + log.Printf("[stt] telnyx in-house engine emits finals only: no interim results, so barge-in and live captions are off for this session") + return newTelnyxUtteranceClient(ctx, cfg.APIKey, onResult), nil + } + return newTelnyxSessionClient(ctx, engine, cfg.APIKey, onResult) +} + +// telnyxSessionClient is the streaming-engine half: one long-lived socket +// per session, the shape deepgram.go uses. +type telnyxSessionClient struct { conn *websocket.Conn cancel context.CancelFunc - // closeOnFinal tears the connection down after the first final, the - // in-house engine's one-shot contract above. - closeOnFinal bool + // writeMu serialises writes: gorilla/websocket allows only one + // concurrent writer, and Close races SendAudio when a call ends + // mid-utterance. + writeMu sync.Mutex + closed bool +} + +func newTelnyxSessionClient(ctx context.Context, engine, apiKey string, onResult func(TranscriptResult)) (Client, error) { + sttCtx, cancel := context.WithCancel(ctx) + + conn, err := dialTelnyxSTT(sttCtx, engine, apiKey) + if err != nil { + cancel() + return nil, err + } + + c := &telnyxSessionClient{conn: conn, cancel: cancel} + + log.Printf("[stt] connected to Telnyx STT (engine %s)", engine) + go c.readLoop(sttCtx, onResult) - // writeMu serialises writes: gorilla/websocket allows only one concurrent - // writer, and Close races SendAudio when a call ends mid-utterance. - writeMu sync.Mutex - closed bool - closeReason string + return c, nil } -// NewTelnyxClient dials the Telnyx transcription endpoint and starts a -// goroutine that forwards transcripts to onResult. Empty stt_engine defaults -// to the in-house Telnyx engine. -func NewTelnyxClient(ctx context.Context, cfg config.TelnyxConfig, onResult func(TranscriptResult)) (Client, error) { +// readLoop forwards transcript frames to onResult until the connection +// ends. Partials and finals both flow: barge-in and live captions read +// the partials, the turn buffer reads the finals. +func (c *telnyxSessionClient) readLoop(ctx context.Context, onResult func(TranscriptResult)) { + for { + _, msg, err := c.conn.ReadMessage() + if err != nil { + if ctx.Err() == nil { + log.Printf("[stt] telnyx read: %v", err) + } + return + } + + result, emit, rerr := routeTelnyxSTTFrame(msg) + if rerr != nil { + // The server closes the connection right after an error + // frame; there is nothing left to read. + log.Printf("[stt] telnyx: %v", rerr) + return + } + if !emit { + continue + } + if result.IsFinal { + log.Printf("[stt] final: %q", result.Text) + } + onResult(result) + } +} + +// SendAudio forwards a chunk of PCM. The pipeline calls this every 20ms. +func (c *telnyxSessionClient) SendAudio(data []byte) error { + c.writeMu.Lock() + defer c.writeMu.Unlock() + if c.closed { + return fmt.Errorf("telnyx stt: connection closed") + } + return c.conn.WriteMessage(websocket.BinaryMessage, data) +} + +// Close tears down the session connection. +func (c *telnyxSessionClient) Close() { + c.writeMu.Lock() + if c.closed { + c.writeMu.Unlock() + return + } + c.closed = true + c.writeMu.Unlock() + + _ = c.conn.Close() + c.cancel() +} + +// telnyxUtteranceClient is the in-house-engine half: one socket per +// utterance with client-side endpointing. +type telnyxUtteranceClient struct { + ctx context.Context + cancel context.CancelFunc + apiKey string + onResult func(TranscriptResult) + + // mu guards every field below it and the in-flight socket registry. + mu sync.Mutex + detector *vad.Detector + speechActive bool + // utterance holds the frames of the utterance being spoken, from the + // onset lookback to the frame that ends it. + utterance [][]byte + uttsBytes int + lookback [][]byte + lookbackBytes int + // conns holds the utterance sockets still in flight so Close can tear + // them down and unblock their reads. + conns map[*websocket.Conn]struct{} + wg sync.WaitGroup +} + +func newTelnyxUtteranceClient(ctx context.Context, apiKey string, onResult func(TranscriptResult)) Client { sttCtx, cancel := context.WithCancel(ctx) + return &telnyxUtteranceClient{ + ctx: sttCtx, + cancel: cancel, + apiKey: apiKey, + onResult: onResult, + detector: vad.New(telnyxVADThreshold, telnyxVADSpeechFrames, telnyxVADSilentFrames), + conns: make(map[*websocket.Conn]struct{}), + } +} - engine := strings.TrimSpace(cfg.STTEngine) - if engine == "" { - engine = telnyxInHouseEngine +// SendAudio runs the endpointing and never fails across utterance +// boundaries: an error would stop the pipeline's inbound loop, so a spent +// socket (there is one per utterance now) must simply be replaced. +// +// The frames are copied because the pipeline reuses its send buffer, and +// they are buffered rather than written through: the utterance socket is +// only dialed after the detector reports the end of speech, the pattern +// openai.go uses for the same problem. Streaming the audio live instead +// would buy nothing, since this engine emits no partials, and it would +// mean dialing (about 600ms) against the first syllable. +func (c *telnyxUtteranceClient) SendAudio(data []byte) error { + c.mu.Lock() + defer c.mu.Unlock() + + if c.ctx.Err() != nil { + return fmt.Errorf("telnyx stt: client closed") } - header := http.Header{} - header.Set("Authorization", "Bearer "+cfg.APIKey) + started, ended := c.detector.Process(audio.Linear16BytesToPCM(data)) + if started { + c.speechActive = true + // The lookback becomes the utterance prefix, so the onset + // frames the detector needed to see before firing are + // transcribed too. + c.utterance = append(c.utterance, c.lookback...) + c.uttsBytes += c.lookbackBytes + c.lookback = nil + c.lookbackBytes = 0 + } - dialer := websocket.Dialer{HandshakeTimeout: 10 * time.Second} - conn, resp, err := dialer.DialContext(sttCtx, telnyxSTTEndpoint(engine), header) + frame := append([]byte(nil), data...) + if c.speechActive { + c.utterance = append(c.utterance, frame) + c.uttsBytes += len(frame) + } else { + c.lookback = append(c.lookback, frame) + c.lookbackBytes += len(frame) + for c.lookbackBytes > telnyxOnsetLookbackBytes { + c.lookbackBytes -= len(c.lookback[0]) + c.lookback = c.lookback[1:] + } + } + + var frames [][]byte + if ended { + c.speechActive = false + frames = c.utterance + c.utterance = nil + c.uttsBytes = 0 + } else if c.uttsBytes >= telnyxMaxUtteranceBytes { + // Force-flush mid-speech, like openai.go's buffer cap: speech + // continues into a fresh utterance, and the detector keeps its + // state. + frames = c.utterance + c.utterance = nil + c.uttsBytes = 0 + } + if len(frames) > 0 { + c.wg.Add(1) + go c.transcribe(frames) + } + return nil +} + +// transcribe dials one socket, uploads the finished utterance at the +// verified pace, waits for the engine's single final, and closes the +// socket: the server holds it open after the final, so the client has to +// end the conversation itself. +func (c *telnyxUtteranceClient) transcribe(frames [][]byte) { + defer c.wg.Done() + + conn, err := dialTelnyxSTT(c.ctx, telnyxInHouseEngine, c.apiKey) if err != nil { - cancel() - if resp != nil { - // The body carries the actual reason — a rejected key reads as a - // bare 401 — which the status code alone does not convey. - body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) - resp.Body.Close() - return nil, fmt.Errorf("telnyx stt dial: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) + if c.ctx.Err() == nil { + log.Printf("[stt] telnyx utterance dial: %v", err) } - return nil, fmt.Errorf("telnyx stt dial: %w", err) + return } + defer conn.Close() - c := &TelnyxClient{ - conn: conn, - cancel: cancel, - closeOnFinal: engine == telnyxInHouseEngine, + c.mu.Lock() + if c.ctx.Err() != nil { + c.mu.Unlock() + return } + c.conns[conn] = struct{}{} + c.mu.Unlock() + defer func() { + c.mu.Lock() + delete(c.conns, conn) + c.mu.Unlock() + }() - log.Printf("[stt] connected to Telnyx STT (engine %s)", engine) - go c.readLoop(sttCtx, onResult) + for i, frame := range frames { + if c.ctx.Err() != nil { + return + } + if err := conn.WriteMessage(websocket.BinaryMessage, frame); err != nil { + if c.ctx.Err() == nil { + log.Printf("[stt] telnyx utterance write: %v", err) + } + return + } + if (i+1)%telnyxSendBatchFrames == 0 { + select { + case <-time.After(telnyxSendPace): + case <-c.ctx.Done(): + return + } + } + } - return c, nil + // The final only lands after the audio stops, which it just did: the + // deadline only guards a socket that never answers. + _ = conn.SetReadDeadline(time.Now().Add(telnyxFinalWait)) + for { + _, msg, err := conn.ReadMessage() + if err != nil { + if c.ctx.Err() == nil { + log.Printf("[stt] telnyx utterance read: %v", err) + } + return + } + result, emit, rerr := routeTelnyxSTTFrame(msg) + if rerr != nil { + log.Printf("[stt] telnyx: %v", rerr) + return + } + if !emit || !result.IsFinal { + // The engine's contract is one final and nothing before + // it; an empty marker frame is skipped, and an interim, + // which the contract says cannot arrive, is dropped + // rather than double-echoed inside the final. + continue + } + log.Printf("[stt] final: %q", result.Text) + c.onResult(result) + return + } +} + +// Close cancels the client and waits for in-flight utterances to finish +// uploading or give up. The sockets are closed explicitly so a blocked +// read ends immediately instead of waiting out its deadline. +func (c *telnyxUtteranceClient) Close() { + c.cancel() + + c.mu.Lock() + for conn := range c.conns { + _ = conn.Close() + } + c.mu.Unlock() + + c.wg.Wait() } // telnyxSTTEndpoint builds the dial URL. The engine value is URL-encoded -// because it reaches the server as a query parameter, and interim_results is -// appended for every engine except the in-house one: only the hosted engines -// emit partials, and the pipeline's barge-in hard-depends on them. +// because it reaches the server as a query parameter, and interim_results +// is appended for every engine except the in-house one: only the hosted +// engines stream partials, and the pipeline's barge-in hard-depends on +// them. The in-house engine ignores the parameter, so it is not sent. func telnyxSTTEndpoint(engine string) string { endpoint := fmt.Sprintf("%s?transcription_engine=%s&input_format=linear16&sample_rate=16000", telnyxSTTURL, url.QueryEscape(engine)) @@ -128,17 +394,39 @@ func telnyxSTTEndpoint(engine string) string { return endpoint } +// dialTelnyxSTT opens one transcription socket with the API key in the +// Authorization header. +func dialTelnyxSTT(ctx context.Context, engine, apiKey string) (*websocket.Conn, error) { + header := http.Header{} + header.Set("Authorization", "Bearer "+apiKey) + + dialer := websocket.Dialer{HandshakeTimeout: 10 * time.Second} + conn, resp, err := dialer.DialContext(ctx, telnyxSTTEndpoint(engine), header) + if err != nil { + if resp != nil { + // The body carries the actual reason, a rejected key reads as + // a bare 401, which the status code alone does not convey. + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + resp.Body.Close() + return nil, fmt.Errorf("telnyx stt dial: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) + } + return nil, fmt.Errorf("telnyx stt dial: %w", err) + } + return conn, nil +} + // telnyxSTTFrame is one transcript message. The in-house engine reports // confidence as null; hosted engines report a 0-1 float. type telnyxSTTFrame struct { Transcript string `json:"transcript"` - // Confidence is a pointer so null — the in-house engine — stays - // distinguishable from a genuine 0.0. Null maps to 0, which + // Confidence is a pointer so null, the in-house engine's value, + // stays distinguishable from a genuine 0.0. Null maps to 0, which // TranscriptResult defines as "unknown", not "low". Confidence *float64 `json:"confidence"` IsFinal bool `json:"is_final"` - // Errors carries a rejected request (an unsupported engine, say) instead - // of a transcript; the server closes the connection right after. + // Errors carries a rejected request (an unsupported engine, say) + // instead of a transcript; the server closes the connection right + // after. Errors []telnyxSTTError `json:"errors"` } @@ -151,11 +439,11 @@ type telnyxSTTError struct { Detail string `json:"detail"` } -// routeTelnyxSTTFrame classifies one server frame: emit reports whether the -// frame carries a transcript the pipeline should see, and err is set for -// error frames and undecodable messages. Frames with no transcript text — -// utterance-end markers, keepalives — come back with emit false so the -// caller skips them; emitting one would surface an empty final turn. +// routeTelnyxSTTFrame classifies one server frame: emit reports whether +// the frame carries a transcript the pipeline should see, and err is set +// for error frames and undecodable messages. Frames with no transcript +// text, utterance-end markers, keepalives, come back with emit false so +// the caller skips them; emitting one would surface an empty final turn. func routeTelnyxSTTFrame(msg []byte) (result TranscriptResult, emit bool, err error) { var f telnyxSTTFrame if jerr := json.Unmarshal(msg, &f); jerr != nil { @@ -163,9 +451,10 @@ func routeTelnyxSTTFrame(msg []byte) (result TranscriptResult, emit bool, err er } if len(f.Errors) > 0 { // The detail names the actual problem ("Unsupported - // transcription_engine 'telnyx'. Supported engines: …"), which beats - // the code and title for an operator reading the log. It is the - // first entry: the service reports one error per frame. + // transcription_engine 'telnyx'. Supported engines: ..."), + // which beats the code and title for an operator reading the + // log. It is the first entry: the service reports one error per + // frame. e := f.Errors[0] detail := e.Detail if detail == "" { @@ -183,73 +472,3 @@ func routeTelnyxSTTFrame(msg []byte) (result TranscriptResult, emit bool, err er } return TranscriptResult{Text: text, IsFinal: f.IsFinal, Confidence: confidence}, true, nil } - -// readLoop forwards transcript frames to onResult until the connection ends. -func (c *TelnyxClient) readLoop(ctx context.Context, onResult func(TranscriptResult)) { - for { - _, msg, err := c.conn.ReadMessage() - if err != nil { - if ctx.Err() == nil { - log.Printf("[stt] telnyx read: %v", err) - } - return - } - - result, emit, rerr := routeTelnyxSTTFrame(msg) - if rerr != nil { - // The server closes the connection right after an error frame; - // there is nothing left to read. - log.Printf("[stt] telnyx: %v", rerr) - return - } - if !emit { - continue - } - if result.IsFinal { - log.Printf("[stt] final: %q", result.Text) - } - onResult(result) - - // The in-house engine's contract is one final per connection, and - // the server holds the socket open afterwards — it never closes. - // Tear it down from this side so the spent session does not linger - // until call teardown, and let SendAudio report why. - if result.IsFinal && c.closeOnFinal { - c.shutdown(telnyxFinalSpentReason) - return - } - } -} - -// SendAudio forwards a chunk of PCM. The pipeline calls this every 20ms. -func (c *TelnyxClient) SendAudio(data []byte) error { - c.writeMu.Lock() - defer c.writeMu.Unlock() - if c.closed { - return fmt.Errorf("telnyx stt: %s", c.closeReason) - } - return c.conn.WriteMessage(websocket.BinaryMessage, data) -} - -// Close tears down the connection. There is no flush to wait for: the -// service settles an utterance on its own silence detection, so a -// client-side close mid-utterance simply drops the stream. -func (c *TelnyxClient) Close() { - c.shutdown("connection closed") -} - -// shutdown closes the connection once, recording why so SendAudio can tell a -// spent finals-only session apart from a normal teardown. -func (c *TelnyxClient) shutdown(reason string) { - c.writeMu.Lock() - if c.closed { - c.writeMu.Unlock() - return - } - c.closed = true - c.closeReason = reason - c.writeMu.Unlock() - - _ = c.conn.Close() - c.cancel() -} diff --git a/internal/stt/telnyx_test.go b/internal/stt/telnyx_test.go index 0e67b32..02e912e 100644 --- a/internal/stt/telnyx_test.go +++ b/internal/stt/telnyx_test.go @@ -1,12 +1,23 @@ package stt import ( + "context" + "encoding/binary" + "net/http" + "net/http/httptest" "strings" + "sync/atomic" "testing" + "time" + + "github.com/gorilla/websocket" + + "github.com/streamcoreai/streamcore-server/internal/audio" + "github.com/streamcoreai/streamcore-server/internal/config" ) // The in-house engine's final frame carries confidence null. A nil pointer -// must map to 0, which TranscriptResult defines as "unknown" — a naive +// must map to 0, which TranscriptResult defines as "unknown", so a naive // dereference would panic, and a misread null would look like a failed // recognition. func TestTelnyxSTTFinalFrameMapsNullConfidence(t *testing.T) { @@ -92,10 +103,10 @@ func TestTelnyxSTTRejectsMalformedJSON(t *testing.T) { } } -// The engine rides the dial URL as a query parameter, so it must be encoded; -// and interim_results decides whether barge-in works at all. The in-house -// engine ignores the parameter, so it is not sent there; every hosted engine -// streams partials only when it is present. +// The engine rides the dial URL as a query parameter, so it must be +// encoded, and interim_results decides whether barge-in works at all. The +// in-house engine ignores the parameter, so it is not sent there; every +// hosted engine streams partials only when it is present. func TestTelnyxSTTEndpointBuildsQuery(t *testing.T) { cases := []struct { name string @@ -113,7 +124,7 @@ func TestTelnyxSTTEndpointBuildsQuery(t *testing.T) { telnyxSTTURL + "?transcription_engine=Deepgram&input_format=linear16&sample_rate=16000&interim_results=true", }, { - "engine value is URL-encoded", + "engine value is URL-encoded with casing preserved", "Two Words", telnyxSTTURL + "?transcription_engine=Two+Words&input_format=linear16&sample_rate=16000&interim_results=true", }, @@ -127,18 +138,330 @@ func TestTelnyxSTTEndpointBuildsQuery(t *testing.T) { } } -// After the in-house engine's single final, SendAudio must refuse rather -// than pretend to accept audio, and the message must point at the fix — -// "connection closed" would read as a network failure when the connection -// is behaving exactly as documented. -func TestTelnyxSTTSendAudioReportsSpentSession(t *testing.T) { - c := &TelnyxClient{closed: true, closeReason: telnyxFinalSpentReason} +// fakeTelnyxSTT is a stand-in for the Telnyx transcription endpoint. It +// records the query string of every dial, checks the auth header, and +// hands each connection to a test-supplied handler, so both socket +// lifecycles can be exercised without network access or an API key. +type fakeTelnyxSTT struct { + t *testing.T + server *httptest.Server + queries chan string + handler func(conn *websocket.Conn) +} + +func newFakeTelnyxSTT(t *testing.T, handler func(conn *websocket.Conn)) *fakeTelnyxSTT { + t.Helper() + f := &fakeTelnyxSTT{ + t: t, + queries: make(chan string, 16), + handler: handler, + } - err := c.SendAudio([]byte{0x00, 0x01}) - if err == nil { - t.Fatal("send after close accepted, want the spent-session reason") + upgrader := websocket.Upgrader{} + f.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if got := r.Header.Get("Authorization"); got != "Bearer k" { + t.Errorf("Authorization = %q, want %q", got, "Bearer k") + } + select { + case f.queries <- r.URL.RawQuery: + default: + } + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + return + } + defer conn.Close() + f.handler(conn) + })) + t.Cleanup(f.server.Close) + return f +} + +func (f *fakeTelnyxSTT) wsURL() string { + return "ws" + strings.TrimPrefix(f.server.URL, "http") +} + +func (f *fakeTelnyxSTT) dialQuery(t *testing.T) string { + t.Helper() + select { + case q := <-f.queries: + return q + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for a dial") + return "" + } +} + +// overrideTelnyxSTTURL points the client at the fake server for the +// duration of one test. +func overrideTelnyxSTTURL(t *testing.T, url string) { + t.Helper() + orig := telnyxSTTURL + telnyxSTTURL = url + t.Cleanup(func() { telnyxSTTURL = orig }) +} + +// pcmFrame returns one 20ms frame of constant linear16 PCM at the given +// amplitude: 20000 sits far above the VAD threshold, 0 is silence. +func pcmFrame(amp int16) []byte { + b := make([]byte, audio.FrameSize*2) + for i := range audio.FrameSize { + binary.LittleEndian.PutUint16(b[i*2:], uint16(amp)) } - if !strings.Contains(err.Error(), "stt_engine") { - t.Errorf("error = %v, want it to point at the stt_engine fix", err) + return b +} + +func waitResult(t *testing.T, ch <-chan TranscriptResult) TranscriptResult { + t.Helper() + select { + case r := <-ch: + return r + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for a transcript result") + return TranscriptResult{} + } +} + +// The hosted engines stream over one long-lived socket: partials and finals +// flow, and the final does not end the session, a later partial still +// arrives on the same connection. +func TestTelnyxSessionClientStreamsPartialsThenFinal(t *testing.T) { + results := make(chan TranscriptResult, 8) + f := newFakeTelnyxSTT(t, func(conn *websocket.Conn) { + _ = conn.WriteMessage(websocket.TextMessage, + []byte(`{"transcript":"The quick","confidence":0.8442383,"is_final":false}`)) + _ = conn.WriteMessage(websocket.TextMessage, + []byte(`{"transcript":"The quick brown fox jumps over the lazy dog.","confidence":0.95,"is_final":true}`)) + // A second partial after the final proves the socket stayed open. + _ = conn.WriteMessage(websocket.TextMessage, + []byte(`{"transcript":"Testing one","confidence":0.9,"is_final":false}`)) + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + }) + overrideTelnyxSTTURL(t, f.wsURL()) + + client, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k", TranscriptionEngine: "Deepgram"}, func(r TranscriptResult) { + results <- r + }) + if err != nil { + t.Fatalf("client: %v", err) + } + defer client.Close() + + if got, want := f.dialQuery(t), "transcription_engine=Deepgram&input_format=linear16&sample_rate=16000&interim_results=true"; got != want { + t.Errorf("dial query = %q, want %q", got, want) + } + + r := waitResult(t, results) + if r.IsFinal || r.Text != "The quick" || r.Confidence != 0.8442383 { + t.Errorf("first result = %+v, want the partial", r) + } + r = waitResult(t, results) + if !r.IsFinal || r.Text != "The quick brown fox jumps over the lazy dog." || r.Confidence != 0.95 { + t.Errorf("second result = %+v, want the final", r) + } + r = waitResult(t, results) + if r.IsFinal || r.Text != "Testing one" { + t.Errorf("third result = %+v, want the post-final partial on the same socket", r) + } + + if err := client.SendAudio(pcmFrame(0)); err != nil { + t.Errorf("session SendAudio: %v", err) + } +} + +// The engine string is sent verbatim: no trimming, no case folding. It is +// case-sensitive on the server side, so any normalization here would turn +// a working configuration into a rejected dial. +func TestTelnyxEngineStringReachesURLVerbatim(t *testing.T) { + f := newFakeTelnyxSTT(t, func(conn *websocket.Conn) { + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + }) + overrideTelnyxSTTURL(t, f.wsURL()) + + client, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k", TranscriptionEngine: "MixedCaseEngine"}, func(TranscriptResult) {}) + if err != nil { + t.Fatalf("client: %v", err) + } + defer client.Close() + + if got, want := f.dialQuery(t), "transcription_engine=MixedCaseEngine&input_format=linear16&sample_rate=16000&interim_results=true"; got != want { + t.Errorf("dial query = %q, want %q", got, want) + } +} + +// An unset engine defaults to Deepgram, the streaming engine, so barge-in +// and live captions work without any extra configuration. +func TestTelnyxEmptyEngineDefaultsToDeepgram(t *testing.T) { + f := newFakeTelnyxSTT(t, func(conn *websocket.Conn) { + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + }) + overrideTelnyxSTTURL(t, f.wsURL()) + + client, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k"}, func(TranscriptResult) {}) + if err != nil { + t.Fatalf("client: %v", err) + } + defer client.Close() + + if got, want := f.dialQuery(t), "transcription_engine=Deepgram&input_format=linear16&sample_rate=16000&interim_results=true"; got != want { + t.Errorf("dial query = %q, want %q", got, want) + } +} + +// The in-house engine emits one final per utterance, so the client runs +// one socket per utterance: the engine string rides the URL without +// interim_results, every frame of the utterance reaches the server, the +// final (confidence null) is forwarded exactly once, the client closes +// the socket afterwards because the server holds it open, and SendAudio +// keeps accepting audio for the next utterance instead of erroring a +// spent session. +func TestTelnyxUtteranceClientOneSocketPerUtterance(t *testing.T) { + results := make(chan TranscriptResult, 4) + // socketClosed fires when the server's read ends, which only happens + // once the client has closed the socket after its single final. + socketClosed := make(chan int, 4) + var dials atomic.Int32 + + f := newFakeTelnyxSTT(t, func(conn *websocket.Conn) { + dials.Add(1) + frames := 0 + // Answer with the single final once the audio starts arriving; + // the client reads it only after it has finished sending. + sentFinal := false + for { + msgType, _, err := conn.ReadMessage() + if err != nil { + socketClosed <- frames + return + } + if msgType != websocket.BinaryMessage { + continue + } + frames++ + if !sentFinal { + sentFinal = true + _ = conn.WriteMessage(websocket.TextMessage, + []byte(`{"transcript":"The quick brown fox jumps over the lazy dog. Testing one two three, this is the Telnyx STT capture for StreamCore.","confidence":null,"is_final":true}`)) + } + } + }) + overrideTelnyxSTTURL(t, f.wsURL()) + + client, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k", TranscriptionEngine: "Telnyx"}, func(r TranscriptResult) { + results <- r + }) + if err != nil { + t.Fatalf("client: %v", err) + } + defer client.Close() + + // Utterance one: 15 loud frames open it, 30 silent frames close it. + for i := 0; i < 15; i++ { + if err := client.SendAudio(pcmFrame(20000)); err != nil { + t.Fatalf("send loud frame: %v", err) + } + } + for i := 0; i < 30; i++ { + if err := client.SendAudio(pcmFrame(0)); err != nil { + t.Fatalf("send silent frame: %v", err) + } + } + + r := waitResult(t, results) + if !r.IsFinal { + t.Errorf("result = %+v, want the single final", r) + } + if r.Text != "The quick brown fox jumps over the lazy dog. Testing one two three, this is the Telnyx STT capture for StreamCore." { + t.Errorf("text = %q, want the capture sentence", r.Text) + } + if r.Confidence != 0 { + t.Errorf("confidence = %v, want 0 (null maps to unknown)", r.Confidence) + } + + // Utterance two: SendAudio across the boundary must keep working, and + // it must produce a second socket, not an error. + for i := 0; i < 15; i++ { + if err := client.SendAudio(pcmFrame(20000)); err != nil { + t.Fatalf("send loud frame after first final: %v", err) + } + } + for i := 0; i < 30; i++ { + if err := client.SendAudio(pcmFrame(0)); err != nil { + t.Fatalf("send silent frame after first final: %v", err) + } + } + + r = waitResult(t, results) + if !r.IsFinal { + t.Errorf("second result = %+v, want the second final", r) + } + + if got := dials.Load(); got != 2 { + t.Errorf("dials = %d, want one socket per utterance (2)", got) + } + for i := 0; i < 2; i++ { + if got, want := f.dialQuery(t), "transcription_engine=Telnyx&input_format=linear16&sample_rate=16000"; got != want { + t.Errorf("utterance dial query = %q, want %q", got, want) + } + } + for i := 0; i < 2; i++ { + select { + case frames := <-socketClosed: + // 15 loud + 30 silent frames = the 45-frame utterance, plus + // the onset lookback when there is one. + if frames < 45 { + t.Errorf("socket closed after %d frames, want the full 45-frame utterance", frames) + } + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for the client to close the socket after its final") + } + } +} + +// Silence alone must never dial: the socket only exists per utterance, and +// there is no utterance until the detector hears speech. +func TestTelnyxUtteranceClientSilenceNeverDials(t *testing.T) { + f := newFakeTelnyxSTT(t, func(conn *websocket.Conn) { + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + }) + overrideTelnyxSTTURL(t, f.wsURL()) + + client, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k", TranscriptionEngine: "Telnyx"}, func(TranscriptResult) {}) + if err != nil { + t.Fatalf("client: %v", err) + } + + for i := 0; i < 200; i++ { + if err := client.SendAudio(pcmFrame(0)); err != nil { + t.Fatalf("send silent frame: %v", err) + } + } + client.Close() + + select { + case q := <-f.queries: + t.Errorf("dialed with query %q on silence, want no connection", q) + default: } }