# CoreDeskAI Voice Bridge — Phase 2 (full conversational loop + barge-in)
# ========================================================================
# Telnyx Media Streaming (WS) <-> Voxtral Realtime STT <-> Mistral LLM <-> Supertonic TTS.
#   greet   : on stream start, the bridge speaks a Supertonic greeting (owns all out-audio)
#   inbound : caller audio -> μ/A-law decode -> 16k PCM -> Voxtral -> transcript
#   endpoint: ~ENDPOINT_SILENCE_MS of no new words = end of caller's turn
#   respond : turn text -> Mistral (streamed, per-sentence) -> Supertonic -> 8k G.711
#             -> send back over the bidirectional WS as media events (caller hears it)
#   barge-in: if Voxtral transcribes caller speech WHILE the bot is playing, cancel the
#             playback task immediately (+ best-effort Telnyx `clear`).
#
# Run: MISTRAL_API_KEY=... SUPERTONIC_URL=http://127.0.0.1:7788 python bridge.py
# Telnyx: streaming_start with stream_bidirectional_mode="rtp" (set in the webhook).
#
# Reply source (set via env): CONVEX_VOICE_TURN_URL = full app orchestration
# (/voice/turn: workflows/tools/KB/compliance + persistence) > CONVEX_LLM_URL =
# org LLM + naturalness (streamed) > direct Mistral. "The existing app orchestrates."

import asyncio
import audioop  # audioop-lts on Python 3.13+
import base64
import io
import json
import os
import re
import sys
import wave
from urllib.parse import parse_qs, urlparse

import httpx
import websockets
from mistralai.client import Mistral
from mistralai.client.models import AudioFormat

sys.stdout.reconfigure(encoding="utf-8", errors="replace")

PORT = int(os.environ.get("PORT", "8080"))
MODEL = os.environ.get("VOXTRAL_MODEL", "voxtral-mini-transcribe-realtime-2602")
TARGET_DELAY_MS = int(os.environ.get("VOXTRAL_DELAY_MS", "300"))  # 480→300: STT finalizes ~180ms sooner (slightly higher WER)
MISTRAL_API_KEY = os.environ.get("MISTRAL_API_KEY")
SUPERTONIC_URL = os.environ.get("SUPERTONIC_URL", "http://127.0.0.1:7788")
TTS_LANG = os.environ.get("TTS_LANG", "fr")
TTS_VOICE = os.environ.get("TTS_VOICE", "M1")
ENDPOINT_SILENCE_MS = int(os.environ.get("ENDPOINT_SILENCE_MS", "450"))  # 800→600→450: faster turn-end (watch for comma-pause splits)
# Echo guard: without acoustic echo cancellation, the caller's handset (especially on
# speakerphone) bleeds the bot's own TTS back on the inbound track. STT transcribes it
# and (pre-fix) fired a false "barge-in" that cut the bot off mid-sentence and polluted
# the next turn with echo fragments. Half-duplex gates the mic while the bot speaks so
# it can never hear itself. VOICE_HALF_DUPLEX=0 restores full-duplex barge-in (only
# safe with real AEC upstream). Trade-off: the caller can't interrupt mid-reply.
HALF_DUPLEX = os.environ.get("VOICE_HALF_DUPLEX", "1") != "0"
CHAT_MODEL = os.environ.get("MISTRAL_CHAT_MODEL", "mistral-small-latest")
# Filler / acknowledgment speech: the real reply can be many seconds away (deep
# /voice/turn pipeline + transatlantic Convex hop + TTS), so we play short
# pre-rendered waiting phrases instead of dead silence. The first plays
# ~FILLER_DELAY_MS after the caller finishes; if the real reply still hasn't
# started, the next phrase plays every FILLER_REPEAT_MS (the last phrase repeats)
# until it does — so even a 10–30s deep turn stays acknowledged. Phrases are
# pre-rendered once per call (~0ms to start) and "|"-separated. FILLER_DELAY_MS=0
# disables. Phrases are non-committal so they fit any pending reply.
FILLER_DELAY_MS = int(os.environ.get("FILLER_DELAY_MS", "400"))
FILLER_TEXTS = [s.strip() for s in os.environ.get(
    "VOICE_FILLER",
    "Alors,|Un instant, je vérifie ça.|Je regarde, un instant.",
).split("|") if s.strip()]
FILLER_REPEAT_MS = int(os.environ.get("FILLER_REPEAT_MS", "2500"))  # gap between waiting phrases
FILLER_MAX = int(os.environ.get("FILLER_MAX", "10"))               # safety cap on repeats
SYSTEM_PROMPT = os.environ.get(
    "VOICE_SYSTEM_PROMPT",
    "Tu es l'assistant vocal de CoreDeskAI. Réponds en français, de manière concise "
    "et naturelle, en une ou deux phrases courtes, comme au téléphone.",
)
SENTENCE_END = re.compile(r"[.!?…]+[\"»)\]]?\s")  # flush + speak each sentence as it streams in
GREETING = os.environ.get(
    "VOICE_GREETING",
    "Bonjour, ici l'assistant vocal de CoreDeskAI. Comment puis-je vous aider ?",
)  # bridge owns ALL outbound audio now (no Telnyx speak) — greet on stream start

