"""Audio dashboard routes: transcription upload, voice config, ElevenLabs voices, TTS speak/lease and the speak-stream WebSocket.

Extracted from ``hermes_cli.web_server``; helpers/state that tests monkeypatch on
``web_server`` stay there and are late-bound (cycle-safe).
"""

import base64
import binascii
import contextlib
import logging
import queue
import tempfile
import threading
import asyncio
import json
import os
import urllib.parse
import urllib.request
from fastapi import APIRouter
from hermes_cli.web_routers._common import http_failure
from hermes_cli.web_deps import late
from hermes_cli.web_server_chat import _ws_auth_ok, _ws_request_is_allowed
from hermes_cli.web_server_gateway import _split_text_for_speak_stream
from fastapi import HTTPException, WebSocket, WebSocketDisconnect
from hermes_cli.web_models import (
    AudioTranscriptionRequest,
    STTLeaseRequest,
    TTSSpeakRequest,
    TTSLeaseRequest,
    VoiceLiveSessionRequest,
)
from typing import Any, Dict, Optional

_log = logging.getLogger("hermes_cli.web_server")
router = APIRouter()

# Late-bound so a test's monkeypatch on the owning module wins at call time.
_config_profile_scope = late("_config_profile_scope", "hermes_cli.web_server_profiles")
_voice_list_error_logged_once = late("_voice_list_error_logged_once")
load_env = late("load_env", "hermes_cli.config")

_AUDIO_MIME_EXTENSIONS: Dict[str, str] = {
    "audio/aac": ".aac", "audio/flac": ".flac", "audio/m4a": ".m4a", "audio/mp3": ".mp3",
    "audio/mp4": ".mp4", "audio/mpeg": ".mp3", "audio/ogg": ".ogg", "audio/wav": ".wav",
    "audio/wave": ".wav", "audio/webm": ".webm", "audio/x-m4a": ".m4a", "audio/x-wav": ".wav",
    "video/webm": ".webm",
}

_MAX_TRANSCRIPTION_UPLOAD_BYTES = 25 * 1024 * 1024

_SPEAK_MIME_BY_EXT = {
    ".mp3": "audio/mpeg", ".ogg": "audio/ogg", ".opus": "audio/ogg", ".wav": "audio/wav",
    ".flac": "audio/flac",
}


def _unlink_quietly(path: str) -> None:
    try:
        os.unlink(path)
    except OSError:
        pass


async def _run_config_scoped(profile: Optional[str], fn):
    """Run ``fn()`` on a worker thread under the config-only profile scope.

    Home-only contextvar scope, NOT ``_profile_scope``: these calls block for a
    provider round-trip and only need config/.env resolution, while
    ``_profile_scope`` holds a process-global skills lock for its entire body.
    """
    def _scoped():
        with _config_profile_scope(profile):
            return fn()

    return await asyncio.get_running_loop().run_in_executor(None, _scoped)


def _audio_extension_for_mime(mime_type: str) -> str:
    normalized = (mime_type or "").split(";", 1)[0].strip().lower()
    return _AUDIO_MIME_EXTENSIONS.get(normalized, ".webm")


