#!/usr/bin/env python3
"""
Podcast TTS — custom MCP server (streamable HTTP) for Claude.

Exposes one tool, generate_podcast_audio: takes a (possibly segmented) spoken
script, synthesizes it with Gemini TTS, stitches the audio with real silences
at segment boundaries, encodes an mp3 with ffmpeg, uploads it to a public GCS
bucket, and returns ONLY metadata: {url, duration_seconds, size_bytes}.

The audio itself never travels through the MCP channel — that is the whole
point of this design: tool results stay tiny, storage does the heavy lifting.

Env vars:
  GEMINI_API_KEY    required — aistudio.google.com API key
  GCS_BUCKET        required — public bucket for the mp3s
  SECRET_PATH       recommended — random string; server only answers under /<SECRET_PATH>/mcp
  GEMINI_TTS_MODEL  default "gemini-3.1-flash-tts-preview"
  DEFAULT_VOICE     default "Charon"
"""
import base64
import io
import json
import os
import re
import subprocess
import tempfile
import time

import httpx
from google.cloud import storage
from mcp.server.fastmcp import FastMCP
from mcp.server.transport_security import TransportSecuritySettings

GEMINI_API_KEY = os.environ.get("GEMINI_API_KEY", "")
GCS_BUCKET = os.environ.get("GCS_BUCKET", "")
SECRET_PATH = os.environ.get("SECRET_PATH", "").strip("/")
MODEL = os.environ.get("GEMINI_TTS_MODEL", "gemini-3.1-flash-tts-preview")
DEFAULT_VOICE = os.environ.get("DEFAULT_VOICE", "Achird")

DIALOGUE_VOICE_A = os.environ.get("DIALOGUE_VOICE_A", "Puck")
DIALOGUE_VOICE_B = os.environ.get("DIALOGUE_VOICE_B", "Kore")
RESEND_API_KEY = os.environ.get("RESEND_API_KEY", "")
DEFAULT_STYLE = os.environ.get(
    "DEFAULT_STYLE",
    "Speak at a calm, deliberate professional pace - unhurried, with a clear beat between sentences, but still moving forward rather than dragging. Ease off on numbers so each figure registers. Keep the delivery understated and even-toned: minimal pitch variation, restrained emphasis, no dramatic swings or theatrical highs and lows. Minimize audible breaths.",
)
CDN = "https://cdn.jsdelivr.net/npm"

SAMPLE_RATE = 24000  # Gemini TTS returns 16-bit mono PCM at 24 kHz
BYTES_PER_SEC = SAMPLE_RATE * 2
MAX_CHUNK_CHARS = 2200  # smaller requests: the preview model returns empty
                        # responses more often as the payload grows

# A dialogue line looks like "BUYER: November Santos, what have you got?"
SPEAKER_RE = re.compile(r"^([A-Z][A-Z0-9 _-]{1,14}):\s*(.+)$")
GEMINI_URL = (
    "https://generativelanguage.googleapis.com/v1beta/models/"
    "{model}:generateContent?key={key}"
)

# DNS-rebinding protection is designed for local servers ("Host must be
# localhost"); on Cloud Run the Host is the public run.app domain, so the
# default check rejects every request with 421. The URL secret path is this
# server's actual access control, so we disable the host check explicitly.
mcp = FastMCP(
    "podcast-tts",
    stateless_http=True,
    transport_security=TransportSecuritySettings(enable_dns_rebinding_protection=False),
)


# ----------------------------------------------------------------- script ----

def parse_segments(script: str):
    """Accepts 'TEXT ||| PAUSE' lines or plain text. Returns [(text, pause)]."""
    segments = []
    for raw in script.splitlines():
        line = raw.strip()
        if not line or line.startswith("#"):
            continue
        if "|||" in line:
            text, _, pause_s = line.rpartition("|||")
            try:
                pause = max(0.0, min(float(pause_s.strip()), 3.0))
            except ValueError:
                pause = 0.4
            text = text.strip()
        else:
            text, pause = line, 0.4
        if text:
            segments.append((text, pause))
    return segments


def is_dialogue(text: str) -> bool:
    return bool(SPEAKER_RE.match(text.strip()))


def group_chunks(segments):
    """Group segments into API-sized chunks, breaking at strong pauses.

    Returns [(chunk_text, trailing_pause, kind)] where kind is "narration" or
    "dialogue". Dialogue lines ("SPEAKER: text") are never mixed with narration
    in the same chunk, so they can be rendered with distinct voices.
    """
    chunks, buf, buf_len, buf_kind = [], [], 0, None

    def flush(pause):
        if buf:
            joined = "\n".join(buf) if buf_kind == "dialogue" else " ".join(buf)
            chunks.append((joined, pause, buf_kind))

    for i, (text, pause) in enumerate(segments):
        kind = "dialogue" if is_dialogue(text) else "narration"
        if buf and kind != buf_kind:
            flush(0.35)  # small beat when switching between narration and dialogue
            buf, buf_len = [], 0
        buf.append(text)
        buf_kind = kind
        buf_len += len(text)
        if pause >= 0.6 or buf_len >= MAX_CHUNK_CHARS or i == len(segments) - 1:
            flush(pause)
            buf, buf_len = [], 0
    return chunks


# -------------------------------------------------------------------- tts ----

