"""CoreDeskAI LiveKit voice agent worker.

Replaces services/voice-bridge/bridge.py per LIVEKIT_VOICE_MIGRATION.md:
LiveKit Agents supplies VAD / turn detection / barge-in / noise cancellation;
a realtime speech-to-speech model is the conversation brain; Convex stays the
tool + compliance + persistence layer (see convex_tools.py).

Engine selection (env VOICE_REALTIME_PROVIDER):
  deepgram — CHAINED pipeline (NOT speech-to-speech): Deepgram Nova-3 STT +
            an OpenAI-compatible LLM (Groq by default) + Deepgram Aura-2 TTS,
            with Silero VAD. The path that runs WITHOUT the Azure gpt-realtime
            quota, and EU-resident via Deepgram's EU endpoint. ⚠️ Aura-2 has no
            Arabic voice — see build_chained_session().
  azure   — Azure OpenAI `gpt-realtime`, Sweden Central (EU-resident; the locked
            production speech-to-speech choice). Needs AZURE_OPENAI_* env.
  openai  — plain OpenAI Realtime (US egress — dev/testing fallback ONLY; using
            it for real calls reverses the EU-sovereignty rationale).
  gemini  — Google Gemini Live via AI Studio (US/global egress — dev/testing
            fallback ONLY, same EU caveat as openai). Needs GOOGLE_API_KEY and
            the livekit-plugins-google plugin installed.

Call routing context (org / call / campaign) comes from job metadata when
dispatched, with env-var fallbacks for console/dev testing.

Run (dev):    python agent.py console     (talk to it from the terminal)
              python agent.py dev         (register with LiveKit Cloud, get dispatched)
"""

from __future__ import annotations

import asyncio
import json
import logging
import os
import re
import time

from dotenv import load_dotenv

from livekit import agents
from livekit.agents import Agent, AgentSession, RunContext, function_tool
from livekit.agents import metrics as lk_metrics
from livekit.plugins import openai as lk_openai
from livekit.plugins import google as lk_google  # registered on main thread (required by LiveKit)
from livekit.plugins import deepgram as lk_deepgram
from livekit.plugins import silero as lk_silero
from openai.types.beta.realtime.session import TurnDetection

load_dotenv()

logger = logging.getLogger("voice-livekit")
logging.basicConfig(level=logging.INFO)

from convex_tools import ConvexVoiceClient, VoiceSession  # noqa: E402
from recording import start_room_recording, stop_room_recording  # noqa: E402


def _openai_realtime_turn_detection(turn_taking: dict | None = None) -> TurnDetection:
    """VAD / barge-in — session turnTaking from Agent Builder, else env defaults.

    interrupt_response=True cancels in-flight assistant audio on user speech.
    We do NOT reset AgentSession or Convex conversation on barge-in (context retained).

    Must return a real openai.types.beta.realtime.session.TurnDetection instance,
    not a plain dict — livekit-plugins-openai's to_turn_detection() only converts
    inputs that pass `isinstance(x, TurnDetection)`; anything else (e.g. a dict) is
    forwarded unchanged and later crashes with AttributeError on `.create_response`.
    """
    tt = turn_taking or {}
    interrupt = tt.get("interruptResponse", True)
    if interrupt is None:
        interrupt = True
    vad_type = (os.environ.get("OPENAI_REALTIME_VAD_TYPE") or "server_vad").lower()
    if vad_type == "semantic_vad":
        return TurnDetection(
            type="semantic_vad",
            eagerness=os.environ.get("OPENAI_REALTIME_VAD_EAGERNESS", "low"),
            create_response=True,
            interrupt_response=bool(interrupt),
        )
    silence_ms = int(
        tt.get("silenceDurationMs")
        or os.environ.get("OPENAI_REALTIME_SILENCE_MS", "1200")
    )
    threshold = float(
        tt.get("vadThreshold")
        or os.environ.get("OPENAI_REALTIME_VAD_THRESHOLD", "0.6")
    )
    return TurnDetection(
        type="server_vad",
        threshold=threshold,
        prefix_padding_ms=300,
        silence_duration_ms=silence_ms,
        create_response=True,
        interrupt_response=bool(interrupt),
    )


def _openai_realtime_transcription(transcription: dict | None = None) -> dict:
    """Whisper language hint — prefer /voice/session transcription, else env."""
    if transcription and isinstance(transcription, dict):
        model = (transcription.get("model") or "whisper-1").strip() or "whisper-1"
        lang = (transcription.get("language") or "").strip()
        if lang:
            return {"model": model, "language": lang}
        return {"model": model}
    # No env override and no Convex session (e.g. console/dev mode) — omit
    # `language` entirely so Whisper auto-detects rather than assuming one.
    lang = (os.environ.get("OPENAI_REALTIME_TRANSCRIPTION_LANGUAGE") or "").strip()
    if not lang:
        return {"model": "whisper-1"}
    return {"model": "whisper-1", "language": lang}