@router.post("/api/audio/transcribe")
async def transcribe_audio_upload(
    payload: AudioTranscriptionRequest, profile: Optional[str] = None
):
    data_url = (payload.data_url or "").strip()
    if not data_url.startswith("data:") or "," not in data_url:
        raise HTTPException(status_code=400, detail="Invalid audio payload")

    header, encoded = data_url.split(",", 1)
    if ";base64" not in header:
        raise HTTPException(status_code=400, detail="Audio payload must be base64 encoded")

    mime_type = (payload.mime_type or header[5:].split(";", 1)[0] or "audio/webm").strip()
    normalized_mime_type = mime_type.split(";", 1)[0].lower()
    if not (normalized_mime_type.startswith("audio/") or normalized_mime_type == "video/webm"):
        raise HTTPException(status_code=400, detail="Payload must be an audio recording")

    try:
        audio_bytes = base64.b64decode(encoded, validate=True)
    except (binascii.Error, ValueError):
        raise HTTPException(status_code=400, detail="Audio payload is not valid base64")

    if not audio_bytes:
        raise HTTPException(status_code=400, detail="Audio recording is empty")
    if len(audio_bytes) > _MAX_TRANSCRIPTION_UPLOAD_BYTES:
        raise HTTPException(status_code=413, detail="Audio recording is too large")

    temp_path = ""
    try:
        with http_failure("Desktop voice transcription failed", 500, "Transcription failed"):
            with tempfile.NamedTemporaryFile(
                prefix="hermes-desktop-voice-", suffix=_audio_extension_for_mime(mime_type),
                delete=False,
            ) as tmp:
                tmp.write(audio_bytes)
                temp_path = tmp.name

            # transcribe_recording (not raw transcribe_audio): filters Whisper
            # hallucinations and maps provider "empty transcript" errors to a
            # successful empty result — the live voice loop treats "" as silence
            # and re-listens instead of surfacing a 400 on every quiet turn.
            from tools.voice_mode import transcribe_recording

            result = await _run_config_scoped(profile, lambda: transcribe_recording(temp_path))
    finally:
        if temp_path:
            _unlink_quietly(temp_path)

    if not result.get("success"):
        err = result.get("error") or "Transcription failed"
        # No speech detected is a normal outcome for VAD/continuous voice loops
        # (re-listening on silence), not an error: return an empty transcript so
        # the client quietly re-listens instead of showing a failure toast.
        if "empty transcript" in err.lower():
            return {"ok": True, "transcript": "", "provider": result.get("provider")}
        raise HTTPException(status_code=400, detail=err)

    return {
        "ok": True, "transcript": str(result.get("transcript") or "").strip(),
        "provider": result.get("provider"),
    }


@router.get("/api/audio/voice-config")
async def get_client_voice_config(profile: Optional[str] = None):
    """The active profile's STT/TTS config for CLIENT-DIRECT voice.

    Lets the desktop cut the audio relay hop: mic audio goes straight to the
    profile's STT provider and reply text is synthesized on the client with
    the profile's TTS provider — the desktop↔gateway link carries only text.
    Providers that can only run on this host (local whisper, edge-tts,
    command/plugin providers) resolve to ``{"mode": "relay"}`` and the
    desktop keeps using the /api/audio/* relay endpoints.

    Same trust boundary as every profile-scoped route: the caller is an
    authenticated client that can already drive the agent. Keys in the
    response are held in client memory only, never persisted client-side.
    Gate: ``voice.client_direct`` in config.yaml (default true).
    """
    from tools.voice_client_config import resolve_client_voice_config
    try:
        result = await _run_config_scoped(profile, resolve_client_voice_config)
    except HTTPException:
        raise  # an unknown ?profile= is the scope's 404, not a reason to fall back to relay
    except Exception:
        _log.exception("Client voice-config resolution failed")
        fallback = {"mode": "relay", "reason": "resolution error"}
        return {"ok": True, "stt": fallback, "tts": dict(fallback)}

    return {"ok": True, **result}


@router.get("/api/audio/voice-live/status")
async def get_voice_live_status(profile: Optional[str] = None):
    """Which voice chat mode the profile selected (``chained`` | ``gpt-live``) and whether GPT-Live
    can start. Non-secret: the desktop decides which conversation engine to mount from this."""
    from tools.voice_live import resolve_gpt_live_status
    with http_failure("GPT-Live status resolution failed", 500, "GPT-Live status failed"):
        result = await _run_config_scoped(profile, resolve_gpt_live_status)
    return {"ok": True, **result}