def speech_config(text: str, voice: str, kind: str, voice_a: str, voice_b: str):
    """Single-voice config for narration; two-voice config for dialogue chunks."""
    if kind != "dialogue":
        return {"voiceConfig": {"prebuiltVoiceConfig": {"voiceName": voice}}}, None

    speakers = []
    for line in text.splitlines():
        m = SPEAKER_RE.match(line.strip())
        if m and m.group(1) not in speakers:
            speakers.append(m.group(1))
    if not speakers:  # defensive: fall back to narration
        return {"voiceConfig": {"prebuiltVoiceConfig": {"voiceName": voice}}}, None

    # Gemini supports up to 2 distinct speakers; extras alternate between them.
    voices = [voice_a, voice_b]
    configs = [
        {
            "speaker": name,
            "voiceConfig": {"prebuiltVoiceConfig": {"voiceName": voices[i % 2]}},
        }
        for i, name in enumerate(speakers[:2])
    ]
    return {"multiSpeakerVoiceConfig": {"speakerVoiceConfigs": configs}}, speakers[:2]


def tts_chunk(text: str, voice: str, style: str, kind: str = "narration",
              voice_a: str = "", voice_b: str = "") -> bytes:
    """One Gemini TTS call -> raw PCM bytes. Retries transient failures."""
    cfg, speakers = speech_config(
        text, voice, kind, voice_a or DIALOGUE_VOICE_A, voice_b or DIALOGUE_VOICE_B
    )
    if speakers:
        lead = (
            f"TTS the following short trading-desk exchange between "
            f"{' and '.join(speakers)}. Deliver it naturally, as real market chatter: "
            f"brisk, clipped, matter-of-fact, no theatrical acting.\n\n"
        )
        prompt = lead + text
    else:
        prompt = f"{style.strip()}\n\n{text}" if style.strip() else text
    body = {
        "contents": [{"parts": [{"text": prompt}]}],
        "generationConfig": {
            "responseModalities": ["AUDIO"],
            "speechConfig": cfg,
        },
    }
    url = GEMINI_URL.format(model=MODEL, key=GEMINI_API_KEY)
    last_error = None
    for attempt in range(7):
        try:
            resp = httpx.post(url, json=body, timeout=120)
            if resp.status_code == 429:
                time.sleep(15 * (attempt + 1))
                last_error = "429 rate limited"
                continue
            if resp.status_code >= 400:
                # Surface the API's own message — essential for diagnosis.
                last_error = f"HTTP {resp.status_code}: {resp.text[:700]}"
                # The preview TTS model returns sporadic 400 INVALID_ARGUMENT on
                # perfectly valid input, so retry those too. Only auth errors are
                # genuinely terminal.
                if resp.status_code in (401, 403):
                    break
                time.sleep(4 * (attempt + 1))
                continue
            data = resp.json()
            cands = data.get("candidates") or []
            if not cands:
                # 200 with no candidates: the preview model does this
                # intermittently. Retryable.
                last_error = f"empty response (no candidates): {str(data)[:300]}"
                time.sleep(6 * (attempt + 1))
                continue
            part = cands[0]["content"]["parts"][0]
            return base64.b64decode(part["inlineData"]["data"])
        except Exception as exc:  # noqa: BLE001
            last_error = f"{type(exc).__name__}: {str(exc)[:500]}"
            time.sleep(6 * (attempt + 1))

    if speakers:
        # A dialogue must never cost us the episode: fall back to the narrator
        # reading the exchange, speaker labels stripped.
        plain = "\n".join(
            (m.group(2) if (m := SPEAKER_RE.match(l.strip())) else l)
            for l in text.splitlines()
        )
        return tts_chunk(plain, voice, style, "narration", voice_a, voice_b)

    raise RuntimeError(f"Gemini TTS chunk failed after 7 attempts — {last_error}")


def encode_mp3(pcm: bytes, bitrate: str = "96k") -> bytes:
    """Raw s16le PCM -> mp3 via ffmpeg."""
    with tempfile.NamedTemporaryFile(suffix=".mp3", delete=False) as out:
        out_path = out.name
    try:
        subprocess.run(
            [
                "ffmpeg", "-y", "-loglevel", "error",
                "-f", "s16le", "-ar", str(SAMPLE_RATE), "-ac", "1", "-i", "pipe:0",
                "-b:a", bitrate, out_path,
            ],
            input=pcm,
            check=True,
        )
        return open(out_path, "rb").read()
    finally:
        os.unlink(out_path)


def upload_public(data: bytes, object_name: str, content_type: str = "audio/mpeg",
                  cache_control: str = "public, max-age=3600") -> str:
    client = storage.Client()
    blob = client.bucket(GCS_BUCKET).blob(object_name)
    blob.cache_control = cache_control
    blob.upload_from_string(data, content_type=content_type)
    return f"https://storage.googleapis.com/{GCS_BUCKET}/{object_name}"


def read_meta(show: str, episode: str):
    """Read the sidecar status json for an episode, or None."""
    client = storage.Client()
    blob = client.bucket(GCS_BUCKET).blob(f"{show}/{episode}.json")
    if not blob.exists():
        return None
    return json.loads(blob.download_as_bytes())


def write_meta(show: str, episode: str, meta: dict):
    upload_public(json.dumps(meta).encode(), f"{show}/{episode}.json",
                  content_type="application/json")


# ------------------------------------------------------------------- tool ----

