From 7c54e46ea75a7324b4a2691d9072a2e48f8d0802 Mon Sep 17 00:00:00 2001 From: a692570 Date: Mon, 21 Sep 2026 14:03:09 -0700 Subject: [PATCH] Add EmitsPartials capability flag for finals-only STT providers --- README.md | 2 +- README.zh-CN.md | 2 +- config.toml.example | 6 +- docs/configuration.md | 4 +- docs/configuration.zh-CN.md | 2 +- docs/providers.md | 8 +- docs/providers.zh-CN.md | 8 +- internal/config/config.go | 8 +- internal/pipeline/inbound.go | 19 ++- internal/pipeline/inbound_test.go | 259 ++++++++++++++++++++++++++++++ internal/stt/openai.go | 6 + internal/stt/openai_test.go | 20 +++ internal/stt/stt.go | 18 +++ internal/stt/stt_test.go | 23 +++ internal/stt/telnyx.go | 18 ++- internal/stt/telnyx_test.go | 38 +++++ 16 files changed, 417 insertions(+), 24 deletions(-) create mode 100644 internal/pipeline/inbound_test.go create mode 100644 internal/stt/stt_test.go diff --git a/README.md b/README.md index 4273846..5c0b8a0 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.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. +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 live captions show finals only and barge-in waits out the full backchannel window on VAD alone, and a startup log line says so. ## Documentation diff --git a/README.zh-CN.md b/README.zh-CN.md index 2a0df4c..208daf9 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.transcription_engine`,默认 `Deepgram`,打断与实时字幕开箱即用);自研 `Telnyx` 引擎只出最终结果,选中它时两者关闭,启动日志会写明。 +Telnyx STT 用一个 key 前置十余种引擎(`telnyx.transcription_engine`,默认 `Deepgram`,打断与实时字幕开箱即用);自研 `Telnyx` 引擎只出最终结果,实时字幕只显示最终结果,打断降级为等满整个回应窗口后仅凭 VAD 触发,启动日志会写明。 ## 文档 diff --git a/config.toml.example b/config.toml.example index 3146d68..26765eb 100644 --- a/config.toml.example +++ b/config.toml.example @@ -221,9 +221,9 @@ voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # Any catalog vo voice_speed = 1.0 # Playback-rate multiplier, clamped to 0.8-1.2 like the per-utterance delivery tags 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 + # work) or "Telnyx" (in-house, finals-only: live captions show finals only and + # barge-in waits out the backchannel window on VAD alone; a startup log line + # says so). Case-sensitive, sent verbatim; other hosted engines pass through untested [mimo] api_key = "" # Required if tts.provider = "mimo" diff --git a/docs/configuration.md b/docs/configuration.md index d25574e..045c04f 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -159,8 +159,8 @@ 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 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 + # "Telnyx" (in-house, finals-only: live captions show finals only and barge-in waits + # out the backchannel window on VAD alone; 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 737face..c855463 100644 --- a/docs/configuration.zh-CN.md +++ b/docs/configuration.zh-CN.md @@ -145,7 +145,7 @@ api_key = "" voice = "Telnyx.Qwen3TTS.d9348e0d-988a-42cc-a64e-18093fe45c03" # GET /v2/text-to-speech/voices 目录中的任意音色;可用性因账号而异 voice_speed = 1.0 # 播放速率倍数,限制在 0.8-1.2 transcription_engine = "Deepgram" # STT 引擎,已验证取值:"Deepgram"(有中间结果,打断可用;默认值)或 - # "Telnyx"(自研,只出最终结果:打断与实时字幕关闭,启动日志会写明)。大小写敏感,原样透传 + # "Telnyx"(自研,只出最终结果:实时字幕只显示最终结果,打断降级为等满回应窗口后仅凭 VAD 触发;启动日志会写明)。大小写敏感,原样透传 [minimax] api_key = "" diff --git a/docs/providers.md b/docs/providers.md index 8845470..ae995c6 100644 --- a/docs/providers.md +++ b/docs/providers.md @@ -12,8 +12,8 @@ Notes: -- `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). +- `stt.provider = "openai"` uses batch final transcription instead of streaming partials, so live captions show finals only and barge-in runs in a degraded mode: the interrupt fires on VAD alone once the full 600ms backchannel window has elapsed, since there is no text to classify a short burst with. 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 live captions show finals only and barge-in waits out the backchannel window on VAD alone. 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. @@ -164,8 +164,8 @@ transcription_engine = "Deepgram" # verified: "Deepgram" (partial results, bar Three things to know: -- **`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. +- **`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 degrade both for anyone who just sets the provider and starts talking: barge-in to VAD-only, live captions to finals only. +- **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`. Live captions therefore show nothing until the final lands. Barge-in cannot confirm an interruption from partial text, so it degrades instead of switching off: speech over the agent opens the suppression window on VAD alone, and the interrupt fires once the speech has lasted past the full 600ms backchannel window. A burst that ends inside the window is treated as backchannel — with no text, it cannot be told apart from "mm-hm" — so sustained interruptions work and quick ones do not. 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 the session runs barge-in VAD-only on this engine, so the operator learns it from the log. 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. diff --git a/docs/providers.zh-CN.md b/docs/providers.zh-CN.md index 756da9f..4c4c1f5 100644 --- a/docs/providers.zh-CN.md +++ b/docs/providers.zh-CN.md @@ -12,8 +12,8 @@ 注意: -- `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)。 +- `stt.provider = "openai"` 使用批量最终转写而不是流式中间结果,因此实时字幕只显示最终结果,打断降级为等满 600ms 回应窗口后仅凭 VAD 触发(没有文本,无法对短促语做分类);可通过 `openai.stt_model` 选择 `whisper-1`、`gpt-4o-transcribe` 或 `gpt-4o-mini-transcribe`。 +- `stt.provider = "telnyx"` 通过一条 WebSocket、一个 key 前置自研与十余种托管转写引擎。`transcription_engine` 默认 `Deepgram`,打断与实时字幕开箱即用;自研 `Telnyx` 引擎只出最终结果,实时字幕只显示最终结果,打断降级为等满回应窗口后仅凭 VAD 触发。见 [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 边车进程。 @@ -164,8 +164,8 @@ transcription_engine = "Deepgram" # 已验证取值:"Deepgram"(有中间 有三件事必须弄对: -- **默认是 `Deepgram`。** 它像内置的 Deepgram 服务商一样流式输出中间结果,打断(barge-in)与实时字幕开箱即用。一个只出最终结果的默认引擎,会让任何只改了 provider 就开始说话的人悄无声息地失去这两项能力。 -- **自研 `Telnyx` 引擎只出最终结果。** 它在来电者停止说话后才发出唯一一帧 final —— 没有中间结果、没有时间戳、置信度为 `null`。打断与实时字幕依赖中间结果(流水线以部分文本来判定打断),因此**这两项在该引擎下不工作**:打断永远不会触发,客户端字幕也要等到 final 落地才有内容。因为该引擎只在整个话语结束后才应答,客户端自己做端点检测(`internal/vad`,即流水线在用的同一个检测器,静音窗口与 `openai.go` 相同),每个话语开一条新连接,final 落地后由客户端关闭(服务端会一直握着连接不放)。启动时日志里会写明本会话的打断与实时字幕已关闭,运维从日志就能知道,而不用等到来电者对着智能体说话却毫无反应。取舍讨论见[设计讨论](https://github.com/streamcoreai/streamcore-server/issues/75)。 +- **默认是 `Deepgram`。** 它像内置的 Deepgram 服务商一样流式输出中间结果,打断(barge-in)与实时字幕开箱即用。一个只出最终结果的默认引擎,会让任何只改了 provider 就开始说话的人悄无声息地让这两项能力降级:打断退化为仅凭 VAD,实时字幕只剩最终结果。 +- **自研 `Telnyx` 引擎只出最终结果。** 它在来电者停止说话后才发出唯一一帧 final —— 没有中间结果、没有时间戳、置信度为 `null`。客户端字幕要等到 final 落地才有内容。打断没有中间文本可判定,因此是降级而不是关闭:来电者压过智能体说话时,抑制窗口仅凭 VAD 打开,只有说话持续超过整个 600ms 回应窗口后才确认打断;窗口内结束的短促语一律按回应词处理 —— 没有文本就无法把它和 "嗯嗯" 区分开。持续抢话可以打断,短促抢话不能。因为该引擎只在整个话语结束后才应答,客户端自己做端点检测(`internal/vad`,即流水线在用的同一个检测器,静音窗口与 `openai.go` 相同),每个话语开一条新连接,final 落地后由客户端关闭(服务端会一直握着连接不放)。启动时日志里会写明本会话的打断以仅凭 VAD 的降级模式运行,运维从日志就能知道。取舍讨论见[设计讨论](https://github.com/streamcoreai/streamcore-server/issues/75)。 - **引擎名大小写敏感,原样透传。** `telnyx` 会被一帧结构化错误拒绝,错误里列出支持的引擎。此处只验证了 `Deepgram` 与 `Telnyx` 两个取值;该端点前置的其他托管引擎(AssemblyAI、Azure 及列表中的其余引擎)可透传但未经测试。 置信度:自研引擎返回 `null`,托管引擎返回 0-1 浮点数;流水线把 `null` 视为未知而不是低置信。 diff --git a/internal/config/config.go b/internal/config/config.go index 921d78c..0949eb6 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -425,10 +425,10 @@ type TelnyxConfig struct { // "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. + // emits one final per utterance and no interims, so live captions show + // finals only and barge-in runs VAD-only after the backchannel window — + // a startup log line says so. Other engines the endpoint fronts + // (AssemblyAI, Azure, ...) pass through untested. TranscriptionEngine string `toml:"transcription_engine"` } diff --git a/internal/pipeline/inbound.go b/internal/pipeline/inbound.go index e2c7a59..157982a 100644 --- a/internal/pipeline/inbound.go +++ b/internal/pipeline/inbound.go @@ -129,6 +129,15 @@ func (p *Pipeline) runInbound() { } defer sttClient.Close() + // Finals-only providers never confirm speech with partial text + // (issue #75). Asserted once here: a provider that does not implement + // stt.PartialsEmitter keeps emitsPartials true and the partials-driven + // path below exactly as it was. + emitsPartials := true + if ep, ok := sttClient.(stt.PartialsEmitter); ok { + emitsPartials = ep.EmitsPartials() + } + // Backchannel suppression state machine var bargeInPending bool var bargeInStart time.Time @@ -172,6 +181,10 @@ func (p *Pipeline) runInbound() { } } else if !p.bargeInVAD.IsSpeaking() { // Speech ended within the window — check for backchannel. + // A finals-only provider has no partial text here, so + // every burst that ends inside the window classifies as + // backchannel: with no text, there is no basis to cut + // the agent off mid-word. partial, _ := latestPartial.Load("text") partialStr, _ := partial.(string) if !isMeaningfulBargeInTranscript(partialStr) { @@ -200,8 +213,12 @@ func (p *Pipeline) runInbound() { } } // else: still speaking within window, keep waiting - } else if p.bargeInVAD.IsSpeaking() && p.speaking.Load() && hasPartialText.Load() { + } else if p.bargeInVAD.IsSpeaking() && p.speaking.Load() && (!emitsPartials || hasPartialText.Load()) { // Conditions met — start backchannel suppression window. + // A finals-only provider opens the window on VAD alone, + // since partial text never arrives to confirm the speech; + // the interrupt still cannot fire until the window has + // fully elapsed. bargeInPending = true bargeInStart = time.Now() hasPartialText.Store(false) diff --git a/internal/pipeline/inbound_test.go b/internal/pipeline/inbound_test.go new file mode 100644 index 0000000..229a921 --- /dev/null +++ b/internal/pipeline/inbound_test.go @@ -0,0 +1,259 @@ +package pipeline + +import ( + "context" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "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" +) + +// Hermetic coverage for runInbound's barge-in gate (issue #75). +// +// A provider that streams partials opens the barge-in window only once +// STT has confirmed real speech with partial text; a finals-only provider +// (stt.PartialsEmitter returning false) opens it on VAD alone but cannot +// fire the interrupt until the 600ms backchannel window has fully +// elapsed, and a burst that ends inside the window is suppressed as +// backchannel — with no text, it cannot be told apart from "mm-hm". +// +// runInbound is driven directly with a hand-built Pipeline because the +// full constructor demands WebRTC tracks and live provider credentials. +// The STT client still comes from stt.NewClient, so the type assertion +// runs on the real provider values: a fake VibeVoice WebSocket sidecar +// plays the plain provider (no PartialsEmitter), and the OpenAI client — +// which dials nothing at construction and only flushes after 600ms of +// silence, a state these tests never reach before their context is +// cancelled — plays the finals-only one. + +// fakeVibeVoiceASR stands in for the local VibeVoice ASR sidecar: it +// accepts the client's WebSocket, optionally sends one partial +// transcript, and discards the audio the pipeline streams at it. +func fakeVibeVoiceASR(t *testing.T, sendPartial bool) *httptest.Server { + t.Helper() + upgrader := websocket.Upgrader{} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + return + } + defer conn.Close() + if sendPartial { + _ = conn.WriteMessage(websocket.TextMessage, + []byte(`{"text":"the billing question please","is_final":false}`)) + } + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + })) + t.Cleanup(server.Close) + return server +} + +// newInboundTestPipeline builds just enough of a Pipeline to run +// runInbound: frames arrive on inPCMCh, a confirmed barge-in lands on +// interruptCh, and the agent is marked speaking — the other half of the +// barge-in trigger. The barge-in VAD is the real detector so the 40ms +// onset and 300ms offset timings are the production ones. +func newInboundTestPipeline(t *testing.T, cfg *config.Config) *Pipeline { + t.Helper() + ctx, cancel := context.WithCancel(context.Background()) + t.Cleanup(cancel) + p := &Pipeline{ + ctx: ctx, + cancel: cancel, + cfg: cfg, + bargeInVAD: vad.NewBargeIn(), + inPCMCh: make(chan PCMFrame, 50), + finalCh: make(chan TranscriptEvent, 10), + transcriptCh: make(chan TranscriptEvent, 10), + interruptCh: make(chan struct{}, 1), + } + p.speaking.Store(true) + return p +} + +// startInbound runs runInbound in the background and tears it down at the +// end of the test, so a hung inbound loop fails the test instead of +// hanging it. +func startInbound(t *testing.T, p *Pipeline) { + t.Helper() + done := make(chan struct{}) + go func() { + defer close(done) + p.runInbound() + }() + t.Cleanup(func() { + p.cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("runInbound did not stop within 5s of cancel") + } + }) +} + +// pushPaced feeds one frame every pace. The backchannel window is +// wall-clock time, so an unpaced burst would never let the 600ms window +// elapse no matter how many frames cross. +func pushPaced(t *testing.T, p *Pipeline, samples []int16, count int, pace time.Duration) { + t.Helper() + for i := 0; i < count; i++ { + select { + case p.inPCMCh <- PCMFrame{Samples: samples}: + case <-time.After(5 * time.Second): + t.Fatalf("inbound stalled after %d frames — did the STT client start?", i) + } + time.Sleep(pace) + } +} + +// pushUntilInterrupt keeps feeding frames until a barge-in lands on +// interruptCh. +func pushUntilInterrupt(t *testing.T, p *Pipeline, samples []int16, pace time.Duration) { + t.Helper() + deadline := time.After(10 * time.Second) + for { + select { + case p.inPCMCh <- PCMFrame{Samples: samples}: + case <-deadline: + t.Fatal("timed out waiting for the barge-in interrupt") + } + select { + case <-p.interruptCh: + return + default: + } + select { + case <-time.After(pace): + case <-deadline: + t.Fatal("timed out waiting for the barge-in interrupt") + } + } +} + +// loudSamples is one 20ms frame of near-full-scale PCM: 20000 RMS sits +// far above the barge-in VAD's 1200 base threshold. silentSamples is the +// same frame at zero. +func loudSamples() []int16 { + s := make([]int16, audio.FrameSize) + for i := range s { + s[i] = 20000 + } + return s +} + +func silentSamples() []int16 { + return make([]int16, audio.FrameSize) +} + +func interruptFired(p *Pipeline) bool { + select { + case <-p.interruptCh: + return true + default: + return false + } +} + +// A provider that does not implement stt.PartialsEmitter keeps the +// current path: sustained loud speech with no partial text never opens +// the barge-in window, because hasPartialText still gates the trigger. +func TestInboundBargeInWaitsForPartialTextByDefault(t *testing.T) { + server := fakeVibeVoiceASR(t, false) + + cfg := &config.Config{} + cfg.STT.Provider = "vibevoice" + cfg.VibeVoice.ASRURL = "ws" + strings.TrimPrefix(server.URL, "http") + bar := true + cfg.Pipeline.BargeIn = &bar + + p := newInboundTestPipeline(t, cfg) + startInbound(t, p) + + // 700ms of loud speech — past the 600ms window — with STT never + // reporting a word. + pushPaced(t, p, loudSamples(), 70, 10*time.Millisecond) + time.Sleep(50 * time.Millisecond) + + if interruptFired(p) { + t.Error("barge-in fired without partial text, want the partials gate to hold for a provider that does not implement PartialsEmitter") + } +} + +// With partials flowing, sustained speech over the agent opens the window +// and the interrupt fires once the speech has outlasted the backchannel +// window — the behaviour every partials provider already had. +func TestInboundBargeInFiresOnSustainedSpeechWithPartials(t *testing.T) { + server := fakeVibeVoiceASR(t, true) + + cfg := &config.Config{} + cfg.STT.Provider = "vibevoice" + cfg.VibeVoice.ASRURL = "ws" + strings.TrimPrefix(server.URL, "http") + bar := true + cfg.Pipeline.BargeIn = &bar + + p := newInboundTestPipeline(t, cfg) + startInbound(t, p) + + pushUntilInterrupt(t, p, loudSamples(), 10*time.Millisecond) +} + +// A finals-only provider opens the window on VAD alone, but the interrupt +// must not fire before the full 600ms backchannel window has elapsed: +// there is no text to classify a short burst with, so the window is the +// only thing standing between "mm-hm" and a cut-off agent. +func TestInboundVADOnlyBargeInWaitsOutBackchannelWindow(t *testing.T) { + cfg := &config.Config{} + cfg.STT.Provider = "openai" + cfg.OpenAI.APIKey = "test-key" + bar := true + cfg.Pipeline.BargeIn = &bar + + p := newInboundTestPipeline(t, cfg) + startInbound(t, p) + + // ~250ms of loud speech: the candidate is open, the window (600ms) + // has not run out — nothing may fire yet. + pushPaced(t, p, loudSamples(), 25, 10*time.Millisecond) + if interruptFired(p) { + t.Fatal("VAD-only barge-in fired inside the 600ms backchannel window") + } + + // Keep talking past the window: the interrupt must fire now. + pushUntilInterrupt(t, p, loudSamples(), 10*time.Millisecond) +} + +// A short burst over the agent, ending inside the window, is suppressed: +// without partial text there is no way to tell a real interruption from a +// backchannel, and the degraded mode must stay quiet rather than guess. +func TestInboundVADOnlyBargeInSuppressesShortBurst(t *testing.T) { + cfg := &config.Config{} + cfg.STT.Provider = "openai" + cfg.OpenAI.APIKey = "test-key" + bar := true + cfg.Pipeline.BargeIn = &bar + + p := newInboundTestPipeline(t, cfg) + startInbound(t, p) + + // ~120ms of speech, then silence: the VAD's 300ms offset ends the + // speech inside the window and the burst must classify as + // backchannel. + pushPaced(t, p, loudSamples(), 12, 10*time.Millisecond) + pushPaced(t, p, silentSamples(), 20, 10*time.Millisecond) + time.Sleep(50 * time.Millisecond) + + if interruptFired(p) { + t.Error("short unclassifiable burst fired barge-in, want it suppressed as backchannel") + } +} diff --git a/internal/stt/openai.go b/internal/stt/openai.go index c034610..5b5591d 100644 --- a/internal/stt/openai.go +++ b/internal/stt/openai.go @@ -190,6 +190,12 @@ func (c *openaiClient) Close() { log.Println("[stt:openai] closed") } +// EmitsPartials reports whether this client streams interim results. The +// transcription API is batch-oriented — only finals ever arrive — so this +// is false and the pipeline runs barge-in in its VAD-only degraded mode +// (issue #75). +func (c *openaiClient) EmitsPartials() bool { return false } + // rmsEnergy calculates the root-mean-square energy of linear16 PCM samples. func rmsEnergy(data []byte) float64 { if len(data) < 2 { diff --git a/internal/stt/openai_test.go b/internal/stt/openai_test.go index af77d79..2903fa6 100644 --- a/internal/stt/openai_test.go +++ b/internal/stt/openai_test.go @@ -47,3 +47,23 @@ func TestOpenAITranscriptionUsesConfiguredModel(t *testing.T) { t.Errorf("transcript = %q, want %q", gotTranscript, "hello") } } + +// The transcription API is batch-oriented: only finals ever arrive, so the +// client must report finals-only and let the pipeline fall back to +// VAD-only barge-in (issue #75). +func TestOpenAIClientEmitsNoPartials(t *testing.T) { + client, err := NewOpenAIClient(context.Background(), "test-key", "", func(TranscriptResult) {}) + if err != nil { + t.Fatalf("client: %v", err) + } + defer client.Close() + + var c Client = client + ep, ok := c.(PartialsEmitter) + if !ok { + t.Fatal("openaiClient does not satisfy PartialsEmitter") + } + if ep.EmitsPartials() { + t.Error("openaiClient EmitsPartials = true, want false (batch transcription is finals-only)") + } +} diff --git a/internal/stt/stt.go b/internal/stt/stt.go index 0186e63..b8eb9f7 100644 --- a/internal/stt/stt.go +++ b/internal/stt/stt.go @@ -23,6 +23,24 @@ type Client interface { Close() } +// PartialsEmitter is an optional capability interface for STT providers. +// A provider that streams interim (partial) transcripts while the caller +// is still talking implements it returning true; a finals-only provider — +// one whose transcripts arrive only once the caller has finished — returns +// false. +// +// The pipeline type-asserts a Client against this interface exactly once, +// at construction, and defaults to true when the provider does not +// implement it, so a provider silent on the question keeps the +// partials-driven behaviour. Returning false changes two things: live +// captions show finals only, and barge-in, which can no longer confirm +// real speech from partial text, degrades to firing on VAD alone once the +// full backchannel window has elapsed — with no text, a short burst +// inside the window cannot be classified as anything but backchannel. +type PartialsEmitter interface { + EmitsPartials() bool +} + // NewClient returns an STT client for the configured provider. func NewClient(ctx context.Context, cfg *config.Config, onResult func(TranscriptResult)) (Client, error) { switch cfg.STT.Provider { diff --git a/internal/stt/stt_test.go b/internal/stt/stt_test.go new file mode 100644 index 0000000..718a519 --- /dev/null +++ b/internal/stt/stt_test.go @@ -0,0 +1,23 @@ +package stt + +import "testing" + +// noopClient is a Client that carries no opinion about partials. It is the +// shape of every provider that predates the capability flag. +type noopClient struct{} + +func (noopClient) SendAudio([]byte) error { return nil } +func (noopClient) Close() {} + +// A provider that does not implement PartialsEmitter must not satisfy it: +// runInbound's assertion defaults to true, which is what keeps every +// existing provider on the partials-driven path. If a future edit makes +// the base Client interface carry EmitsPartials (or the method moves to a +// type every client embeds), this test fails before that default silently +// flips for providers that never opted in. +func TestPlainClientDoesNotSatisfyPartialsEmitter(t *testing.T) { + var c Client = noopClient{} + if _, ok := c.(PartialsEmitter); ok { + t.Error("plain Client satisfies PartialsEmitter, want absence to keep the partials default") + } +} diff --git a/internal/stt/telnyx.go b/internal/stt/telnyx.go index 79ada34..253082e 100644 --- a/internal/stt/telnyx.go +++ b/internal/stt/telnyx.go @@ -33,8 +33,9 @@ const ( telnyxInHouseEngine = "Telnyx" // 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. + // any extra configuration; a finals-only default would silently degrade + // both (barge-in to VAD-only, captions to finals only) for anyone who + // just sets the provider and starts talking. telnyxDefaultEngine = "Deepgram" // Client-side endpointing for the in-house engine, which transcribes @@ -98,7 +99,7 @@ func NewTelnyxClient(ctx context.Context, cfg config.TelnyxConfig, onResult func } 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") + log.Printf("[stt] telnyx in-house engine emits finals only: no interim results, so live captions show finals only and barge-in runs VAD-only after the backchannel window") return newTelnyxUtteranceClient(ctx, cfg.APIKey, onResult), nil } return newTelnyxSessionClient(ctx, engine, cfg.APIKey, onResult) @@ -188,6 +189,11 @@ func (c *telnyxSessionClient) Close() { c.cancel() } +// EmitsPartials reports whether this client streams interim results. Every +// hosted engine streams interims exactly like deepgram.go does, which is +// what barge-in's partial-text gate and live captions read. +func (c *telnyxSessionClient) EmitsPartials() bool { return true } + // telnyxUtteranceClient is the in-house-engine half: one socket per // utterance with client-side endpointing. type telnyxUtteranceClient struct { @@ -380,6 +386,12 @@ func (c *telnyxUtteranceClient) Close() { c.wg.Wait() } +// EmitsPartials reports whether this client streams interim results. The +// in-house engine emits exactly one final per utterance and nothing before +// it, so this is false: the pipeline falls back to VAD-only barge-in +// (issue #75). +func (c *telnyxUtteranceClient) EmitsPartials() bool { return false } + // 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 diff --git a/internal/stt/telnyx_test.go b/internal/stt/telnyx_test.go index 02e912e..0da77af 100644 --- a/internal/stt/telnyx_test.go +++ b/internal/stt/telnyx_test.go @@ -465,3 +465,41 @@ func TestTelnyxUtteranceClientSilenceNeverDials(t *testing.T) { default: } } + +// EmitsPartials splits along the same line as the socket lifecycle: the +// hosted engines stream interims, so barge-in keeps its partial-text gate; +// the in-house engine is finals-only, so the pipeline falls back to +// VAD-only barge-in. The exact-case engine string decides, consistent +// with how NewTelnyxClient routes. +func TestTelnyxClientEmitsPartialsByEngine(t *testing.T) { + f := newFakeTelnyxSTT(t, func(conn *websocket.Conn) { + for { + if _, _, err := conn.ReadMessage(); err != nil { + return + } + } + }) + overrideTelnyxSTTURL(t, f.wsURL()) + + session, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k", TranscriptionEngine: "Deepgram"}, func(TranscriptResult) {}) + if err != nil { + t.Fatalf("session client: %v", err) + } + defer session.Close() + + if ep, ok := session.(PartialsEmitter); !ok || !ep.EmitsPartials() { + t.Error("hosted-engine client EmitsPartials = false, want true (it streams interims)") + } + + utterance, err := NewTelnyxClient(context.Background(), + config.TelnyxConfig{APIKey: "k", TranscriptionEngine: "Telnyx"}, func(TranscriptResult) {}) + if err != nil { + t.Fatalf("utterance client: %v", err) + } + defer utterance.Close() + + if ep, ok := utterance.(PartialsEmitter); !ok || ep.EmitsPartials() { + t.Error("in-house-engine client EmitsPartials = true, want false (one final per utterance, nothing before it)") + } +}