From 68e1fedd2d0d0bbc2ac5c40feeee3345a883bf42 Mon Sep 17 00:00:00 2001 From: Brooklyn Nicholson Date: Wed, 22 Jul 2026 17:47:15 -0500 Subject: [PATCH] feat(voice): stream turn deltas through the TUI gateway, with barge-in MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Turn deltas feed a per-turn TTS pipeline; the post-complete speak_text call survives only as a fallback. session.interrupt, /voice toggles, and new turns cut in-flight speech. VAD barge-in emits voice.interrupted at detection, then the captured interruption goes out as voice.transcript — the same event the TUI already submits as a spoken turn. --- tests/test_tui_gateway_server.py | 129 ++++++++++++++++++++++++++++++- tui_gateway/server.py | 128 ++++++++++++++++++++++++++++-- 2 files changed, 251 insertions(+), 6 deletions(-) diff --git a/tests/test_tui_gateway_server.py b/tests/test_tui_gateway_server.py index 9a197f13acd6..e5169339c8c2 100644 --- a/tests/test_tui_gateway_server.py +++ b/tests/test_tui_gateway_server.py @@ -1209,7 +1209,10 @@ def test_voice_toggle_tts_branch_also_carries_record_key(monkeypatch): ), ) monkeypatch.setenv("HERMES_VOICE", "1") - monkeypatch.delenv("HERMES_VOICE_TTS", raising=False) + # setenv (not delenv) — the handler writes HERMES_VOICE_TTS directly, and + # delenv on an absent var registers no teardown, leaking TTS=1 into every + # later test in the file (which now spins up the streaming TTS pipeline). + monkeypatch.setenv("HERMES_VOICE_TTS", "0") tts_resp = server.dispatch( {"id": "voice-tts", "method": "voice.toggle", "params": {"action": "tts"}} @@ -11733,3 +11736,127 @@ def test_get_usage_clamps_post_compression_sentinel(): usage = server._get_usage(agent) assert "context_used" not in usage assert "context_percent" not in usage + + +# --------------------------------------------------------------------------- +# Streaming TTS — per-turn pipeline + barge-in +# --------------------------------------------------------------------------- + +def _fake_tts_modules(monkeypatch, *, requirements=True, playback_stops=None, listen=None, transcribe=None): + """Install lightweight tools.tts_tool / tools.voice_mode fakes.""" + started = {} + + def fake_stream(text_queue, stop, done, **_kw): + started["queue"] = text_queue + stop.wait(5) + done.set() + + def default_listen(should_stop, capture=False, on_trigger=None, **_kw): + return None if capture else False + + monkeypatch.setitem( + sys.modules, + "tools.tts_tool", + types.SimpleNamespace( + check_tts_requirements=lambda: requirements, + stream_tts_to_speaker=fake_stream, + ), + ) + monkeypatch.setitem( + sys.modules, + "tools.voice_mode", + types.SimpleNamespace( + stop_playback=lambda: (playback_stops.append(True) if playback_stops is not None else None), + listen_for_speech=listen or default_listen, + transcribe_recording=transcribe or (lambda path, model=None: {"success": True, "transcript": ""}), + ), + ) + return started + + +def test_tts_stream_begin_requires_voice_tts(monkeypatch): + monkeypatch.setenv("HERMES_VOICE_TTS", "0") + assert server._tts_stream_begin() is None + + +def test_tts_stream_begin_requires_working_provider(monkeypatch): + monkeypatch.setenv("HERMES_VOICE_TTS", "1") + _fake_tts_modules(monkeypatch, requirements=False) + assert server._tts_stream_begin() is None + + +def test_tts_stream_begin_and_stop_lifecycle(monkeypatch): + """begin() spawns the consumer; stop() cuts it and clears the slot.""" + monkeypatch.setenv("HERMES_VOICE_TTS", "1") + monkeypatch.setenv("HERMES_VOICE", "0") # no barge-in monitor (no mic) + playback_stops: list = [] + started = _fake_tts_modules(monkeypatch, playback_stops=playback_stops) + + text_queue = server._tts_stream_begin() + assert text_queue is not None + assert started["queue"] is text_queue + + with server._tts_stream_lock: + state = server._tts_stream_state + assert state is not None and not state["stop"].is_set() + + server._tts_stream_stop() + assert state["stop"].is_set() + assert playback_stops == [True] + with server._tts_stream_lock: + assert server._tts_stream_state is None + + +def test_tts_stream_begin_barges_in_on_previous_pipeline(monkeypatch): + """A new turn's pipeline stops the previous turn's speech (one speaker).""" + monkeypatch.setenv("HERMES_VOICE_TTS", "1") + monkeypatch.setenv("HERMES_VOICE", "0") + _fake_tts_modules(monkeypatch) + + server._tts_stream_begin() + with server._tts_stream_lock: + first = server._tts_stream_state + server._tts_stream_begin() + assert first is not None and first["stop"].is_set() + server._tts_stream_stop() + + +def test_tts_stream_vad_barge_in_cuts_pipeline_and_submits_capture(monkeypatch, tmp_path): + """User speech during playback cuts TTS at the moment of detection + (voice.interrupted), then the captured interruption is transcribed and + emitted as voice.transcript so the TUI submits it — complete from its + first syllable, no re-record round trip.""" + monkeypatch.setenv("HERMES_VOICE_TTS", "1") + monkeypatch.setenv("HERMES_VOICE", "1") + monkeypatch.setattr(server, "_load_cfg", lambda: {"voice": {"barge_in": True}}) + events: list = [] + monkeypatch.setattr( + server, "_voice_emit", lambda event, payload=None: events.append((event, payload)) + ) + + wav = tmp_path / "barge.wav" + wav.write_bytes(b"RIFF") + + def fake_listen(should_stop, capture=False, on_trigger=None, **_kw): + assert capture is True + on_trigger() # playback cut happens at detection, not after endpointing + return str(wav) + + _fake_tts_modules( + monkeypatch, + listen=fake_listen, + transcribe=lambda path, model=None: {"success": True, "transcript": "stop, actually—"}, + ) + + server._tts_stream_begin() + with server._tts_stream_lock: + state = server._tts_stream_state + assert state is not None + assert state["stop"].wait(2.0) + deadline = time.monotonic() + 2.0 + while time.monotonic() < deadline and wav.exists(): + time.sleep(0.01) # unlink (finally) runs after the transcript emit + assert ("voice.interrupted", None) in events + assert ("voice.transcript", {"text": "stop, actually—"}) in events + assert not wav.exists() # capture temp file cleaned up + server._tts_stream_stop() diff --git a/tui_gateway/server.py b/tui_gateway/server.py index f10eba7fdddc..ba7372673835 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -9372,6 +9372,9 @@ def _(rid, params: dict) -> dict: @method("session.interrupt") def _(rid, params: dict) -> dict: + # Keypress barge-in: stopping the turn also silences its streaming TTS + # (voice is process-global, so no per-session scoping is needed). + _tts_stream_stop() session, err = _sess_nowait(params, rid) if err: return err @@ -10302,6 +10305,7 @@ def _run_prompt_submit(rid, sid: str, session: dict, text: Any) -> None: session_tokens = [] home_token = None # per-turn HERMES_HOME override for a resumed remote profile goal_followup = None # set by the post-turn goal hook below + tts_queue = None # streaming-TTS feed for this turn (voice mode) one_turn_restore = session.pop("one_turn_model_restore", None) try: from tools.approval import ( @@ -10424,12 +10428,18 @@ def _run_prompt_submit(rid, sid: str, session: dict, text: Any) -> None: else: run_message = _enrich_with_attached_images(prompt, images) + # Streaming TTS: voice-mode replies are spoken sentence-by-sentence + # as tokens arrive (CLI parity) instead of after the full turn. + tts_queue = _tts_stream_begin() + def _stream(delta): with session["history_lock"]: _append_inflight_delta(session, delta) payload = {"text": delta} if streamer and (r := streamer.feed(delta)) is not None: payload["rendered"] = r + if tts_queue is not None and isinstance(delta, str): + tts_queue.put(delta) _emit("message.delta", sid, payload) # Surface interim assistant text (commentary emitted alongside @@ -10698,12 +10708,13 @@ def _run_prompt_submit(rid, sid: str, session: dict, text: Any) -> None: except Exception: pass - # CLI parity: when voice-mode TTS is on, speak the agent reply - # (cli.py:_voice_speak_response). Only the final text — tool - # calls / reasoning already stream separately and would be - # noisy to read aloud. + # Voice TTS fallback: when the streaming pipeline couldn't start + # (no provider / missing deps probed at turn start), speak the + # final text whole (cli.py:_voice_speak_response parity). The + # streaming path already spoke everything via tts_queue. if ( status == "complete" + and tts_queue is None and isinstance(raw, str) and raw.strip() and _voice_tts_enabled() @@ -10738,6 +10749,8 @@ def _run_prompt_submit(rid, sid: str, session: dict, text: Any) -> None: ) _emit("error", sid, {"message": str(e)}) finally: + if tts_queue is not None: + tts_queue.put(None) # end-of-text sentinel — flush + finish speaking if one_turn_restore: try: _restore_agent_model_runtime(agent, one_turn_restore) @@ -15512,6 +15525,107 @@ def _voice_tts_enabled() -> bool: return os.environ.get("HERMES_VOICE_TTS", "").strip() == "1" +# ── Streaming TTS (one active pipeline per process — one speaker) ────────── +# Token deltas from the running turn feed a sentence-buffering consumer +# (tools.tts_tool.stream_tts_to_speaker) so speech starts on the first +# sentence instead of after the full reply. Voice is process-global, so a +# single slot suffices; starting a new turn's pipeline barges in on the +# previous one. + +_tts_stream_lock = threading.Lock() +_tts_stream_state: Optional[dict] = None + + +def _tts_stream_begin() -> Optional[queue.Queue]: + """Start a per-turn streaming TTS consumer; None when TTS can't stream.""" + if not _voice_tts_enabled(): + return None + try: + from tools.tts_tool import check_tts_requirements, stream_tts_to_speaker + + if not check_tts_requirements(): + return None + except Exception: + return None + + _tts_stream_stop() + text_queue: queue.Queue = queue.Queue() + stop = threading.Event() + done = threading.Event() + threading.Thread( + target=stream_tts_to_speaker, args=(text_queue, stop, done), daemon=True + ).start() + + global _tts_stream_state + with _tts_stream_lock: + _tts_stream_state = {"stop": stop, "done": done} + + if _voice_mode_enabled() and _voice_cfg_dict().get("barge_in", True): + threading.Thread( + target=_tts_stream_barge_in_monitor, args=(stop, done), daemon=True + ).start() + + return text_queue + + +def _tts_stream_stop() -> None: + """Barge-in: cut any in-flight streaming TTS (new turn, interrupt, /voice off).""" + global _tts_stream_state + with _tts_stream_lock: + state, _tts_stream_state = _tts_stream_state, None + if state is None: + return + state["stop"].set() + try: + from tools.voice_mode import stop_playback + + stop_playback() + except Exception: + pass + + +def _tts_stream_barge_in_monitor(stop: threading.Event, done: threading.Event) -> None: + """VAD barge-in: cut streaming TTS when the user starts talking. + + Playback is cut at the moment of detection while the monitor keeps + capturing (with pre-roll) until the user goes quiet — the interruption is + then transcribed and emitted as ``voice.transcript``, which the TUI + submits like any spoken turn. Without capture the opening words would be + lost between detection and the next recording start. + """ + try: + from tools.voice_mode import listen_for_speech, stop_playback, transcribe_recording + + barged = threading.Event() + + def _cut_playback(): + if not done.is_set(): + barged.set() + stop.set() + stop_playback() + _voice_emit("voice.interrupted") + + wav_path = listen_for_speech( + lambda: stop.is_set() or done.is_set(), + capture=True, + on_trigger=_cut_playback, + ) + if not (wav_path and barged.is_set()): + return + try: + result = transcribe_recording(wav_path) + text = (result.get("transcript") or "").strip() if result.get("success") else "" + if text: + _voice_emit("voice.transcript", {"text": text}) + finally: + try: + os.unlink(wav_path) + except OSError: + pass + except Exception as e: + logger.debug("TTS barge-in monitor failed: %s", e) + + def _voice_cfg_dict() -> dict: """Shape-safe accessor for the ``voice:`` block in config.yaml. @@ -15599,8 +15713,10 @@ def _(rid, params: dict) -> dict: except Exception as e: logger.warning("voice: stop_continuous failed during toggle off: %s", e) - # Clear TTS so it can be toggled independently after voice is off. + # Clear TTS so it can be toggled independently after voice is off, + # and silence any in-flight streaming speech. os.environ["HERMES_VOICE_TTS"] = "0" + _tts_stream_stop() return _ok( rid, @@ -15617,6 +15733,8 @@ def _(rid, params: dict) -> dict: new_value = not _voice_tts_enabled() # Runtime-only flag (CLI parity) — see voice.toggle on/off above. os.environ["HERMES_VOICE_TTS"] = "1" if new_value else "0" + if not new_value: + _tts_stream_stop() # Include ``record_key`` on every branch so a /voice tts toggle # doesn't reset the TUI's cached shortcut to the default when a # user has a custom binding configured (Copilot review, round 2