def build_openai_realtime_model(
    voice_id: str | None = None,
    turn_taking: dict | None = None,
    transcription: dict | None = None,
):
    """OpenAI Realtime S2S — same brain as the `/realtime-voice` browser demo.

    Persona/dialect/scripts still come from Convex `/voice/session` (Agent Builder).
    Voice: session voiceId → OPENAI_REALTIME_VOICE → marin (Story 2).
    Phase 3: turn_taking + transcription language from session (not env-only).
    """
    api_key = os.environ.get("OPENAI_API_KEY")
    if not api_key:
        raise ValueError("OPENAI_API_KEY is required when VOICE_REALTIME_PROVIDER=openai")

    model = os.environ.get("OPENAI_REALTIME_MODEL", "gpt-realtime")
    voice = (voice_id or "").strip() or os.environ.get("OPENAI_REALTIME_VOICE", "marin") or "marin"
    kwargs: dict = {
        "api_key": api_key,
        "model": model,
        "voice": voice,
        "turn_detection": _openai_realtime_turn_detection(turn_taking),
        "input_audio_transcription": _openai_realtime_transcription(transcription),
    }
    logger.info(
        "VOICE_REALTIME_PROVIDER=openai — model=%s voice=%s silence_ms=%s interrupt=%s tx=%s",
        model,
        voice,
        kwargs["turn_detection"].silence_duration_ms,
        kwargs["turn_detection"].interrupt_response,
        (transcription or {}).get("language")
        or os.environ.get("OPENAI_REALTIME_TRANSCRIPTION_LANGUAGE", "ar"),
    )
    # Older livekit-plugins-openai may not accept all kwargs — peel until construct works.
    while True:
        try:
            return lk_openai.realtime.RealtimeModel(**kwargs)
        except TypeError as err:
            msg = str(err)
            dropped = False
            for key in ("input_audio_transcription", "turn_detection", "voice", "model"):
                if key in kwargs and key in msg:
                    kwargs.pop(key, None)
                    dropped = True
                    logger.warning("RealtimeModel: dropping unsupported kwarg %s (%s)", key, err)
                    break
            if not dropped:
                if set(kwargs.keys()) != {"api_key"}:
                    kwargs = {"api_key": api_key}
                    continue
                raise


def build_realtime_model(
    voice_id: str | None = None,
    turn_taking: dict | None = None,
    transcription: dict | None = None,
):
    """Build the provider-specific LiveKit realtime model (return type varies by provider)."""
    provider = os.environ.get("VOICE_REALTIME_PROVIDER", "openai").lower()
    if provider == "azure":
        # Passing api_version uses the stable legacy realtime path
        # (/openai/realtime?api-version=…&deployment=…) — LiveKit's documented Azure
        # config. Omitting it targets the newer /openai/v1/realtime path, which 400s
        # on the Sweden Central gpt-realtime resource. Tunable via env so the version
        # can be changed without a code edit (e.g. 2024-10-01-preview if this fails).
        return lk_openai.realtime.RealtimeModel.with_azure(
            azure_deployment=os.environ.get("AZURE_OPENAI_DEPLOYMENT", "gpt-realtime"),
            azure_endpoint=os.environ["AZURE_OPENAI_ENDPOINT"],
            api_key=os.environ["AZURE_OPENAI_API_KEY"],
            api_version=os.environ.get("AZURE_OPENAI_API_VERSION", "2025-04-01-preview"),
        )
    if provider == "openai":
        return build_openai_realtime_model(
            voice_id=voice_id,
            turn_taking=turn_taking,
            transcription=transcription,
        )
    if provider == "gemini":
        logger.warning(
            "VOICE_REALTIME_PROVIDER=gemini — US/global egress; dev/testing only, NOT the EU production path"
        )
        # NOTE: in some plugin versions the Gemini Live realtime class is at
        # lk_google.realtime.RealtimeModel instead of lk_google.beta.realtime.RealtimeModel —
        # adjust if the call below fails against the installed version.
        return lk_google.beta.realtime.RealtimeModel(
            # Model + voice are env-tunable (Gemini Live model names change often).
            # Default = the plugin's current known-good AI Studio Live model.
            model=os.environ.get(
                "GEMINI_REALTIME_MODEL", "gemini-2.5-flash-native-audio-preview-12-2025"
            ),
            voice=os.environ.get("GEMINI_REALTIME_VOICE", "Puck"),
            api_key=os.environ["GOOGLE_API_KEY"],
        )
    raise ValueError(f"Unknown VOICE_REALTIME_PROVIDER: {provider}")