def _generate_impl(script, voice, style, show, episode, voice_a="", voice_b="") -> str:
    if not GEMINI_API_KEY or not GCS_BUCKET:
        return json.dumps({"error": "server misconfigured: missing GEMINI_API_KEY or GCS_BUCKET"})

    segments = parse_segments(script)
    if not segments:
        return json.dumps({"error": "empty script"})
    chunks = group_chunks(segments)

    safe_show = re.sub(r"[^a-z0-9_-]", "-", show.lower()) or "default"
    safe_ep = re.sub(r"[^a-z0-9_.-]", "-", episode.lower()) or "episode"
    voice = voice or DEFAULT_VOICE

    # Status sidecar: lets a client that timed out on this (long) call recover
    # the result later via check_audio_ready.
    n_dialogue = sum(1 for c in chunks if c[2] == "dialogue")
    write_meta(safe_show, safe_ep,
               {"status": "working", "chunks": len(chunks), "dialogue_chunks": n_dialogue})

    try:
        pcm = io.BytesIO()
        failed = []
        for i, (text, trailing_pause, kind) in enumerate(chunks):
            try:
                pcm.write(tts_chunk(text, voice, style, kind, voice_a, voice_b))
            except Exception as exc:  # noqa: BLE001
                # One stubborn chunk must not cost the whole episode: leave a
                # short gap, record it, and keep going.
                failed.append({"chunk": i + 1, "error": str(exc)[:300],
                               "preview": text[:80]})
                pcm.write(b"\x00" * int(BYTES_PER_SEC * 0.4))
            if i < len(chunks) - 1 and trailing_pause > 0:
                pcm.write(b"\x00" * int(BYTES_PER_SEC * trailing_pause))
        if len(failed) == len(chunks):
            raise RuntimeError(f"every chunk failed; first error: {failed[0]['error']}")

        raw = pcm.getvalue()
        duration = int(len(raw) / BYTES_PER_SEC)
        mp3 = encode_mp3(raw)
        url = upload_public(mp3, f"{safe_show}/{safe_ep}.mp3")
    except Exception as exc:  # noqa: BLE001
        write_meta(safe_show, safe_ep, {"status": "error", "error": str(exc)[:500]})
        raise

    meta = {
        "status": "done",
        "url": url,
        "duration_seconds": duration,
        "size_bytes": len(mp3),
        "chunks": len(chunks),
        "dialogue_chunks": n_dialogue,
        "failed_chunks": failed,
        "voice": voice,
        "model": MODEL,
    }
    write_meta(safe_show, safe_ep, meta)
    return json.dumps(meta)


@mcp.tool()
async def generate_podcast_audio(
    script: str,
    voice: str = "",
    style: str = "Narrate calmly and clearly, like a professional podcast host. Mark a brief pause between ideas.",
    show: str = "default",
    episode: str = "episode",
    dialogue_voice_a: str = "",
    dialogue_voice_b: str = "",
) -> str:
    """Generate podcast audio from a script with Gemini TTS and host it publicly.

    Supports two-voice dialogue: any script line written as "SPEAKER: text"
    (uppercase label, e.g. "BUYER: November Santos, what have you got?") is
    rendered as a two-person exchange with distinct voices, separate from the
    narrator. Consecutive dialogue lines are kept together; up to two distinct
    speakers per exchange.

    Args:
        script: The spoken script. Either plain text, or one segment per line in
            the format "TEXT ||| PAUSE" where PAUSE (seconds) is inserted as real
            silence after the segment (pauses >= 0.6 s become chunk boundaries).
            Dialogue lines use "SPEAKER: text ||| PAUSE".
        voice: Narrator voice — Gemini prebuilt name (Achird, Charon, Kore,
            Puck, Fenrir...). Empty = server default.
        style: Natural-language delivery instruction for narration.
        show: Show slug — becomes the storage folder (one per podcast).
        episode: Episode slug — becomes the file name, e.g. "ep02".
        dialogue_voice_a: Voice for the first speaker in an exchange (default Puck).
        dialogue_voice_b: Voice for the second speaker (default Kore).

    Returns JSON: {"url", "duration_seconds", "size_bytes", "chunks", "model"}.
    The audio itself is uploaded to public storage; only this metadata returns.

    Runs the blocking work (HTTP to Gemini, ffmpeg, GCS upload) in a worker
    thread so concurrent MCP calls are served without blocking the event loop.
    """
    import anyio

    return await anyio.to_thread.run_sync(
        _generate_impl, script, voice, style, show, episode,
        dialogue_voice_a, dialogue_voice_b
    )


@mcp.tool()
async def generate_episode_audio(package: str, version: str, show: str,
                                 episode: str, script_file: str = "script.txt") -> str:
    """Generate the episode audio from the script already published to npm.

    Identical to generate_podcast_audio, except the server fetches the script
    itself and uses the show's configured narrator voice, dialogue voices and
    delivery style from its own environment. Nothing about the voice or the
    script has to be restated by the caller, so it cannot drift between runs.

    Same asynchronous behaviour: the call may exceed the client timeout, in
    which case the work continues server-side — poll check_audio_ready. Never
    call it twice for the same episode.

    Args:
        package: npm package the toolkit published, e.g. "@sdelsad/commodity-desk-daily"
        version: the exact version
        show: storage folder, e.g. "commodity-desk-daily"
        episode: episode slug, e.g. "ep07"
        script_file: file inside the package, default "script.txt"
    """
    import anyio

    def work():
        try:
            script = _fetch(f"{CDN}/{package}@{version}/{script_file}", 2).decode("utf-8")
        except Exception as exc:  # noqa: BLE001
            return json.dumps({"status": "error",
                               "error": f"could not read the script: {str(exc)[:200]}"})
        return _generate_impl(script, DEFAULT_VOICE, DEFAULT_STYLE,
                              show, episode, DIALOGUE_VOICE_A, DIALOGUE_VOICE_B)

    return await anyio.to_thread.run_sync(work)