# Convex orchestration: if CONVEX_LLM_URL is set, replies route through the existing
# app (org-selected LLM + §14.6 naturalness + tools/compliance) via the OpenAI-compatible
# /webhooks/vapi/llm SSE endpoint — "the existing app orchestrates." Else: direct Mistral.
# Two levels of "the app orchestrates":
#   CONVEX_VOICE_TURN_URL → FULL pipeline (intent, workflows, tools, KB/MCP, compliance,
#                           persistence) via /voice/turn. One-shot reply per turn.
#   CONVEX_LLM_URL        → LLM-reply layer only (org model + naturalness), streamed.
#   neither               → direct Mistral.
CONVEX_VOICE_TURN_URL = os.environ.get("CONVEX_VOICE_TURN_URL")
CONVEX_LLM_URL = os.environ.get("CONVEX_LLM_URL")
VOICE_ORG_ID = os.environ.get("VOICE_ORG_ID", "")          # which org this call belongs to
VAPI_WEBHOOK_SECRET = os.environ.get("VAPI_WEBHOOK_SECRET")  # x-vapi-secret, if the endpoint enforces it
MAX_HISTORY = int(os.environ.get("VOICE_MAX_HISTORY", "12"))  # messages of context kept per call

# Hybrid routing (#2): when BOTH the deep (/voice/turn) and light (CONVEX_LLM_URL)
# endpoints are set, divert only clearly-trivial conversational turns (greetings,
# acks, smalltalk) to the fast streaming light path; everything substantive stays
# on the deep pipeline (tools/KB/persistence + authoritative Convex context).
# Conservative on purpose — when in doubt, go deep — so tool/booking turns never
# lose context. VOICE_HYBRID=0 forces every turn deep (original behaviour).
HYBRID_ENABLED = os.environ.get("VOICE_HYBRID", "1") != "0"
# Trailing matches ONLY whitespace/punctuation — the turn must be *entirely* a
# trivial phrase, so "oui réserve jeudi à 15h" (a booking) still goes deep.
LIGHT_TURN_RE = re.compile(
    r"^(bonjour|bonsoir|salut|coucou|all[oô]|merci( beaucoup)?|"
    r"oui|non|ok|okay|d'accord|parfait|tr[èe]s bien|nickel|super|"
    r"[çc]a va|comment [çc]a va|au revoir|[àa] bient[ôo]t|bonne journ[ée]e)"
    r"[\s.,!?…]*$",
    re.IGNORECASE,
)


def turn_uses_light_path(turn: str) -> bool:
    """True for short social turns that don't need tools/KB → route to the fast
    streaming light path. Requires the light endpoint to be configured."""
    return bool(CONVEX_LLM_URL) and HYBRID_ENABLED and bool(LIGHT_TURN_RE.match(turn.strip()))
TAG_RE = re.compile(r"<[^>]+>")  # strip SSML/markup before Supertonic (it would read tags aloud)
# Arabic script (main block + supplement + presentation forms). Supertonic-3 speaks
# Arabic with any voice, but only when told lang="ar" — otherwise an Arabic reply is
# phonemized as the default language and comes out garbled. We pick the TTS language
# per sentence from its script so a multilingual reply is voiced correctly.
# Built from code-point ranges (ASCII source, no literal RTL chars): Arabic +
# Arabic Supplement + Extended-A + Presentation Forms-A/B.
_AR_RANGES = ((0x0600, 0x06FF), (0x0750, 0x077F), (0x08A0, 0x08FF),
              (0xFB50, 0xFDFF), (0xFE70, 0xFEFC))