def build_chained_llm():
    """LLM for the chained pipeline.

    Reuses the product's existing Groq stack (fast, already wired) through the
    OpenAI-compatible plugin; OpenAI selectable as a fallback. This is the
    'brain' in the chained design — the realtime model is NOT used here.
    """
    llm_provider = os.environ.get("VOICE_CHAINED_LLM", "groq").lower()
    if llm_provider == "groq":
        # Groq via its OpenAI-compatible endpoint. (Equivalent to
        # openai.LLM.with_groq() on newer plugin versions; 1.5.17 lacks that
        # helper, so we point base_url at Groq directly.)
        return lk_openai.LLM(
            model=os.environ.get("VOICE_CHAINED_LLM_MODEL", "llama-3.3-70b-versatile"),
            api_key=os.environ["GROQ_API_KEY"],
            base_url=os.environ.get("GROQ_BASE_URL", "https://api.groq.com/openai/v1"),
        )
    if llm_provider == "openai":
        return lk_openai.LLM(
            model=os.environ.get("VOICE_CHAINED_LLM_MODEL", "gpt-4o-mini"),
            api_key=os.environ["OPENAI_API_KEY"],
        )
    raise ValueError(f"Unknown VOICE_CHAINED_LLM: {llm_provider}")


def build_chained_session() -> AgentSession:
    """Chained STT → LLM → TTS pipeline (Deepgram speech + Groq brain).

    The path we can run WITHOUT the Azure gpt-realtime quota. Deepgram Nova-3 STT
    on the EU endpoint (audio stays in the EU for GDPR) + an OpenAI-compatible LLM
    (Groq by default) + Deepgram Aura-2 TTS + Silero VAD for turn detection.

    Language notes:
      • STT language defaults to `multi` so FR/EN/AR code-switching (Tunisian
        speech) is transcribed on one connection.
      • ⚠️ Aura-2 TTS has NO Arabic voice (English + French + a few EU langs only).
        DEEPGRAM_TTS_MODEL must be a valid Aura-2 voice id (default is English —
        set a French voice for the French-default product). For Arabic *output*
        a different TTS (ElevenLabs / Azure Neural) is required — wire that in
        when Arabic voice replies are needed.
    """
    # DEEPGRAM_BASE_URL is the HOST (no path). The STT/TTS plugins bake the full
    # path into their default base_url (/v1/listen, /v1/speak), so a bare host
    # must have the right path appended per service — otherwise the STT websocket
    # connects to `wss://host/?…` and 404s. Accept either a host or a full
    # listen/speak URL and normalize back to the host.
    dg_host = re.sub(
        r"/v1/(listen|speak)/?$",
        "",
        os.environ.get("DEEPGRAM_BASE_URL", "https://api.eu.deepgram.com").rstrip("/"),
    )
    dg_key = os.environ["DEEPGRAM_API_KEY"]

    stt = lk_deepgram.STT(
        model=os.environ.get("DEEPGRAM_STT_MODEL", "nova-3"),
        language=os.environ.get("DEEPGRAM_STT_LANGUAGE", "multi"),
        api_key=dg_key,
        base_url=f"{dg_host}/v1/listen",
    )
    tts = lk_deepgram.TTS(
        model=os.environ.get("DEEPGRAM_TTS_MODEL", "aura-2-thalia-en"),
        api_key=dg_key,
        base_url=f"{dg_host}/v1/speak",
    )
    return AgentSession(
        stt=stt,
        llm=build_chained_llm(),
        tts=tts,
        vad=lk_silero.VAD.load(),
    )


def build_agent_session(
    voice_id: str | None = None,
    turn_taking: dict | None = None,
    transcription: dict | None = None,
) -> AgentSession:
    """Build the AgentSession for the configured engine.

    openai (default for demo) → OpenAI Realtime S2S (same as /realtime-voice lab).
    deepgram/chained → Deepgram STT + LLM + Deepgram TTS (no Azure dependency).
    azure/gemini → other realtime speech-to-speech models.
    voice_id: Story 2 — per-agent OpenAI Realtime voice from /voice/session.
    turn_taking / transcription: Phase 3 — from /voice/session Agent Builder settings.
    """
    provider = os.environ.get("VOICE_REALTIME_PROVIDER", "openai").lower()
    if provider in ("deepgram", "chained"):
        return build_chained_session()
    return AgentSession(
        llm=build_realtime_model(
            voice_id=voice_id,
            turn_taking=turn_taking,
            transcription=transcription,
        )
    )