def _probe_impl(text: str, voice_a: str, voice_b: str) -> str:
    """Diagnostic: try a multi-speaker call and return the raw API outcome."""
    cfg, speakers = speech_config(text, "Achird", "dialogue",
                                  voice_a or DIALOGUE_VOICE_A, voice_b or DIALOGUE_VOICE_B)
    body = {
        "contents": [{"parts": [{"text": text}]}],
        "generationConfig": {"responseModalities": ["AUDIO"], "speechConfig": cfg},
    }
    url = GEMINI_URL.format(model=MODEL, key=GEMINI_API_KEY)
    try:
        resp = httpx.post(url, json=body, timeout=120)
        ok = resp.status_code == 200
        return json.dumps({
            "speakers": speakers,
            "status_code": resp.status_code,
            "ok": ok,
            "body": "" if ok else resp.text[:1500],
            "config_sent": cfg,
        })
    except Exception as exc:  # noqa: BLE001
        return json.dumps({"error": f"{type(exc).__name__}: {str(exc)[:500]}"})


@mcp.tool()
async def probe_multispeaker(
    text: str = "BUYER: November Santos, what have you got?\nSELLER: I make you plus eighty-five.",
    voice_a: str = "",
    voice_b: str = "",
) -> str:
    """Diagnostic tool: attempt one multi-speaker TTS call and report the raw result.

    Returns the HTTP status, the exact config sent, and the API's error body if
    it failed. Use this to debug two-voice dialogue without generating audio.
    """
    import anyio

    return await anyio.to_thread.run_sync(_probe_impl, text, voice_a, voice_b)


def _check_impl(show: str, episode: str) -> str:
    safe_show = re.sub(r"[^a-z0-9_-]", "-", show.lower()) or "default"
    safe_ep = re.sub(r"[^a-z0-9_.-]", "-", episode.lower()) or "episode"
    meta = read_meta(safe_show, safe_ep)
    if meta:
        return json.dumps(meta)
    return json.dumps({"status": "unknown"})


def _publish_feed_impl(show: str, feed_xml: str) -> str:
    safe_show = re.sub(r"[^a-z0-9_-]", "-", show.lower()) or "default"
    url = upload_public(
        feed_xml.encode("utf-8"),
        f"{safe_show}/feed.xml",
        content_type="application/rss+xml; charset=utf-8",
        cache_control="public, max-age=300",  # 5 min: podcast apps see updates fast
    )
    return json.dumps({"feed_url": url})


@mcp.tool()
async def publish_feed(show: str, feed_xml: str) -> str:
    """Host the podcast RSS feed on public storage with a 5-minute cache.

    Pass the full feed.xml content. Returns {"feed_url": ...} — a stable URL
    podcast apps can subscribe to, refreshed within minutes of each publish
    (unlike CDN-cached copies which can lag by hours).
    """
    import anyio

    return await anyio.to_thread.run_sync(_publish_feed_impl, show, feed_xml)


def _publish_page_impl(show: str, episode: str, page_html: str) -> str:
    safe_show = re.sub(r"[^a-z0-9_-]", "-", show.lower()) or "default"
    safe_ep = re.sub(r"[^a-z0-9_.-]", "-", episode.lower()) or "episode"
    url = upload_public(
        page_html.encode("utf-8"),
        f"{safe_show}/{safe_ep}.html",
        content_type="text/html; charset=utf-8",
        cache_control="public, max-age=300",
    )
    return json.dumps({"page_url": url})


@mcp.tool()
async def publish_page(show: str, episode: str, page_html: str) -> str:
    """Host an episode's HTML page on public storage and return its URL.

    Pass the full HTML document (as produced by the toolkit's build_page.py).
    Returns {"page_url": ...} — a stable link to put in the episode email.
    Re-publishing the same episode replaces the page (5-minute cache).
    """
    import anyio

    return await anyio.to_thread.run_sync(_publish_page_impl, show, episode, page_html)


def _publish_asset_impl(path: str, data_b64: str, content_type: str,
                        cache_seconds: int) -> str:
    import base64

    safe = re.sub(r"[^a-zA-Z0-9_./-]", "-", path.lstrip("/")) or "asset"
    if ".." in safe:
        return json.dumps({"error": "path may not contain '..'"})
    url = upload_public(
        base64.b64decode(data_b64),
        safe,
        content_type=content_type,
        cache_control=f"public, max-age={max(int(cache_seconds), 0)}",
    )
    return json.dumps({"url": url})


@mcp.tool()
async def publish_asset(path: str, data_b64: str,
                        content_type: str = "image/png",
                        cache_seconds: int = 86400) -> str:
    """Host any small file on public storage and return its URL.

    Use this for episode chart images, cover art, or a static page. `path` is
    the object path inside the bucket, e.g. "commodity-desk-daily/ep09-1.png"
    or "index.html". `data_b64` is the file's bytes, base64-encoded.
    Re-publishing the same path replaces the file. Returns {"url": ...}.

    Charts belong in the e-mail as PNG, because mail clients do not render
    inline SVG: publish the PNG here and point an <img src> at the URL.
    """
    import anyio

    return await anyio.to_thread.run_sync(
        _publish_asset_impl, path, data_b64, content_type, cache_seconds
    )


def _publish_from_url_impl(path: str, source_url: str, content_type: str,
                           cache_seconds: int) -> str:
    safe = re.sub(r"[^a-zA-Z0-9_./-]", "-", path.lstrip("/")) or "asset"
    if ".." in safe:
        return json.dumps({"error": "path may not contain '..'"})
    if not source_url.startswith(("https://cdn.jsdelivr.net/", "https://registry.npmjs.org/",
                                  "https://unpkg.com/", "https://storage.googleapis.com/")):
        return json.dumps({"error": "source_url must be on jsdelivr, unpkg, "
                                    "the npm registry or this bucket"})
    try:
        resp = httpx.get(source_url, timeout=120, follow_redirects=True)
        resp.raise_for_status()
    except Exception as exc:  # noqa: BLE001
        return json.dumps({"error": f"fetch failed: {str(exc)[:200]}"})
    body = resp.content
    if len(body) > 25 * 1024 * 1024:
        return json.dumps({"error": "source is larger than 25 MB"})
    ctype = content_type or resp.headers.get("content-type", "application/octet-stream")
    url = upload_public(body, safe, content_type=ctype,
                        cache_control=f"public, max-age={max(int(cache_seconds), 0)}")
    return json.dumps({"url": url, "bytes": len(body)})