AR_SCRIPT = re.compile("[" + "".join(f"{chr(a)}-{chr(b)}" for a, b in _AR_RANGES) + "]")


def tts_lang_for(text: str) -> str:
    """Pick the Supertonic language for a sentence: 'ar' when it contains Arabic
    script, else the configured TTS_LANG default (fr/en/…)."""
    return "ar" if AR_SCRIPT.search(text) else TTS_LANG

client = Mistral(api_key=MISTRAL_API_KEY)
AUDIO_FORMAT = AudioFormat(encoding="pcm_s16le", sample_rate=16000)


async def voice_turn_deltas(turn: str, call_id: str, org: str, caller: str | None = None,
                            campaign: str | None = None):
    """Full app orchestration via /voice/turn (persisted; tools/workflows/KB). One-shot reply.
    Yielded as a single chunk — respond()'s sentence splitter still does per-sentence TTS.
    `campaign` (outbound campaign id) makes /voice/turn run the qualification script."""
    headers = {"Content-Type": "application/json", "x-organization-id": org}
    if VAPI_WEBHOOK_SECRET:
        headers["x-vapi-secret"] = VAPI_WEBHOOK_SECRET
    body = {"callId": call_id, "text": turn}
    if caller:
        body["callerNumber"] = caller
    if campaign:
        body["campaignId"] = campaign
    async with httpx.AsyncClient(timeout=90.0) as h:
        r = await h.post(CONVEX_VOICE_TURN_URL, headers=headers, json=body)
        r.raise_for_status()
        reply = (r.json().get("reply") or "").strip()
    if reply:
        yield reply


async def convex_sse_deltas(messages: list[dict], org: str):
    """Stream reply text deltas from the Convex app (org LLM + naturalness)."""
    headers = {"Content-Type": "application/json", "x-organization-id": org}
    if VAPI_WEBHOOK_SECRET:
        headers["x-vapi-secret"] = VAPI_WEBHOOK_SECRET
    body = {"model": "voice", "messages": messages, "stream": True,
            "temperature": 0.5, "max_tokens": 160}
    async with httpx.AsyncClient(timeout=60.0) as h:
        async with h.stream("POST", CONVEX_LLM_URL, headers=headers, json=body) as r:
            r.raise_for_status()
            async for line in r.aiter_lines():
                if not line.startswith("data:"):
                    continue
                payload = line[5:].strip()
                if payload == "[DONE]":
                    return
                try:
                    delta = json.loads(payload)["choices"][0]["delta"].get("content") or ""
                except Exception:
                    delta = ""
                if delta:
                    yield delta


async def mistral_deltas(messages: list[dict]):
    """Direct Mistral streaming (fallback when CONVEX_LLM_URL is unset)."""
    async with await client.chat.stream_async(
        model=CHAT_MODEL, messages=messages, temperature=0.5, max_tokens=160,
    ) as stream:
        async for event in stream:
            try:
                delta = event.data.choices[0].delta.content or ""
            except Exception:
                delta = ""
            if delta:
                yield delta


# ── G.711 codec helpers (Telnyx negotiates PCMU/μ-law or PCMA/A-law) ──
def g711_decode(raw: bytes, codec: str) -> bytes:
    if codec == "PCMA":
        return audioop.alaw2lin(raw, 2)
    if codec == "PCMU":
        return audioop.ulaw2lin(raw, 2)
    return raw  # L16/linear


def g711_encode(pcm: bytes, codec: str) -> bytes:
    if codec == "PCMA":
        return audioop.lin2alaw(pcm, 2)
    if codec == "PCMU":
        return audioop.lin2ulaw(pcm, 2)
    return pcm


async def synthesize(text: str, codec: str, lang: str | None = None) -> bytes:
    """Supertonic TTS -> 8kHz G.711 bytes in the call's codec.
    `lang` overrides the TTS_LANG default (used to voice Arabic replies as 'ar')."""
    async with httpx.AsyncClient(timeout=30.0) as h:
        r = await h.post(
            f"{SUPERTONIC_URL}/v1/tts",
            json={"text": text, "voice": TTS_VOICE, "lang": lang or TTS_LANG, "response_format": "wav"},
        )
        r.raise_for_status()
        content = r.content
    wf = wave.open(io.BytesIO(content), "rb")
    sr, ch = wf.getframerate(), wf.getnchannels()
    pcm = wf.readframes(wf.getnframes())
    wf.close()
    if ch == 2:
        pcm = audioop.tomono(pcm, 2, 0.5, 0.5)
    pcm8, _ = audioop.ratecv(pcm, 2, 1, sr, 8000, None)  # -> 8kHz
    return g711_encode(pcm8, codec)