@router.post("/api/audio/voice-live/session")
async def create_voice_live_session(payload: VoiceLiveSessionRequest, profile: Optional[str] = None):
    """Exchange the renderer's WebRTC SDP offer for a GPT-Live session answer.

    The project API key stays on this host; the renderer only receives the session id and the
    SDP answer. Client delegation is fixed at creation: every ``session.delegation.created`` the
    renderer receives becomes a Hermes turn on the session it belongs to.
    """
    from tools.voice_live import create_webrtc_session
    # Validate emptiness only: the vendor's SDP parser needs the offer byte-exact, including the
    # trailing CRLF (a stripped offer answers 400 "failed to unmarshal SDP: EOF").
    sdp = payload.sdp or ""
    if not sdp.strip():
        raise HTTPException(status_code=400, detail="An SDP offer is required")
    try:
        result = await _run_config_scoped(profile, lambda: create_webrtc_session(sdp, payload.history))
    except ValueError as exc:
        raise HTTPException(status_code=503, detail=str(exc))
    except RuntimeError as exc:
        raise HTTPException(status_code=502, detail=str(exc))
    return {"ok": True, **result}


def _elevenlabs_voice_label(voice: Dict[str, Any]) -> str:
    name = str(voice.get("name") or voice.get("voice_id") or "Voice").strip()
    category = str(voice.get("category") or "").strip()

    return f"{name} ({category})" if category else name


@router.get("/api/audio/elevenlabs/voices")
async def get_elevenlabs_voices(profile: Optional[str] = None):
    """Return ElevenLabs voices when an API key is configured.

    The desktop UI uses this for the ``tts.elevenlabs.voice_id`` dropdown.
    Only non-secret voice metadata is returned; the API key stays server-side.
    """
    # Config-only scope (await-safe): the key lookup reads the requested
    # profile's .env, matching the profile the settings UI writes to.
    with _config_profile_scope(profile):
        api_key = (load_env().get("ELEVENLABS_API_KEY") or "").strip()
    if not api_key:
        # Fallback for env-only deployments — scope-aware: under multiplex
        # os.environ may hold another profile's key, so honor the installed
        # scope's verdict. Only the unscoped default-profile path
        # (UnscopedSecretError) reads the env; any other failure stays empty.
        try:
            from agent.secret_scope import UnscopedSecretError, get_secret

            try:
                api_key = (get_secret("ELEVENLABS_API_KEY") or "").strip()
            except UnscopedSecretError:
                api_key = (os.environ.get("ELEVENLABS_API_KEY") or "").strip()
        except Exception:
            pass
    if not api_key:
        return {"available": False, "voices": []}

    request = urllib.request.Request(
        "https://api.elevenlabs.io/v1/voices",
        headers={"Accept": "application/json", "xi-api-key": api_key},
    )

    try:
        loop = asyncio.get_running_loop()

        def _fetch() -> Dict[str, Any]:
            with urllib.request.urlopen(request, timeout=10) as response:
                return json.loads(response.read().decode("utf-8"))

        payload = await loop.run_in_executor(None, _fetch)
    except urllib.error.HTTPError as exc:
        # An auth failure (bad/expired/scoped key) is a persistent, user-fixable
        # state and the desktop polls this on every settings open/focus, so
        # treat 401/403 as "integration unavailable": 200 to the UI and log at
        # most once until the error signature changes.
        if exc.code in (401, 403):
            if _voice_list_error_logged_once(f"http-{exc.code}"):
                _log.info("ElevenLabs voices unavailable: %s — check ELEVENLABS_API_KEY", exc)
            return {"available": False, "voices": [], "error": "unauthorized"}
        if _voice_list_error_logged_once(f"http-{exc.code}"):
            _log.warning("ElevenLabs voice list failed: %s", exc)
        raise HTTPException(status_code=502, detail="Could not load ElevenLabs voices")
    except Exception as exc:
        if _voice_list_error_logged_once(str(exc)):
            _log.warning("ElevenLabs voice list failed: %s", exc)
        raise HTTPException(status_code=502, detail="Could not load ElevenLabs voices")
    _voice_list_error_logged_once(None)  # success — re-arm logging for next failure

    voices = []
    for voice in payload.get("voices") or []:
        if not isinstance(voice, dict):
            continue

        voice_id = str(voice.get("voice_id") or "").strip()
        if not voice_id:
            continue

        voices.append({
            "voice_id": voice_id, "name": str(voice.get("name") or voice_id),
            "label": _elevenlabs_voice_label(voice),
        })

    voices.sort(key=lambda item: str(item.get("label") or "").lower())
    return {"available": True, "voices": voices}