def _opening_reply_instructions(cfg: VoiceSession) -> str:
    """First-turn instructions.

    The dialect lock always comes first and applies from the very first word —
    an opening script only supplies the wording, never the language. Previously
    a non-empty `scripts.opening` (which falls back to `greetMessage`, often
    left at an English default) short-circuited straight to "speak this
    verbatim" with no language enforcement, so a Tunisian Derja agent with an
    empty *opening script* field but a leftover English `greetMessage` would
    open in English and only switch once the caller spoke Arabic.
    """
    opening = ""
    try:
        opening = (cfg.scripts or {}).get("opening") or ""
    except Exception:
        opening = ""

    dialect = (cfg.dialect or "").lower()
    lang = (cfg.language or "").lower()

    if cfg.campaign:
        if dialect == "tunisian_derja" or lang.startswith("ar"):
            forced = "ابدأ فوراً بالدارجة التونسية من أول كلمة. متحكيش إنقليزي ولا فرنسي. "
            if opening:
                return forced + f'قول السكريبت التالي تقريباً بالحرف: "{opening}"'
            return forced + "سلّم باختصار، قدّم سبب المكالمة، وابدأ السكريبت."
        forced = "Démarre immédiatement en français, dès le premier mot. "
        if opening:
            return forced + f'Dis ce script d\'ouverture (quasi mot pour mot) : "{opening}"'
        return forced + "Salue brièvement, présente la raison de l'appel, puis commence le script."

    if dialect == "tunisian_derja" or (lang.startswith("ar") and dialect not in ("msa", "gulf")):
        forced = (
            "ابدأ فوراً بالدارجة التونسية من أول كلمة. متحكيش إنقليزي ولا فرنسي "
            "حتى لو السكريبت مكتوب بلغة أخرى. "
        )
        if opening:
            return forced + f'قول السكريبت التالي (ترجمو للدارجة التونسية إذا لزم): "{opening}"'
        return forced + "رحّب بالمتصل بالدارجة التونسية باختصار واسأل شنية يحب."

    if dialect == "gulf":
        forced = "ابدأ فوراً باللهجة الخليجية من أول كلمة. "
        if opening:
            return forced + f'قول السكريبت التالي: "{opening}"'
        return forced + "رحّب بالمتصل باللهجة الخليجية باختصار واسأله كيف تقدر تساعده."

    if dialect == "msa" or lang.startswith("ar"):
        forced = "ابدأ فوراً بالعربية الفصحى من أول كلمة. "
        if opening:
            return forced + f'قل السكريبت التالي: "{opening}"'
        return forced + "رحّب بالمتصل باختصار بالعربية الفصحى واسأله كيف يمكن مساعدته."

    if dialect in ("parisian", "canadian_fr") or lang.startswith("fr"):
        forced = "Démarre immédiatement en français, dès le premier mot. "
        if opening:
            return forced + f'Dis ce script d\'ouverture : "{opening}"'
        return forced + "Salue brièvement l'appelant en français et demande comment tu peux aider."

    if dialect.startswith("en") or lang.startswith("en"):
        forced = "Start immediately in English, from the very first word. "
        if opening:
            return forced + f'Say this opening script: "{opening}"'
        return forced + "Greet the caller briefly in English and ask how you can help."

    # No configured language/dialect — open with a short, neutral greeting and
    # let the caller's reply set the language; continue in whatever language
    # they use from then on.
    if opening:
        return f'Speak this opening first (near-verbatim), then continue: "{opening}"'
    return (
        "Greet the caller briefly with a short, language-neutral greeting "
        "(e.g. \"Hello / Bonjour\") and ask how you can help. Then continue "
        "the whole conversation in whichever language the caller replies in."
    )