async def handle_telnyx(telnyx_ws):
    # Per-call routing context from the WS URL query (the Telnyx webhook appends
    # ?org=<org>&from=<caller>&call=<ccid> to stream_url so the bridge is multi-org).
    try:
        _q = parse_qs(urlparse(getattr(getattr(telnyx_ws, "request", None), "path", "") or "").query)
    except Exception:
        _q = {}
    org_id = [(_q.get("org", [VOICE_ORG_ID])[0] or VOICE_ORG_ID)]
    caller_num = [_q.get("from", [None])[0]]
    call_q = _q.get("call", [None])[0]
    campaign_id = _q.get("campaign", [None])[0]   # set for outbound campaign calls → run the script
    print(f"[telnyx] media stream connected (org={org_id[0] or '-'})", flush=True)
    loop = asyncio.get_event_loop()
    audio_q: asyncio.Queue = asyncio.Queue()
    pending: list[str] = []     # transcript deltas since the last bot reply
    history: list[dict] = []    # conversation context for multi-turn (user/assistant)
    last_delta = [loop.time()]  # ts of most recent delta (mutable for closures)
    speaking = [False]          # True while the bot is playing audio back
    playback_task = [None]      # the in-flight greet()/respond() task (cancel = barge-in)
    codec = ["PCMU"]            # set from `start`; mutable for closures
    src_rate = [8000]
    call_id = ["?"]             # Telnyx call_control_id — conversation key for /voice/turn
    out_lock = asyncio.Lock()   # serialize outbound audio (filler vs real reply never overlap)
    filler_audio = [None]       # list of pre-rendered G.711 filler clips (synthesized once per call)

    async def barge_in():
        # Caller talked over the bot → cut the bot's audio immediately.
        t = playback_task[0]
        if t and not t.done():
            t.cancel()
        try:  # best-effort: ask Telnyx to drop anything already buffered
            await telnyx_ws.send(json.dumps({"event": "clear"}))
        except Exception:
            pass
        speaking[0] = False
        print("[barge-in] caller spoke over the bot — stopped playback", flush=True)

    async def audio_iter():
        while True:
            chunk = await audio_q.get()
            if chunk is None:
                return
            yield chunk

    async def run_stt():
        try:
            async for ev in client.audio.realtime.transcribe_stream(
                audio_stream=audio_iter(), model=MODEL,
                audio_format=AUDIO_FORMAT, target_streaming_delay_ms=TARGET_DELAY_MS,
            ):
                t = getattr(ev, "type", None)
                if t == "transcription.text.delta":
                    txt = getattr(ev, "text", None) or getattr(ev, "delta", "") or ""
                    # Echo guard: while the bot is speaking, half-duplex drops everything
                    # STT hears — without AEC it's almost certainly the bot's own audio
                    # echoing back, which would otherwise fire a false barge-in and feed
                    # echo fragments into the next turn.
                    if speaking[0] and HALF_DUPLEX:
                        continue
                    pending.append(txt)
                    last_delta[0] = loop.time()
                    print(f"[stt+] {txt}", flush=True)
                    if speaking[0]:  # full-duplex only: caller talked over the bot → barge-in
                        await barge_in()
                elif t == "error":
                    print(f"[stt ERROR] {ev}", flush=True)
        except Exception as e:
            print(f"[stt] stream error: {e}", flush=True)

    async def send_audio_back(g711: bytes):
        # 20ms frames @ 8kHz G.711 = 160 bytes; pace at real-time.
        for i in range(0, len(g711), 160):
            await telnyx_ws.send(json.dumps(
                {"event": "media", "media": {"payload": base64.b64encode(g711[i:i + 160]).decode()}}
            ))
            await asyncio.sleep(0.02)

    async def respond(turn: str, speech_end: float):
        # 2d: stream the LLM; synthesize + play each sentence as soon as it's ready,
        # so time-to-first-audio is ~one short sentence instead of the whole reply.
        # `speech_end` = ts of the caller's last STT word; we log the wall-clock gap
        # from there to the FIRST reply-audio frame (the latency the caller perceives).
        speaking[0] = True
        full: list[str] = []
        buf = ""
        first_audio = [False]    # any audio reached the caller (filler OR real reply)
        first_content = [False]  # the first REAL reply sentence reached the caller

        # Pick the orchestration source (most → least "the app orchestrates").
        # Hybrid (#2): trivial social turns take the fast streaming light path even
        # when the deep endpoint is set; everything substantive stays deep. Campaign
        # calls ALWAYS stay deep — the qualification script only runs on /voice/turn,
        # so a trivial greeting must not be diverted to the light path.
        divert_light = turn_uses_light_path(turn) and not campaign_id
        if CONVEX_VOICE_TURN_URL and not divert_light:
            reply_src = "voice-turn"  # full pipeline + persistence (Convex owns history)
            deltas = voice_turn_deltas(turn, call_id[0], org_id[0], caller_num[0], campaign_id)
        elif CONVEX_LLM_URL:
            reply_src = "convex"      # org LLM + naturalness, streamed
            deltas = convex_sse_deltas([{"role": "system", "content": SYSTEM_PROMPT}, *history,
                                        {"role": "user", "content": turn}], org_id[0])
        else:
            reply_src = "mistral"     # direct Mistral fallback
            deltas = mistral_deltas([{"role": "system", "content": SYSTEM_PROMPT}, *history,
                                     {"role": "user", "content": turn}])

        async def play(seg: str):
            seg = TAG_RE.sub("", seg).strip()  # drop any SSML/markup before TTS
            if not seg:
                return
            g711 = await synthesize(seg, codec[0], tts_lang_for(seg))  # TTS off-lock
            async with out_lock:  # serialize vs the filler so audio never overlaps
                if not first_audio[0]:
                    first_audio[0] = True
                    print(f"[lat] end-of-speech -> first audio (reply) = {loop.time() - speech_end:.2f}s "
                          f"(endpoint {ENDPOINT_SILENCE_MS}ms)", flush=True)
                if not first_content[0]:
                    first_content[0] = True
                    print(f"[lat-content] end-of-speech -> first reply audio = {loop.time() - speech_end:.2f}s "
                          f"({reply_src} + TTS)", flush=True)
                await send_audio_back(g711)

        async def filler_watchdog():
            # If the real reply is slow (deep pipeline + transatlantic hop + TTS), play
            # short pre-rendered waiting phrases so the caller hears acknowledgment
            # instead of dead silence. First phrase plays after ~FILLER_DELAY_MS; if the
            # real reply still hasn't started, the next phrase plays every
            # FILLER_REPEAT_MS (last phrase repeats) up to FILLER_MAX times. Stops the
            # instant the real reply's first audio lands. Skipped if disabled or the
            # clips didn't pre-render in time.
            clips = filler_audio[0]
            if FILLER_DELAY_MS <= 0 or not clips:
                return
            await asyncio.sleep(FILLER_DELAY_MS / 1000)
            for i in range(FILLER_MAX):
                if first_content[0]:
                    return
                async with out_lock:
                    if first_content[0]:  # real reply won the race while we waited on the lock
                        return
                    if not first_audio[0]:
                        first_audio[0] = True
                        print(f"[lat] end-of-speech -> first audio (filler) = {loop.time() - speech_end:.2f}s "
                              f"(endpoint {ENDPOINT_SILENCE_MS}ms)", flush=True)
                    await send_audio_back(clips[min(i, len(clips) - 1)])
                # Wait before the next reassurance, bailing fast once the reply arrives.
                for _ in range(max(1, FILLER_REPEAT_MS // 100)):
                    if first_content[0]:
                        return
                    await asyncio.sleep(0.1)

        wd_task = asyncio.create_task(filler_watchdog())
        try:
            async for delta in deltas:
                full.append(delta)
                buf += delta
                while True:
                    m = SENTENCE_END.search(buf)
                    if not m:
                        break
                    seg = buf[:m.end()].strip()
                    buf = buf[m.end():]
                    if seg:
                        await play(seg)
            if buf.strip():  # final partial sentence
                await play(buf.strip())
            reply = "".join(full).strip()
            print(f"[bot] [{reply_src}] {reply}", flush=True)
            # Keep a local history buffer for the light/Mistral paths. Append on EVERY
            # turn (incl. deep) so a later trivial turn routed to the light path still
            # has recent context — Convex keeps its own authoritative history for the
            # deep /voice/turn path (which looks it up by call_id, not from here).
            history.append({"role": "user", "content": turn})
            history.append({"role": "assistant", "content": reply})
            del history[:-MAX_HISTORY]  # keep only the most recent MAX_HISTORY messages
        except Exception as e:
            print(f"[bot] error: {e}", flush=True)
        finally:
            wd_task.cancel()  # stop the filler watchdog (also on barge-in cancel)
            speaking[0] = False
            last_delta[0] = loop.time()  # fresh silence baseline; echo tail can't fire a turn

    async def turn_monitor():
        # End-of-turn = pending text + a gap of silence + bot not already talking.
        while True:
            await asyncio.sleep(0.2)
            if pending and not speaking[0] and (loop.time() - last_delta[0]) * 1000 >= ENDPOINT_SILENCE_MS:
                turn = "".join(pending).strip()
                speech_end = last_delta[0]  # when the caller's last word landed
                pending.clear()
                if turn:
                    speaking[0] = True  # claim playback now so this turn can't double-fire
                    playback_task[0] = asyncio.create_task(respond(turn, speech_end))

    stt_task = asyncio.create_task(run_stt())
    mon_task = asyncio.create_task(turn_monitor())
    rate_state = None
    try:
        async for raw in telnyx_ws:
            try:
                evt = json.loads(raw)
            except Exception:
                continue
            kind = evt.get("event")
            if kind == "start":
                start_obj = evt.get("start", {})
                fmt = start_obj.get("media_format", {})
                codec[0] = (fmt.get("encoding") or "PCMU").upper()
                src_rate[0] = int(fmt.get("sample_rate") or 8000)
                call_id[0] = (call_q or start_obj.get("call_control_id") or evt.get("stream_id")
                              or f"telnyx-{id(telnyx_ws)}")
                print(f"[telnyx] start; codec={codec[0]} rate={src_rate[0]} call={call_id[0]}", flush=True)

                async def greet():
                    # Greet immediately so the caller isn't met with silence and knows
                    # who answered (replaces the removed Telnyx `speak`). Same Supertonic
                    # voice as replies. `speaking` guards the endpoint monitor.
                    speaking[0] = True
                    try:
                        await send_audio_back(await synthesize(GREETING, codec[0], tts_lang_for(GREETING)))
                        print(f"[bot] (greeting) {GREETING}", flush=True)
                    except Exception as e:
                        print(f"[greet] error: {e}", flush=True)
                    finally:
                        last_delta[0] = loop.time()  # don't let a stale gap fire a turn
                        speaking[0] = False
                playback_task[0] = asyncio.create_task(greet())

                async def prep_filler():
                    # Render the filler phrases once per call (codec is known now) so the
                    # watchdog can play them with ~0ms synth latency mid-turn.
                    if FILLER_DELAY_MS <= 0 or not FILLER_TEXTS:
                        return
                    try:
                        filler_audio[0] = [
                            await synthesize(phrase, codec[0], tts_lang_for(phrase))
                            for phrase in FILLER_TEXTS
                        ]
                    except Exception as e:
                        print(f"[filler] prep error: {e}", flush=True)
                asyncio.create_task(prep_filler())
            elif kind == "media":
                media = evt.get("media", {})
                if media.get("track") and media["track"] != "inbound":
                    continue  # only the caller's audio feeds STT
                if HALF_DUPLEX and speaking[0]:
                    continue  # echo guard: don't feed the bot's own audio back into STT
                payload = media.get("payload")
                if not payload:
                    continue
                pcm8 = g711_decode(base64.b64decode(payload), codec[0])
                pcm16, rate_state = audioop.ratecv(pcm8, 2, 1, src_rate[0], 16000, rate_state)
                await audio_q.put(pcm16)
            elif kind == "stop":
                print("[telnyx] stop", flush=True)
                break
    finally:
        await audio_q.put(None)
        for task in (stt_task, mon_task, playback_task[0]):
            if task:
                task.cancel()
        print("[telnyx] closed", flush=True)


async def main():
    if not MISTRAL_API_KEY:
        print("WARNING: MISTRAL_API_KEY not set — Voxtral will fail.", flush=True)
    async with websockets.serve(handle_telnyx, "0.0.0.0", PORT):
        print(f"[bridge] listening on :{PORT}  (STT in + Supertonic TTS back)", flush=True)
        await asyncio.Future()


if __name__ == "__main__":
    asyncio.run(main())