@router.post("/api/audio/speak")
async def speak_text(payload: TTSSpeakRequest, profile: Optional[str] = None):
    """Synthesize speech and return audio as base64 data URL.

    Used by the desktop voice-conversation mode to play back assistant
    responses without exposing the on-disk file path; reuses the TTS provider
    chain configured under ``tts.`` in config.yaml.
    """
    text = (payload.text or "").strip()
    if not text:
        raise HTTPException(status_code=400, detail="Text is required")

    # _config_profile_scope raises 400/404 for a bad profile — pass it
    # through instead of masking it as a 500 synthesis failure.
    with http_failure("Desktop voice TTS failed", 500, "Speech synthesis failed"):
        from tools.tts_tool import text_to_speech_tool

        result_json = await _run_config_scoped(profile, lambda: text_to_speech_tool(text))

    try:
        result = json.loads(result_json) if isinstance(result_json, str) else result_json
    except Exception:
        raise HTTPException(status_code=500, detail="Invalid TTS response")

    if not result.get("success"):
        raise HTTPException(
            status_code=400, detail=result.get("error") or "Speech synthesis failed",
        )

    file_path = result.get("file_path")
    if not file_path or not os.path.isfile(file_path):
        raise HTTPException(status_code=500, detail="Audio file missing")

    mime_type = _SPEAK_MIME_BY_EXT.get(os.path.splitext(file_path)[1].lower(), "audio/mpeg")

    def _read_and_unlink() -> bytes:
        # Off-loop: synthesized audio can be several MB; reading it inline
        # blocks the uvicorn event loop. Unlink rides the same thread hop so
        # the temp file cannot outlive an early return.
        try:
            with open(file_path, "rb") as fh:
                return fh.read()
        finally:
            _unlink_quietly(file_path)

    try:
        audio_bytes = await asyncio.to_thread(_read_and_unlink)
    except OSError as exc:
        raise HTTPException(status_code=500, detail=f"Could not read audio: {exc}")

    encoded = base64.b64encode(audio_bytes).decode("ascii")
    return {
        "ok": True, "data_url": f"data:{mime_type};base64,{encoded}", "mime_type": mime_type,
        "provider": result.get("provider"),
    }


@router.post("/api/audio/tts-lease")
async def tts_lease(payload: TTSLeaseRequest, profile: Optional[str] = None):
    """Desktop TTS-output toggles as warm-up / release signals.

    ``active: true`` registers a lease on the TTS engine and pre-loads the
    configured provider (local model, lazily-installed SDK) so the first spoken
    reply doesn't pay the load as dead air; ``active: false`` drops the lease
    and, once no surface holds one, unloads resident local models. Blocking
    work runs off the event loop. Warm-up failures are reported in the body,
    never as an HTTP error — the toggle must succeed even when preload fails.
    """
    lease = (payload.lease or "").strip()
    if not lease:
        raise HTTPException(status_code=400, detail="lease is required")

    def _apply():
        from tools.tts_tool_lifecycle import acquire_tts_lease, release_tts_lease
        if payload.active:
            with _config_profile_scope(profile):
                return acquire_tts_lease(lease)
        # Release reads the requester's keep_warm_seconds, but must drop the lease even when
        # that profile is gone — a stuck lease pins the local model in memory.
        try:
            with _config_profile_scope(profile):
                return release_tts_lease(lease)
        except HTTPException:
            return release_tts_lease(lease)

    try:
        result = await asyncio.get_running_loop().run_in_executor(None, _apply)
    except HTTPException:
        raise
    except Exception as exc:
        _log.warning("TTS lease %s (%s) failed: %s", lease, payload.active, exc)
        result = {"leases": None, "action": "error", "error": str(exc)}
    return {"ok": True, "lease": lease, "active": payload.active, **result}