@mcp.tool()
async def publish_from_url(path: str, source_url: str, content_type: str = "",
                           cache_seconds: int = 300) -> str:
    """Copy a file from a public URL straight into the bucket, server-side.

    Same result as publish_asset, without the caller ever holding the bytes:
    use it for episode pages, chart images and any other large artefact that
    has already been published to npm/jsDelivr. `path` is the object path in
    the bucket, e.g. "commodity-desk-daily/ep06.html". `source_url` must be on
    jsDelivr, unpkg, the npm registry, or this bucket. Returns {"url", "bytes"}.

    Prefer this over publish_asset for anything above a few kilobytes.
    """
    import anyio

    return await anyio.to_thread.run_sync(
        _publish_from_url_impl, path, source_url, content_type, cache_seconds
    )


CTYPES = {
    ".js": "application/javascript",
    ".html": "text/html; charset=utf-8",
    ".xml": "application/rss+xml; charset=utf-8",
    ".png": "image/png",
    ".jpg": "image/jpeg",
    ".json": "application/json",
    ".txt": "text/plain; charset=utf-8",
    ".mp3": "audio/mpeg",
}


def _fetch(url: str, limit_mb: int = 25) -> bytes:
    resp = httpx.get(url, timeout=120, follow_redirects=True)
    resp.raise_for_status()
    if len(resp.content) > limit_mb * 1024 * 1024:
        raise ValueError(f"{url} is larger than {limit_mb} MB")
    return resp.content


def _release_impl(package: str, version: str, prefix: str) -> str:
    """Copy every file listed in a package's release.json into the bucket."""
    base = f"{CDN}/{package}@{version}"
    try:
        manifest = json.loads(_fetch(f"{base}/release.json", 1))
    except Exception as exc:  # noqa: BLE001
        return json.dumps({"error": f"could not read release.json: {str(exc)[:200]}"})

    published, failed = {}, {}
    for entry in manifest.get("files", []):
        src = entry.get("file")
        dest = entry.get("path") or f"{prefix.strip('/')}/{src}"
        if not src or ".." in dest:
            failed[str(src)] = "bad entry"
            continue
        ext = os.path.splitext(src)[1].lower()
        ctype = entry.get("content_type") or CTYPES.get(ext, "application/octet-stream")
        cache = int(entry.get("cache_seconds", 300))
        try:
            url = upload_public(_fetch(f"{base}/{src}"), dest.lstrip("/"),
                                content_type=ctype,
                                cache_control=f"public, max-age={max(cache, 0)}")
            published[src] = url
        except Exception as exc:  # noqa: BLE001
            failed[src] = str(exc)[:200]

    return json.dumps({"published": published, "failed": failed,
                       "count": len(published)})


@mcp.tool()
async def publish_release(package: str, version: str, prefix: str = "") -> str:
    """Copy a whole episode release from npm into the bucket, server-side.

    The toolkit publishes every artefact of an episode to npm (page, chart
    images, feed, e-mail bodies) together with a `release.json` manifest that
    says where each file belongs. This tool reads that manifest and does all
    the uploads itself, so the caller never handles a single byte.

    One call replaces every publish_page / publish_asset / publish_feed call
    for the day. Returns {"published": {file: url}, "failed": {...}, "count"}.

    Args:
        package: npm package name, e.g. "@sdelsad/commodity-desk-daily"
        version: the exact version the toolkit just published, e.g. "1.0.19"
        prefix: bucket folder for entries whose manifest has no explicit path
    """
    import anyio

    return await anyio.to_thread.run_sync(_release_impl, package, version, prefix)


def _send_email_impl(package: str, version: str, subject: str, to: str,
                     sender: str, attachment_url: str, attachment_name: str) -> str:
    if not RESEND_API_KEY:
        return json.dumps({"error": "RESEND_API_KEY is not set on the server"})
    base = f"{CDN}/{package}@{version}"
    try:
        html = _fetch(f"{base}/email.html", 5).decode("utf-8")
        text = _fetch(f"{base}/email.txt", 5).decode("utf-8")
    except Exception as exc:  # noqa: BLE001
        return json.dumps({"error": f"could not read the e-mail bodies: {str(exc)[:200]}"})

    payload = {"from": sender, "to": [to], "subject": subject,
               "html": html, "text": text}
    if attachment_url:
        payload["attachments"] = [{"path": attachment_url,
                                   "filename": attachment_name or "episode.mp3"}]
    try:
        resp = httpx.post("https://api.resend.com/emails", timeout=120,
                          headers={"Authorization": f"Bearer {RESEND_API_KEY}",
                                   "Content-Type": "application/json"},
                          json=payload)
    except Exception as exc:  # noqa: BLE001
        return json.dumps({"error": f"resend request failed: {str(exc)[:200]}"})
    if resp.status_code >= 300:
        return json.dumps({"error": f"resend {resp.status_code}: {resp.text[:300]}",
                           "html_bytes": len(html)})
    return json.dumps({"sent": True, "id": resp.json().get("id"),
                       "html_bytes": len(html), "text_bytes": len(text)})