class CoreDeskVoiceAgent(Agent):
    """Persona + tools come from the Convex /voice/session bootstrap."""

    def __init__(
        self,
        session_cfg: VoiceSession,
        convex: ConvexVoiceClient,
        call_id: str,
        room_name: str,
        on_transferred=None,
    ):
        super().__init__(instructions=session_cfg.system_prompt)
        self._cfg = session_cfg
        self._convex = convex
        self._call_id = call_id
        self._room_name = room_name
        self._on_transferred = on_transferred

    # ── Convex-backed function tools (contract §B) ──────────────────────────
    # The realtime model selects + fills arguments; endpoints return structured
    # JSON the model verbalizes. Tool latency rule: the persona prompt tells the
    # model to speak a brief acknowledgment before calling.

    @function_tool()
    async def search_knowledge(self, context: RunContext, query: str) -> dict:
        """Search the organization's knowledge base for information that answers
        the caller's question (policies, hours, products, procedures). Use for
        any informational question about the business. Returns text snippets —
        answer ONLY from them; if nothing is found, say you don't have that
        information."""
        return await self._convex.search_knowledge(query)

    @function_tool()
    async def lookup_data(self, context: RunContext, tool_name: str, args_json: str) -> dict:
        """Execute one of the organization's data tools (listed in your
        instructions) to read or write live business data (orders, calendar,
        sheets). `tool_name` must exactly match a listed tool name; `args_json`
        is a JSON object string with that tool's parameters."""
        try:
            tool_args = json.loads(args_json) if args_json else {}
        except json.JSONDecodeError:
            return {"success": False, "error": {"code": "bad_args", "message": "args_json is not valid JSON"}}
        return await self._convex.lookup_data(
            tool_name, tool_args, self._cfg.conversation_id or ""
        )

    @function_tool()
    async def run_workflow(self, context: RunContext, workflow_name: str, inputs_json: str) -> dict:
        """Start one of the organization's predefined workflows (listed in your
        instructions). `inputs_json` is a JSON object with the workflow's
        required fields. If the result is status=needs_input, ask the caller for
        the missing fields and call again with complete inputs."""
        try:
            inputs = json.loads(inputs_json) if inputs_json else {}
        except json.JSONDecodeError:
            return {"status": "failed", "message": "inputs_json is not valid JSON"}
        return await self._convex.run_workflow(
            workflow_name, inputs, self._cfg.conversation_id or ""
        )

    @function_tool()
    async def escalate_to_human(self, context: RunContext, reason: str) -> dict:
        """Escalate this conversation to a human agent. Call when the caller
        explicitly asks for a person, is angry/frustrated, or you cannot resolve
        the issue. Tell the caller a human will take over."""
        return await self._convex.escalate(self._call_id, reason)

    @function_tool()
    async def transfer_call(self, context: RunContext, reason: str, to: str = "") -> dict:
        """Transfer the live call to a human phone number without dropping the
        caller. Use when escalation rules require a live handoff. Optional `to`
        overrides the agent's configured transfer number or ring group."""
        if not self._room_name:
            return {
                "ok": False,
                "error": {"code": "missing_room", "message": "No LiveKit room for transfer"},
            }
        result = await self._convex.transfer_call(
            self._call_id,
            to=to or None,
            room_name=self._room_name,
            reason=reason,
        )
        if result.get("ok") and self._on_transferred:
            try:
                await self._on_transferred(result)
            except Exception as e:  # noqa: BLE001
                logger.warning("post-transfer handoff failed: %s", e)
        return result

    @function_tool()
    async def record_campaign_answer(self, context: RunContext, question_id: str, answer: str) -> dict:
        """Record the caller's answer to one campaign question, immediately
        after they answer it. `question_id` must be one of the campaign question
        ids from your instructions."""
        attempt_id = (self._cfg.campaign or {}).get("attemptId") or ""
        if not attempt_id:
            return {"ok": False, "error": "no campaign attempt on this call"}
        return await self._convex.record_answer(attempt_id, question_id, answer)


def _call_context(ctx: agents.JobContext) -> dict:
    """Org/call/campaign routing context: job metadata (JSON) > env fallbacks."""
    meta: dict = {}
    try:
        if ctx.job.metadata:
            meta = json.loads(ctx.job.metadata)
    except (json.JSONDecodeError, AttributeError):
        pass
    return {
        "organization_id": meta.get("organizationId") or os.environ.get("VOICE_DEV_ORG_ID", ""),
        "call_id": meta.get("callId") or ctx.room.name,
        "caller_number": meta.get("callerNumber"),
        "campaign_id": meta.get("campaignId"),
        "attempt_id": meta.get("attemptId"),
        "agent_id": meta.get("agentId"),
        "dialed_number": meta.get("dialedNumber"),
        "channel": meta.get("channel"),
        # Browser test calls (Agent Builder) — dialect/language currently
        # selected in the form, which may not be saved to the agent yet.
        "dialect_override": meta.get("dialectOverride"),
        "language_override": meta.get("languageOverride"),
    }


def _sip_attr(participant, *keys: str) -> str | None:
    """Read a SIP attribute from a LiveKit remote participant."""
    attrs = getattr(participant, "attributes", None) or {}
    for key in keys:
        val = attrs.get(key) if hasattr(attrs, "get") else None
        if not val and isinstance(attrs, dict):
            val = attrs.get(key)
        if val:
            return str(val)
    return None