class _SyncSentencePCMStreamer:
    """Speak one sentence through the sync TTS tool and yield int16 mono PCM.

    Not a second provider: ``text_to_speech_tool`` is the same stack the CLI
    speaker uses for edge and every other non-chunked provider. The desktop
    socket only plays PCM, so the written file is decoded before it is sent.
    ``sample_rate`` is updated before the first yield, matching the chunked
    streamers whose rate is only final once synthesis has answered.
    """

    sample_rate = 24000
    channels = 1

    def stream(self, text: str):
        pcm, rate = _sync_sentence_to_pcm(text)
        if rate:
            self.sample_rate = rate
        if pcm:
            yield pcm


def _sync_sentence_to_pcm(text: str) -> tuple:
    """Synthesize *text* with the configured sync provider and return PCM."""
    from tools.tts_tool import text_to_speech_tool

    fd, tmp_path = tempfile.mkstemp(suffix=".mp3")
    os.close(fd)
    extra = None
    try:
        raw = text_to_speech_tool(text=text, output_path=tmp_path)
        try:
            payload = json.loads(raw) if isinstance(raw, str) else raw
        except (TypeError, ValueError) as exc:
            raise RuntimeError("Invalid TTS response") from exc
        if not isinstance(payload, dict) or not payload.get("success"):
            detail = str(payload.get("error") or "") if isinstance(payload, dict) else ""
            raise RuntimeError(detail or "Speech synthesis failed")
        written = payload.get("file_path") or tmp_path
        if not isinstance(written, str) or not os.path.isfile(written) or os.path.getsize(written) <= 0:
            raise RuntimeError("Audio file missing")
        if os.path.abspath(written) != os.path.abspath(tmp_path):
            extra = written
        pcm, rate = _audio_file_to_pcm(written)
        if not pcm:
            raise RuntimeError("TTS audio decoded to silence")
        return pcm, rate
    finally:
        _unlink_quietly(tmp_path)
        if extra:
            _unlink_quietly(extra)


def _audio_file_to_pcm(path: str) -> tuple:
    """Decode *path* to int16 mono PCM. WAV via the stdlib; anything else via ffmpeg."""
    with open(path, "rb") as fh:
        head = fh.read(12)
    if len(head) >= 12 and head.startswith(b"RIFF") and head[8:12] == b"WAVE":
        pcm, rate = _wav_s16le_mono(path)
        if pcm:
            return pcm, rate
    return _ffmpeg_s16le_mono(path)


def _wav_s16le_mono(path: str) -> tuple:
    import wave

    try:
        with wave.open(path, "rb") as wf:
            if wf.getsampwidth() != 2 or wf.getnchannels() != 1 or wf.getframerate() <= 0:
                return b"", 0
            return wf.readframes(wf.getnframes()), int(wf.getframerate())
    except (wave.Error, EOFError, OSError):
        return b"", 0


def _ffmpeg_s16le_mono(path: str) -> tuple:
    import shutil

    from tools.tts_tool_delivery import _ffmpeg_run

    ffmpeg = shutil.which("ffmpeg")
    if not ffmpeg:
        raise RuntimeError("ffmpeg is required to decode TTS audio for speak-stream")
    result = _ffmpeg_run(
        ffmpeg,
        ["-i", path, "-f", "s16le", "-ac", "1", "-ar", "24000", "-loglevel", "error", "pipe:1"],
        timeout=60,
    )
    if result.returncode != 0 or not result.stdout:
        stderr = (result.stderr or b"").decode("utf-8", "replace")[:200]
        raise RuntimeError(f"TTS audio decode failed: {stderr}")
    return result.stdout, 24000