@mcp.tool()
async def send_episode_email(package: str, version: str, subject: str,
                             to: str, sender: str = "",
                             attachment_url: str = "",
                             attachment_name: str = "") -> str:
    """Send the episode e-mail, reading both bodies from npm server-side.

    Reads `email.html` and `email.txt` out of the published package and posts
    them to Resend as the HTML and plain-text parts of ONE message, so the
    two can never disagree and neither ever passes through the caller. The
    HTML part is what the reader sees; text is only a fallback.

    Returns {"sent": true, "id": ..., "html_bytes": ..., "text_bytes": ...}.

    Args:
        package: npm package name the toolkit published
        version: the exact version
        subject: e-mail subject line
        to: recipient address
        sender: From header, defaults to the show address
        attachment_url: optional public mp3 URL to attach
        attachment_name: filename for that attachment
    """
    import anyio

    sender = sender or "Soft Commodity Trading <send@luxup.ai>"
    return await anyio.to_thread.run_sync(
        _send_email_impl, package, version, subject, to, sender,
        attachment_url, attachment_name
    )


def _purge_impl(paths) -> str:
    results = {}
    for path in paths:
        try:
            resp = httpx.get(f"https://purge.jsdelivr.net/{path.lstrip('/')}", timeout=30)
            results[path] = resp.status_code
        except Exception as exc:  # noqa: BLE001
            results[path] = str(exc)[:200]
    return json.dumps(results)


@mcp.tool()
async def purge_feed_cache(
    paths: list[str] = ["npm/@sdelsad/commodity-desk-daily@latest/feed.xml"],
) -> str:
    """Purge jsDelivr's CDN cache for the given paths (default: the podcast RSS feed).

    Call this right after publishing a new episode so podcast apps see the
    updated feed immediately instead of waiting out the CDN cache (up to ~12 h).
    Returns per-path HTTP status codes from the purge API.
    """
    import anyio

    return await anyio.to_thread.run_sync(_purge_impl, paths)



# --------------------------------------------------------------- finalize ----
#
# The daily run writes epNN.md + the script and publishes them to npm. Every
# episode that ever broke did so AFTER that point: audio never uploaded, page
# never built, e-mail never sent. finalize_episode does that half entirely on
# the server, from the npm package alone, and is idempotent: anything already
# in the bucket is left alone, so it can run every hour from Cloud Scheduler
# and at the end of every run without doing harm.

TOOLKIT_PKG = "@sdelsad/commodity-desk-toolkit"
TOOLKIT_FILES = ["build_page.py", "chart.py", "build_email.py", "conversions.md"]
EMAIL_TO = os.environ.get("EMAIL_TO", "seb.ge.ed@gmail.com")
EMAIL_FROM = os.environ.get("EMAIL_FROM", "Soft Commodity Trading <send@luxup.ai>")
SHOW_NAME = os.environ.get("SHOW_NAME", "Soft Commodity Trading")


def _blob_exists(object_name: str) -> bool:
    return storage.Client().bucket(GCS_BUCKET).blob(object_name).exists()


def _jsdelivr_files(package: str, version: str):
    """File list of an npm version, from jsDelivr's data API."""
    url = f"https://data.jsdelivr.com/v1/package/npm/{package}@{version}/flat"
    resp = httpx.get(url, timeout=60)
    resp.raise_for_status()
    return [f["name"].lstrip("/") for f in resp.json().get("files", [])]


def _resolve_version(package: str, version: str) -> str:
    if version and version != "latest":
        return version
    resp = httpx.get(f"https://registry.npmjs.org/{package}/latest", timeout=60)
    resp.raise_for_status()
    return resp.json()["version"]


def _parse_first_item(feed: str):
    m = re.search(r"<item>.*?</item>", feed, re.S)
    if not m:
        return None
    it = m.group(0)
    g = lambda pat: (re.search(pat, it, re.S) or [None, None])[1]  # noqa: E731
    title = g(r"<title>(.*?)</title>") or ""
    num = re.search(r"Ep\s+(\d+)", title)
    enc = re.search(r'<enclosure url="([^"]+)" length="(\d+)"', it)
    pub = g(r"<pubDate>(.*?)</pubDate>") or ""
    desc = g(r"<itunes:summary>(.*?)</itunes:summary>") or ""
    return {
        "item": it,
        "number": int(num.group(1)) if num else 0,
        "title": re.sub(r"^Ep\s+\d+\s*[—-]\s*", "", title).strip(),
        "enclosure": enc.group(1) if enc else "",
        "length": int(enc.group(2)) if enc else 0,
        "duration": int(g(r"<itunes:duration>(\d+)</itunes:duration>") or 0),
        "pubdate": pub,
        "summary": desc.split("\n")[0].strip(),
    }


def _pretty_date(pubdate: str) -> str:
    from datetime import datetime
    try:
        d = datetime.strptime(pubdate[:16], "%a, %d %b %Y")
        return d.strftime("%A %-d %B %Y")
    except ValueError:
        return ""