def _extract_sip_numbers(ctx: agents.JobContext) -> tuple[str | None, str | None]:
    """Best-effort (dialed DID, caller ANI) from SIP participants in the room."""
    dialed: str | None = None
    caller: str | None = None
    try:
        for p in ctx.room.remote_participants.values():
            identity = (getattr(p, "identity", None) or "").lower()
            # LiveKit SIP attrs (common keys across versions).
            trunk = _sip_attr(
                p,
                "sip.trunkPhoneNumber",
                "sip.trunkPhoneNumber",
                "sip.to",
                "sip.callTo",
            )
            from_num = _sip_attr(
                p,
                "sip.phoneNumber",
                "sip.from",
                "sip.callFrom",
            )
            if trunk and not dialed:
                dialed = trunk
            if from_num and not caller:
                caller = from_num
            # Fallback: identity like sip_+216...
            if not dialed and "sip" in identity:
                m = re.search(r"\+?\d{8,15}", identity)
                if m:
                    dialed = m.group(0)
    except Exception as e:  # noqa: BLE001
        logger.warning("SIP attr extract failed: %s", e)
    return dialed, caller


async def entrypoint(ctx: agents.JobContext) -> None:
    call = _call_context(ctx)
    if call.get("channel") == "browser-test":
        logger.info(
            "browser-test call room=%s org=%s agent=%s",
            call.get("call_id"),
            call.get("organization_id"),
            call.get("agent_id"),
        )
    await ctx.connect()

    # Story 7 — inbound: resolve org/agent from dialed DID when metadata is thin.
    sip_dialed, sip_caller = _extract_sip_numbers(ctx)
    if sip_dialed and not call.get("dialed_number"):
        call["dialed_number"] = sip_dialed
    if sip_caller and not call.get("caller_number"):
        call["caller_number"] = sip_caller

    convex = ConvexVoiceClient(organization_id=call["organization_id"] or "")

    if (not call["organization_id"] or not call.get("agent_id")) and call.get("dialed_number"):
        try:
            resolved = await convex.resolve_inbound(call["dialed_number"])
            if resolved.get("ok"):
                if not call["organization_id"]:
                    call["organization_id"] = resolved.get("organizationId") or ""
                if not call.get("agent_id") and resolved.get("agentId"):
                    call["agent_id"] = resolved["agentId"]
                if resolved.get("phoneNumber"):
                    call["dialed_number"] = resolved["phoneNumber"]
                convex.set_organization_id(call["organization_id"])
                logger.info(
                    "inbound resolve DID=%s org=%s agent=%s",
                    call["dialed_number"],
                    call["organization_id"],
                    call.get("agent_id"),
                )
            else:
                logger.warning(
                    "inbound resolve failed for %s: %s",
                    call.get("dialed_number"),
                    resolved.get("error"),
                )
        except Exception as e:  # noqa: BLE001
            logger.warning("resolve_inbound error: %s", e)

    if not call["organization_id"]:
        logger.error("no organizationId (metadata, DID resolve, or VOICE_DEV_ORG_ID) — refusing job")
        return

    convex.set_organization_id(call["organization_id"])

    # A. Session bootstrap — gates + persona + tool manifest, one call.
    cfg = await convex.get_session(
        call_id=call["call_id"],
        caller_number=call["caller_number"],
        campaign_id=call["campaign_id"],
        attempt_id=call["attempt_id"],
        agent_id=call.get("agent_id"),
        dialed_number=call.get("dialed_number"),
        dialect_override=call.get("dialect_override"),
        language_override=call.get("language_override"),
    )
    if cfg.agent_id and not call.get("agent_id"):
        call["agent_id"] = cfg.agent_id

    session_voice_id = None
    if isinstance(cfg.voice, dict):
        session_voice_id = (cfg.voice.get("voiceId") or cfg.voice.get("voice_id") or "").strip() or None

    # Phase 3 — Agent Builder barge-in / Whisper dialect hint (same session keeps context).
    session_turn_taking = cfg.turn_taking if isinstance(cfg.turn_taking, dict) else None
    session_transcription = cfg.transcription if isinstance(cfg.transcription, dict) else None

    session = build_agent_session(
        voice_id=session_voice_id,
        turn_taking=session_turn_taking,
        transcription=session_transcription,
    )

    # Suspended / out of credits → speak the refusal and end the call.
    if not cfg.ok:
        logger.warning("session gate=%s — refusing call", cfg.gate)
        refusal = Agent(instructions="Tu es un standard téléphonique.")
        await session.start(agent=refusal, room=ctx.room)
        await session.say(cfg.refusal_message or "Ce service est momentanément indisponible.")
        await session.aclose()
        await convex.aclose()
        return

    # Patch the attemptId through if dispatch knew it but the session didn't.
    if cfg.campaign is not None and call["attempt_id"] and not cfg.campaign.get("attemptId"):
        cfg.campaign["attemptId"] = call["attempt_id"]

    pool = cfg.pool or {}
    pool_status = (pool.get("status") or "no_pool") if isinstance(pool, dict) else "no_pool"

    # ── Story 13: hold when pool is saturated (no tools / no opening script) ──
    if pool_status == "queued":
        position = pool.get("position") or 1
        logger.info(
            "pool queued call=%s position=%s — starting hold loop",
            call["call_id"],
            position,
        )
        hold_agent = Agent(
            instructions=(
                "Tu es un standard téléphonique. Tu ne fais que des messages d'attente. "
                "Ne réponds pas aux questions."
            )
        )
        await session.start(agent=hold_agent, room=ctx.room)
        await session.say(
            f"Tous nos conseillers sont actuellement occupés. "
            f"Vous êtes en position {position}. Merci de patienter."
        )

        hold_started = time.monotonic()
        max_hold_s = float(os.environ.get("VOICE_QUEUE_MAX_HOLD_SECONDS", "600"))
        poll_s = float(os.environ.get("VOICE_QUEUE_POLL_SECONDS", "5"))
        last_announce = hold_started
        promoted = False
        abandoned_hold = False
        room_gone = {"v": False}

        @ctx.room.on("disconnected")
        def _on_hold_disconnect(*_args) -> None:
            room_gone["v"] = True

        while True:
            if room_gone["v"]:
                logger.info("caller left during hold call=%s", call["call_id"])
                try:
                    await convex.queue_abandon(call["call_id"])
                except Exception as e:  # noqa: BLE001
                    logger.warning("queue_abandon on disconnect failed: %s", e)
                abandoned_hold = True
                break

            await asyncio.sleep(poll_s)
            elapsed = time.monotonic() - hold_started
            if elapsed >= max_hold_s:
                logger.info("hold timeout call=%s — abandoning queue", call["call_id"])
                try:
                    await session.say(
                        "Nous sommes désolés, tous nos conseillers restent occupés. "
                        "Veuillez rappeler plus tard."
                    )
                except Exception:  # noqa: BLE001
                    pass
                try:
                    await convex.queue_abandon(call["call_id"])
                except Exception as e:  # noqa: BLE001
                    logger.warning("queue_abandon failed: %s", e)
                abandoned_hold = True
                break

            try:
                st = await convex.queue_status(call["call_id"])
            except Exception as e:  # noqa: BLE001
                logger.warning("queue_status poll failed: %s", e)
                continue

            status = st.get("status")
            if status == "connected":
                logger.info("promoted from queue call=%s", call["call_id"])
                promoted = True
                break
            if status in ("abandoned", "failed"):
                abandoned_hold = True
                break

            pos = st.get("position") or position
            if time.monotonic() - last_announce >= 30:
                try:
                    await session.say(
                        f"Merci de patienter, vous êtes toujours en position {pos}."
                    )
                except Exception:  # noqa: BLE001
                    pass
                last_announce = time.monotonic()

        if abandoned_hold or not promoted:
            try:
                await session.aclose()
            except Exception:  # noqa: BLE001
                pass
            try:
                await convex.archive_call(
                    call_id=call["call_id"],
                    conversation_id=cfg.conversation_id,
                    language=cfg.language,
                )
            except Exception as e:  # noqa: BLE001
                logger.warning("archive after hold abandon failed: %s", e)
            await convex.aclose()
            return

        # Promoted — close hold agent; rebuild session for the real AI agent.
        try:
            await session.aclose()
        except Exception:  # noqa: BLE001
            pass
        session = build_agent_session(
            voice_id=session_voice_id,
            turn_taking=session_turn_taking,
            transcription=session_transcription,
        )
        logger.info("hold ended — starting AI agent for call=%s", call["call_id"])

    # C. Persistence sink — every finalized turn → idempotent /voice/persist.
    turn_counter = {"i": 0}
    session_started_ms = int(time.time() * 1000)

    @session.on("conversation_item_added")
    def _on_item(ev) -> None:
        item = ev.item
        role = getattr(item, "role", None)
        text = (getattr(item, "text_content", None) or "").strip()
        if role not in ("user", "assistant") or not text:
            return
        idx = turn_counter["i"]
        turn_counter["i"] += 1
        now_ms = int(time.time() * 1000)

        async def _persist() -> None:
            try:
                await convex.persist_turn(
                    call["call_id"],
                    idx,
                    role,
                    text,
                    cfg.thread_id or "",
                    timestamp_ms=now_ms,
                    offset_ms=max(0, now_ms - session_started_ms),
                )
            except Exception as e:  # persistence must never kill the call
                logger.warning("persist_turn failed (turn %s): %s", idx, e)

        asyncio.create_task(_persist())

    # ── Latency metrics — per-turn conversational breakdown vs the 600-800ms target.
    _lat: dict[str, dict[str, float]] = {}

    @session.on("metrics_collected")
    def _on_metrics(ev) -> None:
        m = ev.metrics
        lk_metrics.log_metrics(m)
        sid = getattr(m, "speech_id", None)
        if not sid:
            return
        slot = _lat.setdefault(sid, {})
        kind = type(m).__name__
        if kind == "EOUMetrics":
            slot["eou"] = m.end_of_utterance_delay
        elif kind == "LLMMetrics":
            slot["ttft"] = m.ttft
        elif kind == "TTSMetrics":
            slot["ttfb"] = m.ttfb
        if {"eou", "ttft", "ttfb"} <= slot.keys():
            total = slot["eou"] + slot["ttft"] + slot["ttfb"]
            logger.info(
                "[latency] turn %s — EOU %.0f + LLM_TTFT %.0f + TTS_TTFB %.0f = %.0f ms (target 600-800)",
                sid, slot["eou"] * 1000, slot["ttft"] * 1000, slot["ttfb"] * 1000, total * 1000,
            )
            _lat.pop(sid, None)

    room_name = ctx.room.name
    leave_after_transfer = {"done": False}

    async def _on_transferred(_result: dict) -> None:
        """Warm handoff: announce, then shut down AI so human + caller stay in-room."""
        if leave_after_transfer["done"]:
            return
        leave_after_transfer["done"] = True
        try:
            await session.say(
                "Je vous connecte avec un conseiller. Veuillez rester en ligne."
            )
        except Exception as e:  # noqa: BLE001
            logger.warning("transfer announce failed: %s", e)
        try:
            await session.aclose()
        except Exception as e:  # noqa: BLE001
            logger.warning("session close after transfer failed: %s", e)

    agent = CoreDeskVoiceAgent(
        cfg,
        convex,
        call["call_id"],
        room_name=room_name,
        on_transferred=_on_transferred,
    )
    await session.start(agent=agent, room=ctx.room)

    # Story 9 — start RoomCompositeEgress once the real AI session is live.
    egress_state: dict | None = await start_room_recording(room_name, call["call_id"])

    # Open the conversation — script from Agent Builder, else dialect-aware greeting.
    await session.generate_reply(instructions=_opening_reply_instructions(cfg))

    # Archive + post-call eval when the room closes.
    @ctx.room.on("disconnected")
    def _on_disconnect(*_args) -> None:
        async def _archive() -> None:
            recording_url = None
            duration_seconds = None
            try:
                stopped = await stop_room_recording(egress_state)
                recording_url = stopped.get("recording_url")
                duration_seconds = stopped.get("duration_seconds")
            except Exception as e:  # noqa: BLE001
                logger.warning("stop recording failed: %s", e)

            try:
                await convex.archive_call(
                    call_id=call["call_id"],
                    conversation_id=cfg.conversation_id,
                    language=cfg.language,
                    recording_url=recording_url,
                    duration_seconds=duration_seconds,
                )
            except Exception as e:
                logger.warning("archive_call failed: %s", e)

            # Late URL if egress resolved after archive (or archive raced).
            if recording_url:
                try:
                    await convex.patch_recording_url(
                        call["call_id"],
                        recording_url,
                        duration_seconds=duration_seconds,
                    )
                except Exception as e:  # noqa: BLE001
                    logger.warning("patch_recording_url failed: %s", e)

            await convex.aclose()

        asyncio.create_task(_archive())


if __name__ == "__main__":
    # agent_name registers the worker for EXPLICIT dispatch so the Convex outbound
    # route (system/livekitSip.dialViaLiveKit) can dispatch it into a room WITH
    # per-call metadata (org/call/campaign). Must match LIVEKIT_AGENT_NAME on Convex.
    # Inbound dispatch rules should reference the same agent name.
    agents.cli.run_app(
        agents.WorkerOptions(entrypoint_fnc=entrypoint, agent_name="coredesk-voice")
    )