@router.post("/api/audio/stt-lease")
async def stt_lease(payload: STTLeaseRequest, profile: Optional[str] = None):
    """Desktop voice-input sessions as STT warm-up / release signals.

    ``active: true`` registers a lease and pre-loads the configured local STT
    model (first-use download + load) so the transcription request doesn't pay
    the cold cost inside its timeout; ``active: false`` drops the lease. The
    model stays resident after the last release — it is shared with the
    gateway/CLI surfaces in this process, and ``stt.local.unload_after_idle_seconds``
    still governs eviction. Blocking work runs off the event loop. Warm-up
    failures are reported in the body, never as an HTTP error — recording must
    start even when preload fails.
    """
    lease = (payload.lease or "").strip()
    if not lease:
        raise HTTPException(status_code=400, detail="lease is required")

    def _apply():
        from tools.stt_lease import acquire_stt_lease, release_stt_lease
        if payload.active:
            with _config_profile_scope(profile):
                return acquire_stt_lease(lease)
        return release_stt_lease(lease)

    try:
        result = await asyncio.get_running_loop().run_in_executor(None, _apply)
    except HTTPException:
        raise
    except Exception as exc:
        _log.warning("STT lease %s (%s) failed: %s", lease, payload.active, exc)
        result = {"leases": None, "action": "error", "error": str(exc)}
    return {"ok": True, "lease": lease, "active": payload.active, **result}