def _finalize_impl(package: str, version: str, show: str, send_email: bool) -> str:
    report = {"package": package, "steps": []}
    log = report["steps"].append
    try:
        version = _resolve_version(package, version)
        report["version"] = version
        base = f"{CDN}/{package}@{version}"
        files = _jsdelivr_files(package, version)
        feed = _fetch(f"{base}/feed.xml", 2).decode("utf-8")
    except Exception as exc:  # noqa: BLE001
        report["error"] = f"could not read the package: {str(exc)[:200]}"
        return json.dumps(report)

    item = _parse_first_item(feed)
    if not item or not item["number"]:
        report["error"] = "no episode item in feed.xml"
        return json.dumps(report)
    n = item["number"]
    ep = f"ep{n:02d}"
    report["episode"] = ep
    safe_show = re.sub(r"[^a-z0-9_-]", "-", show.lower()) or "default"
    notes = f"{ep}.md" if f"{ep}.md" in files else ""
    script = next((f for f in (f"{ep}.script.txt", "script.txt") if f in files), "")
    report["notes"] = notes
    report["script"] = script

    # -- 1. audio ----------------------------------------------------------
    enc = item["enclosure"]
    enc_obj = enc.split(f"storage.googleapis.com/{GCS_BUCKET}/")[-1] if GCS_BUCKET in enc else ""
    audio_url, duration, size = enc, item["duration"], item["length"]
    feed_changed = False
    if enc_obj and _blob_exists(enc_obj):
        log(f"audio present: {enc_obj}")
    elif script:
        meta = read_meta(safe_show, ep) or {}
        if meta.get("status") == "working":
            log("audio generation already in progress; try again later")
            return json.dumps(report)
        if meta.get("status") == "done" and meta.get("url"):
            log("audio already generated by an earlier call")
        else:
            log(f"audio missing at {enc_obj or enc}: generating from {script}")
            text = _fetch(f"{base}/{script}", 2).decode("utf-8")
            meta = json.loads(_generate_impl(text, DEFAULT_VOICE, DEFAULT_STYLE,
                                             safe_show, ep, DIALOGUE_VOICE_A,
                                             DIALOGUE_VOICE_B))
            if meta.get("status") != "done":
                report["error"] = f"audio generation failed: {meta.get('error', '')[:200]}"
                return json.dumps(report)
            log(f"audio generated: {meta['duration_seconds']} s, "
                f"{len(meta.get('failed_chunks', []))} failed chunk(s)")
        audio_url, duration, size = meta["url"], meta["duration_seconds"], meta["size_bytes"]
        # the feed promised a file that did not exist: point it at the real one
        new_item = item["item"]
        new_item = re.sub(r'<enclosure url="[^"]+" length="\d+"',
                          f'<enclosure url="{audio_url}" length="{size}"', new_item)
        new_item = re.sub(r"<guid isPermaLink=\"false\">[^<]+</guid>",
                          f'<guid isPermaLink="false">{audio_url}</guid>', new_item)
        new_item = re.sub(r"<itunes:duration>\d+</itunes:duration>",
                          f"<itunes:duration>{duration}</itunes:duration>", new_item)
        feed = feed.replace(item["item"], new_item)
        feed_changed = True
    else:
        report["error"] = "audio missing and the package has no script to generate it from"
        return json.dumps(report)

    page_url = f"https://storage.googleapis.com/{GCS_BUCKET}/{safe_show}/{ep}.html"
    if f"<link>{page_url}</link>" not in feed:
        # the item was published without its website link: add it
        new_item = feed[feed.find("<item>"):feed.find("</item>") + 7]
        if "<link>" not in new_item:
            fixed = new_item.replace("</title>", f"</title>\n      <link>{page_url}</link>", 1)
            feed = feed.replace(new_item, fixed)
            feed_changed = True
    if feed_changed:
        _publish_feed_impl(safe_show, feed)
        log("feed republished with the corrected item")
    elif not _blob_exists(f"{safe_show}/feed.xml"):
        _publish_feed_impl(safe_show, feed)
        log("feed published")
    else:
        # the bucket feed may be older than the package's (run died before publish_feed)
        try:
            live = _fetch(f"https://storage.googleapis.com/{GCS_BUCKET}/{safe_show}/feed.xml?f={int(time.time())}", 2).decode("utf-8")
            if f"Ep {n} " not in live and f"Ep {n}—" not in live and f"<title>Ep {n} " not in live:
                _publish_feed_impl(safe_show, feed)
                log("bucket feed was behind the package; republished")
            else:
                log("feed already live")
        except Exception as exc:  # noqa: BLE001
            log(f"could not compare feeds: {str(exc)[:100]}")

    # -- 2. page + charts + e-mail, rendered by the toolkit ------------------
    if not notes:
        log("no notes in the package: page and e-mail skipped")
        return json.dumps(report)
    page_missing = not _blob_exists(f"{safe_show}/{ep}.html")
    charts_missing = not _blob_exists(f"{safe_show}/{ep}_chart1.png")
    email_wanted = send_email and RESEND_API_KEY and not (read_meta(safe_show, ep) or {}).get("email_sent")
    if not (page_missing or charts_missing or email_wanted):
        log("page, charts and e-mail already done")
        return json.dumps(report)

    work = tempfile.mkdtemp(prefix="finalize-")
    try:
        tk_version = _resolve_version(TOOLKIT_PKG, "latest")
        for f in TOOLKIT_FILES:
            open(os.path.join(work, f), "wb").write(_fetch(f"{CDN}/{TOOLKIT_PKG}@{tk_version}/{f}", 5))
        for f in (notes, "covered.md", "glossary.md"):
            if f in files:
                open(os.path.join(work, f), "wb").write(_fetch(f"{base}/{f}", 5))
        title, dek = item["title"], item["summary"]
        date = _pretty_date(item["pubdate"])
        cmd = ["python3", "build_page.py", "--notes", notes, "--number", str(n),
               "--title", title, "--dek", dek, "--audio-url", audio_url,
               "--duration", str(duration), "--date", date,
               "--out", f"{ep}.html", "--charts-prefix", f"{ep}_chart"]
        r = subprocess.run(cmd, cwd=work, capture_output=True, text=True, timeout=300)
        if r.returncode != 0:
            report["error"] = f"build_page failed: {r.stderr[-300:]}"
            return json.dumps(report)
        charts = sorted(f for f in os.listdir(work) if re.fullmatch(rf"{ep}_chart\d+\.png", f))
        if page_missing or charts_missing:
            upload_public(open(os.path.join(work, f"{ep}.html"), "rb").read(),
                          f"{safe_show}/{ep}.html", content_type="text/html; charset=utf-8",
                          cache_control="public, max-age=300")
            for c in charts:
                upload_public(open(os.path.join(work, c), "rb").read(), f"{safe_show}/{c}",
                              content_type="image/png", cache_control="public, max-age=86400")
            log(f"page published with {len(charts)} chart(s): {page_url}")
        if email_wanted:
            chart_urls = ",".join(f"https://storage.googleapis.com/{GCS_BUCKET}/{safe_show}/{c}" for c in charts)
            cmd = ["python3", "build_email.py", "--notes", notes, "--number", str(n),
                   "--title", title, "--dek", dek, "--audio-url", audio_url,
                   "--page-url", page_url, "--duration", str(duration), "--date", date,
                   "--charts", chart_urls, "--drill-file", "conversions.md",
                   "--out", "email.html", "--text-out", "email.txt"]
            r = subprocess.run(cmd, cwd=work, capture_output=True, text=True, timeout=300)
            if r.returncode != 0:
                log(f"build_email failed: {r.stderr[-200:]}")
            else:
                html = open(os.path.join(work, "email.html"), encoding="utf-8").read()
                text = open(os.path.join(work, "email.txt"), encoding="utf-8").read()
                payload = {"from": EMAIL_FROM, "to": [EMAIL_TO],
                           "subject": f"{SHOW_NAME} — Ep {n}: {title}",
                           "html": html, "text": text,
                           "attachments": [{"path": audio_url,
                                            "filename": f"{SHOW_NAME.replace(' ', '_')}_{ep.capitalize()}.mp3"}]}
                resp = httpx.post("https://api.resend.com/emails", timeout=120,
                                  headers={"Authorization": f"Bearer {RESEND_API_KEY}",
                                           "Content-Type": "application/json"}, json=payload)
                if resp.status_code < 300:
                    meta = read_meta(safe_show, ep) or {}
                    meta["email_sent"] = resp.json().get("id", True)
                    write_meta(safe_show, ep, meta)
                    log(f"e-mail sent ({len(html)} bytes)")
                else:
                    log(f"e-mail failed: resend {resp.status_code} {resp.text[:120]}")
    except Exception as exc:  # noqa: BLE001
        report["error"] = f"finalize failed: {str(exc)[:300]}"
    finally:
        import shutil
        shutil.rmtree(work, ignore_errors=True)
    report["page_url"] = page_url
    report["audio_url"] = audio_url
    return json.dumps(report)


@mcp.tool()
async def finalize_episode(package: str = "@sdelsad/commodity-desk-daily",
                           version: str = "latest",
                           show: str = "commodity-desk-daily",
                           send_email: bool = True) -> str:
    """Finish the newest episode of a package from the server side, idempotently.

    Reads the package's feed.xml (latest version by default), takes its first
    item and makes sure everything that item promises actually exists: the mp3
    (generated from the script in the package if the bucket has none, and the
    feed corrected to the real file), the episode page and chart images (built
    by the toolkit's build_page.py from epNN.md), and the e-mail (build_email.py
    + Resend, once). Steps already done are skipped, so it is safe to call
    after every run and from a scheduler.

    The audio step can take 10-20 minutes and the call may time out
    client-side; the work continues. Call again later: it reports
    "already in progress" while generating and finishes the rest when done.
    """
    import anyio
    return await anyio.to_thread.run_sync(_finalize_impl, package, version, show, send_email)


@mcp.tool()
async def check_audio_ready(show: str, episode: str) -> str:
    """Check whether a previously requested episode audio is ready.

    Use this when a generate_podcast_audio call timed out client-side: the
    server keeps working and uploads the result anyway. Poll this (it returns
    instantly) until status is "done", then use the returned url,
    duration_seconds and size_bytes exactly as if the original call had
    returned them. Statuses: "working", "done", "error" (with message),
    "unknown" (no such job).
    """
    import anyio

    return await anyio.to_thread.run_sync(_check_impl, show, episode)


# ------------------------------------------------------------------- asgi ----

_inner = mcp.streamable_http_app()


async def app(scope, receive, send):
    """Health endpoint + optional secret path gate around the MCP app."""
    if scope["type"] == "http":
        path = scope.get("path", "")
        if path.startswith("/finalize") and SECRET_PATH and \
                scope.get("query_string", b"").decode() == f"key={SECRET_PATH}":
            import anyio
            body = await anyio.to_thread.run_sync(
                _finalize_impl, "@sdelsad/commodity-desk-daily", "latest",
                "commodity-desk-daily", True)
            await send({"type": "http.response.start", "status": 200,
                        "headers": [(b"content-type", b"application/json")]})
            await send({"type": "http.response.body", "body": body.encode()})
            return
        if path in ("/", "/healthz"):
            await send({"type": "http.response.start", "status": 200,
                        "headers": [(b"content-type", b"text/plain")]})
            await send({"type": "http.response.body", "body": b"ok"})
            return
        if SECRET_PATH:
            prefix = f"/{SECRET_PATH}"
            if not path.startswith(prefix + "/") and path != prefix:
                await send({"type": "http.response.start", "status": 404,
                            "headers": [(b"content-type", b"text/plain")]})
                await send({"type": "http.response.body", "body": b"not found"})
                return
            scope = dict(scope)
            scope["path"] = path[len(prefix):] or "/"
            raw = scope.get("raw_path")
            if raw:
                scope["raw_path"] = raw[len(prefix.encode()):] or b"/"
    await _inner(scope, receive, send)