@router.websocket("/api/audio/speak-stream")
async def speak_stream_ws(ws: "WebSocket") -> None:
    """Streaming TTS for the desktop: text in, raw int16 PCM frames out.

    The socket is a per-reply speech *session*: the client feeds text
    incrementally as LLM deltas arrive, the server cuts sentences
    (``SentenceChunker`` — same cutter as the CLI/TUI speaker pipeline) and
    streams each one's PCM the moment it's ready, so speech overlaps generation.

    Protocol:
      client → ``{"text": "..."}`` frames (incremental; may combine with done),
               ``{"done": true}`` when the reply is complete,
               ``{"stop": true}`` or disconnect = barge-in
      server → ``{"type": "start", "sample_rate": N, "channels": 1}`` (sent
               with the first PCM frame, once the provider's rate is final),
               binary PCM frames, then ``{"type": "end"}``
      server → ``{"type": "fallback"}`` only when sentence synthesis produced
               no audio. Providers with no chunked API (edge, the default)
               still speak per sentence via ``text_to_speech_tool`` and stream
               that PCM — the client POST is the last resort, not the Edge path.
    """
    if not _ws_auth_ok(ws):
        await ws.close(code=4401)
        return
    if not _ws_request_is_allowed(ws):
        await ws.close(code=4403)
        return
    await ws.accept()

    # Profile via query param, like /api/pty and /api/console: the provider
    # chain + API keys must resolve from the requesting profile's config, not
    # the dashboard's own — at resolve time AND in the synthesis thread.
    profile = (ws.query_params.get("profile") or "").strip() or None

    loop = asyncio.get_running_loop()

    def _resolve():
        from tools.tts_streaming import resolve_streaming_provider
        from tools.tts_tool import _get_provider, _load_tts_config, _resolve_max_text_length
        with _config_profile_scope(profile):
            cfg = _load_tts_config()
            streamer = resolve_streaming_provider(cfg)
            cap = _resolve_max_text_length(_get_provider(cfg), cfg)
        return streamer, cap, cfg

    try:
        streamer, cap, cfg = await loop.run_in_executor(None, _resolve)
    except Exception:
        _log.exception("speak-stream provider resolution failed")
        streamer, cap, cfg = None, 0, {}
    if streamer is None:
        # Edge (the default) and every other non-chunked provider still have a
        # documented per-sentence path. type=fallback here is what makes Desktop
        # wait for the whole reply and POST it to /api/audio/speak.
        streamer = _SyncSentencePCMStreamer()

    # The start frame is deferred until the first PCM chunk (or end-of-speech):
    # the OpenAI-compatible streamer only learns the endpoint's real rate from
    # the response headers inside stream(), and the client opens its
    # AudioContext at whatever rate the start frame carries.
    start_sent = False

    async def _send_start():
        nonlocal start_sent
        if start_sent:
            return
        start_sent = True
        await ws.send_json(
            {"type": "start", "sample_rate": streamer.sample_rate, "channels": streamer.channels}
        )

    stop = threading.Event()
    produced_audio = False
    synthesis_failed = False
    text_q: queue.Queue = queue.Queue()  # str deltas; None = end-of-text
    chunks: asyncio.Queue = asyncio.Queue()  # PCM out; None = synthesis done

    def _produce():
        # Every streamer re-resolves its API key on each stream() call (tts_streaming ->
        # resolve_provider_secret), so the whole synthesis body runs under the requesting
        # profile's scope, not only the resolve step above (else the launch profile's key).
        with _config_profile_scope(profile):
            _synthesize()

    def _synthesize():
        nonlocal produced_audio, synthesis_failed
        from tools.tts_streaming import SentenceChunker
        from tools.tts_text_normalize import _strip_markdown_for_tts

        chunker = SentenceChunker.from_config(cfg)  # the requesting profile's tts.streaming.min_len

        # The session stays open for a whole agent turn and no text arrives
        # during tool execution, so without an idle flush a narration line with
        # no trailing whitespace ("Let me check.") sits in the chunker until
        # end-of-turn. Mirror the CLI speaker pipeline: poll with a timeout and
        # flush when the producer goes idle — immediately when the buffer ends
        # on sentence punctuation, after a longer quiet spell otherwise.
        idle_poll_seconds = 0.5
        idle_polls_before_force_flush = 4  # ~2s of silence

        def _sentences():
            idle_polls = 0
            while not stop.is_set():
                try:
                    delta = text_q.get(timeout=idle_poll_seconds)
                except queue.Empty:
                    idle_polls += 1
                    buffered = chunker.buf.strip()
                    if not buffered or ("<think" in chunker.buf and "</think>" not in chunker.buf):
                        continue
                    if buffered.endswith((".", "!", "?", "…", ":")) or idle_polls >= idle_polls_before_force_flush:
                        yield from chunker.flush()
                    continue
                idle_polls = 0
                if delta is None:
                    yield from chunker.flush()
                    return
                yield from chunker.feed(delta)

        try:
            for sentence in _sentences():
                cleaned = _strip_markdown_for_tts(sentence)
                if not cleaned:
                    continue
                for piece in _split_text_for_speak_stream(cleaned, cap):
                    for chunk in streamer.stream(piece):
                        if stop.is_set():
                            return
                        produced_audio = True
                        loop.call_soon_threadsafe(chunks.put_nowait, chunk)
        except Exception as exc:
            _log.warning("speak-stream synthesis failed: %s", exc)
            synthesis_failed = True
        finally:
            loop.call_soon_threadsafe(chunks.put_nowait, None)

    threading.Thread(target=_produce, daemon=True).start()

    async def _pump_client():
        # Text frames feed synthesis; done ends the text; stop/disconnect
        # (or any unparseable frame) is barge-in.
        try:
            while True:
                frame = json.loads(await ws.receive_text())
                if frame.get("text"):
                    text_q.put(str(frame["text"]))
                if frame.get("stop"):
                    break
                if frame.get("done"):
                    text_q.put(None)
        except Exception:
            pass
        stop.set()
        text_q.put(None)  # unblock the producer

    pump = asyncio.ensure_future(_pump_client())
    try:
        while True:
            chunk = await chunks.get()
            if chunk is None:
                break
            await _send_start()
            await ws.send_bytes(chunk)
        if not stop.is_set():
            # Fallback is the last resort: sentence synthesis was asked for and
            # produced nothing. A normal edge reply has already streamed PCM.
            if synthesis_failed and not produced_audio:
                await ws.send_json({"type": "fallback"})
            else:
                await _send_start()
                await ws.send_json({"type": "end"})
    except (WebSocketDisconnect, RuntimeError):
        pass
    finally:
        stop.set()
        text_q.put(None)
        pump.cancel()
        with contextlib.suppress(Exception):
            await ws.close()
