"""
GOLS - AI Avatar Video Generation Pipeline
==========================================
Runs as a FastAPI BackgroundTask (no Celery needed).

Pipeline stages (executed in order):
  1. TTS           - script -> voice.mp3   (OpenAI TTS -> gTTS -> pyttsx3 -> silence)
  2. Base / Liveness - image|video + audio -> base.mp4
                       • image + GPU models OFF -> FFmpeg "liveness" (Ken-Burns
                         motion + subtle head sway) so the avatar looks alive
                       • video source -> looped/clipped to audio length
  3. Face animation (LivePortrait / SadTalker)  - run if model installed
  4. Lip sync (Wav2Lip)                          - run if model installed
  5. Face enhancement (CodeFormer / GFPGAN)      - run if model installed
  6. Watermark (logo overlay, opacity/position)  - if enabled in settings
  7. Thumbnail + save to outputs/ + record in videos.json

GPU-only stages (3-5) are *actually invoked* now (previously TODO stubs). Each
runner is guarded: if the model directory / entrypoint is missing or the run
fails, the pipeline logs it and falls back to the previous stage's output so a
video is always produced. On a machine with the models installed and a CUDA GPU,
they execute and chain their outputs.
"""

import os
import shutil
import subprocess
import sys
import traceback
import uuid
from datetime import datetime, timezone
from typing import Optional

from app.config import cfg
from app.repository.store import projects_store, avatars_store, videos_store, settings_store
from app.utils.logger import log
from app.services import openai_service as ai


# ── FFmpeg / ffprobe discovery (cached) ───────────────────────────────────────
_FFMPEG: Optional[str] = None
_FFPROBE: Optional[str] = None


def _find_ffmpeg() -> str:
    global _FFMPEG
    if _FFMPEG:
        return _FFMPEG
    env_path = (cfg.FFMPEG_PATH or "").strip()
    if env_path and env_path != "ffmpeg":
        if os.path.isfile(env_path):
            _FFMPEG = env_path
            return _FFMPEG
        exe = env_path if env_path.endswith(".exe") else env_path + ".exe"
        if os.path.isfile(exe):
            _FFMPEG = exe
            return _FFMPEG
    found = shutil.which("ffmpeg")
    if found:
        _FFMPEG = found
        return _FFMPEG
    win_paths = [
        r"C:\ffmpeg\bin\ffmpeg.exe",
        r"C:\Program Files\ffmpeg\bin\ffmpeg.exe",
        r"C:\Program Files (x86)\ffmpeg\bin\ffmpeg.exe",
        r"C:\tools\ffmpeg\bin\ffmpeg.exe",
        os.path.join(os.environ.get("USERPROFILE", ""), "ffmpeg", "bin", "ffmpeg.exe"),
        os.path.join(os.environ.get("LOCALAPPDATA", ""), "ffmpeg", "bin", "ffmpeg.exe"),
    ]
    for p in win_paths:
        if os.path.isfile(p):
            _FFMPEG = p
            log.info("[pipeline] Found ffmpeg at: %s", p)
            return _FFMPEG
    raise RuntimeError(
        "ffmpeg not found!\n\n"
        "FIX (Windows):\n"
        "  1. Download from https://www.gyan.dev/ffmpeg/builds/\n"
        "  2. Extract to C:\\ffmpeg\\\n"
        "  3. Add C:\\ffmpeg\\bin to PATH, OR set FFMPEG_PATH in backend/.env\n"
        "FIX (macOS):  brew install ffmpeg\n"
        "FIX (Linux):  sudo apt install ffmpeg\n"
        "Then restart the app."
    )


def _find_ffprobe() -> str:
    global _FFPROBE
    if _FFPROBE:
        return _FFPROBE
    ffmpeg_path = _find_ffmpeg()
    ffprobe = os.path.join(os.path.dirname(ffmpeg_path), "ffprobe")
    if sys.platform == "win32":
        ffprobe += ".exe"
    if os.path.isfile(ffprobe):
        _FFPROBE = ffprobe
    else:
        _FFPROBE = shutil.which("ffprobe") or "ffprobe"
    return _FFPROBE


# ── Small helpers ─────────────────────────────────────────────────────────────
def _now() -> str:
    return datetime.now(timezone.utc).isoformat()


def _set(project_id: str, **fields):
    projects_store.update(project_id, fields)


def _err(project_id: str, stage: str, exc: Exception):
    msg = f"{stage}: {exc}"
    log.error("[pipeline] %s FAILED - %s", project_id, msg)
    _set(project_id, status="failed", error=msg, stage=stage, updated_at=_now())


def _run(cmd: list, label: str, cwd: Optional[str] = None,
         timeout: Optional[int] = None) -> subprocess.CompletedProcess:
    """Run a subprocess; raise RuntimeError with a clean message on failure."""
    log.debug("[pipeline] %s: %s", label, " ".join(str(c) for c in cmd))
    result = subprocess.run(
        cmd, capture_output=True, text=True, encoding="utf-8", errors="replace",
        cwd=cwd, timeout=timeout,
    )
    if result.returncode != 0:
        raise RuntimeError(f"{label} failed (exit {result.returncode}):\n{result.stderr[-700:]}")
    return result


def _audio_duration(audio_path: str) -> float:
    ffprobe = _find_ffprobe()
    try:
        probe = subprocess.run(
            [ffprobe, "-v", "error", "-show_entries", "format=duration",
             "-of", "default=noprint_wrappers=1:nokey=1", audio_path],
            capture_output=True, text=True, encoding="utf-8", errors="replace",
        )
        return float(probe.stdout.strip())
    except Exception:
        return 10.0


RES_MAP = {"480p": "854:480", "720p": "1280:720", "1080p": "1920:1080", "4k": "3840:2160"}


def _is_image(path: str) -> bool:
    return os.path.splitext(path)[1].lower() in (".jpg", ".jpeg", ".png", ".webp", ".bmp")


def _is_video(path: str) -> bool:
    return os.path.splitext(path)[1].lower() in (".mp4", ".avi", ".mov", ".mkv", ".webm")


# ══════════════════════════════════════════════════════════════════════════════
# STAGE 1 — TTS  (OpenAI gendered/multilingual -> gTTS -> pyttsx3 -> silence)
# ══════════════════════════════════════════════════════════════════════════════
def stage_tts(script: str, lang: str, out_path: str, *,
              gender: str = "female", emotion: str = "natural") -> str:
    """Generate voice audio. Returns the engine used: openai|gtts|pyttsx3|silence."""
    # --- OpenAI TTS (primary; gendered + multilingual) ---
    if ai.is_enabled():
        try:
            spoken = script
            try:
                spoken = ai.optimize_for_speech(script, language=lang, emotion=emotion)
            except Exception as e:
                log.warning("[pipeline] speech optimisation skipped (%s)", e)
            if ai.synthesize_speech(spoken, out_path, gender=gender,
                                    language=lang, emotion=emotion):
                log.info("[pipeline] TTS done via OpenAI (%s/%s) -> %s", gender, lang, out_path)
                return "openai"
        except Exception as e:
            log.warning("[pipeline] OpenAI TTS failed (%s), falling back to gTTS", e)
    else:
        log.warning("[pipeline] OpenAI TTS unavailable (no key or 'openai' package not "
                    "installed). Falling back to gTTS — run: pip install openai httpx")

    # --- gTTS (online) ---
    try:
        from gtts import gTTS
        supported = ("en", "hi", "mr", "ta", "te", "kn", "ml", "gu", "bn", "pa",
                     "ur", "es", "fr", "de", "ja", "zh", "ar", "pt", "ru", "it")
        tts = gTTS(text=script, lang=lang if lang in supported else "en")
        tts.save(out_path)
        log.info("[pipeline] TTS done via gTTS (%s) -> %s", lang, out_path)
        return "gtts"
    except Exception as e:
        log.warning("[pipeline] gTTS failed (%s), trying macOS 'say'...", e)

    # --- macOS 'say' (offline, always available on Mac, gendered) ---
    if sys.platform == "darwin" and shutil.which("say"):
        try:
            # Built-in macOS voices. Female/male defaults; many language voices exist.
            mac_voice = {
                "en": ("Samantha", "Alex"), "hi": ("Lekha", "Lekha"),
                "fr": ("Amelie", "Thomas"), "de": ("Anna", "Markus"),
                "es": ("Monica", "Jorge"), "it": ("Alice", "Luca"),
                "ja": ("Kyoko", "Otoya"), "zh": ("Tingting", "Tingting"),
                "pt": ("Luciana", "Felipe"), "ru": ("Milena", "Yuri"),
            }
            fem, male = mac_voice.get((lang or "en").lower(), ("Samantha", "Alex"))
            v = fem if (gender or "").lower() == "female" else male
            aiff = out_path.replace(".mp3", "_say.aiff")
            _run(["say", "-v", v, "-o", aiff, "--", script], "macos-say")
            _run([_find_ffmpeg(), "-y", "-i", aiff, "-b:a", "160k", out_path], "say->mp3")
            if os.path.exists(aiff):
                os.remove(aiff)
            log.info("[pipeline] TTS done via macOS say (%s) -> %s", v, out_path)
            return "say"
        except Exception as e:
            log.warning("[pipeline] macOS say failed (%s), trying pyttsx3...", e)

    # --- pyttsx3 (offline) ---
    try:
        import pyttsx3
        engine = pyttsx3.init()
        engine.setProperty("rate", 160)
        try:
            voices = engine.getProperty("voices")
            want = "female" if (gender or "").lower() == "female" else "male"
            for v in voices:
                meta = f"{getattr(v,'name','')} {getattr(v,'id','')}".lower()
                if want in meta or (want == "female" and "zira" in meta) or \
                   (want == "male" and "david" in meta):
                    engine.setProperty("voice", v.id)
                    break
        except Exception:
            pass
        wav_tmp = out_path.replace(".mp3", "_tmp.wav")
        engine.save_to_file(script, wav_tmp)
        engine.runAndWait()
        _run([_find_ffmpeg(), "-y", "-i", wav_tmp, out_path], "pyttsx3->mp3")
        if os.path.exists(wav_tmp):
            os.remove(wav_tmp)
        log.info("[pipeline] TTS done via pyttsx3 -> %s", out_path)
        return "pyttsx3"
    except Exception as e:
        log.warning("[pipeline] pyttsx3 failed (%s), using silence fallback...", e)

    # --- Silence fallback ---
    duration = max(5, min(len(script) // 14, 120))
    _run([_find_ffmpeg(), "-y", "-f", "lavfi", "-i", "anullsrc=r=44100:cl=mono",
          "-t", str(duration), out_path], "silence-fallback")
    log.warning("[pipeline] NO VOICE: all TTS engines failed — produced SILENT audio (%ds). "
                "Install openai (pip install openai httpx) and set OPENAI_API_KEY, or ensure "
                "internet access for gTTS.", duration)
    return "silence"


# ══════════════════════════════════════════════════════════════════════════════
# STAGE 2 — Base video / Liveness
# ══════════════════════════════════════════════════════════════════════════════
def stage_make_video(source_path, audio_path: str, out_path: str,
                     resolution: str, *, liveness: bool = True) -> float:
    """
    Build the base MP4 from an image or video + audio. Returns duration (s).
    For image sources with `liveness` enabled, applies a Ken-Burns zoom plus a
    gentle sinusoidal pan ("head sway") so the avatar is not a frozen frame even
    without GPU models.
    """
    vf_scale = RES_MAP.get(resolution, "1920:1080")
    w, h = vf_scale.split(":")
    ffmpeg = _find_ffmpeg()
    duration = _audio_duration(audio_path)
    fps = 25

    pad = (f"scale={vf_scale}:force_original_aspect_ratio=decrease,"
           f"pad={w}:{h}:(ow-iw)/2:(oh-ih)/2:black,format=yuv420p")

    if source_path and os.path.exists(source_path) and _is_video(source_path):
        cmd = [
            ffmpeg, "-y", "-stream_loop", "-1", "-i", source_path, "-i", audio_path,
            "-t", str(duration), "-vf", pad,
            "-c:v", "libx264", "-preset", "fast", "-crf", "23",
            "-c:a", "aac", "-b:a", "160k", "-shortest", "-movflags", "+faststart", out_path,
        ]
    elif source_path and os.path.exists(source_path) and _is_image(source_path):
        if liveness:
            # Show the WHOLE photo (no aggressive zoom/crop): fit it inside the
            # frame and fill the empty space with a soft blurred copy of itself,
            # then add only a very gentle Ken-Burns drift so it feels alive while
            # the subject stays fully framed and natural.
            zoom = "min(1.0+0.0006*on,1.05)"          # slow zoom, max +5%
            x = f"iw/2-(iw/zoom/2)+sin(on/{fps*2.0})*(iw*0.006)"   # tiny sway
            y = f"ih/2-(ih/zoom/2)+cos(on/{fps*2.6})*(ih*0.005)"
            fc = (
                f"[0:v]scale={w}:{h}:force_original_aspect_ratio=increase,"
                f"crop={w}:{h},boxblur=26:3,setsar=1[bg];"
                f"[0:v]scale={w}:{h}:force_original_aspect_ratio=decrease[fg];"
                f"[bg][fg]overlay=(W-w)/2:(H-h)/2,"
                f"zoompan=z='{zoom}':x='{x}':y='{y}':d=1:s={w}x{h}:fps={fps},"
                f"format=yuv420p[v]"
            )
            cmd = [
                ffmpeg, "-y", "-loop", "1", "-i", source_path, "-i", audio_path,
                "-t", str(duration), "-r", str(fps),
                "-filter_complex", fc, "-map", "[v]", "-map", "1:a",
                "-c:v", "libx264", "-preset", "fast", "-crf", "20",
                "-c:a", "aac", "-b:a", "160k", "-shortest",
                "-movflags", "+faststart", out_path,
            ]
        else:
            cmd = [
                ffmpeg, "-y", "-loop", "1", "-i", source_path, "-i", audio_path,
                "-t", str(duration), "-vf", pad,
                "-c:v", "libx264", "-preset", "fast", "-crf", "23",
                "-c:a", "aac", "-b:a", "160k", "-movflags", "+faststart", out_path,
            ]
    else:
        cmd = [
            ffmpeg, "-y", "-f", "lavfi", "-i", f"color=c=black:s={w}x{h}:r={fps}",
            "-i", audio_path, "-t", str(duration), "-vf", "format=yuv420p",
            "-c:v", "libx264", "-preset", "fast", "-crf", "23",
            "-c:a", "aac", "-b:a", "160k", "-movflags", "+faststart", out_path,
        ]

    _run(cmd, "make-video")
    log.info("[pipeline] Base video created (%.1fs, liveness=%s) -> %s",
             duration, liveness, out_path)
    return duration


# ══════════════════════════════════════════════════════════════════════════════
# STAGES 3-5 — GPU model runners (executed when installed; graceful fallback)
# ══════════════════════════════════════════════════════════════════════════════
def _entrypoint(model_dir: str, names) -> Optional[str]:
    for n in names:
        p = os.path.join(model_dir, n)
        if os.path.isfile(p):
            return p
    return None


def run_liveportrait(source_image: str, work_dir: str) -> Optional[str]:
    """
    LivePortrait / SadTalker: animate a still portrait into a moving base clip.
    Requires a driving template video at <LIVEPORTRAIT_PATH>/assets/driving.mp4.
    Returns animated (silent) mp4 path or None to skip.
    """
    mdir = cfg.LIVEPORTRAIT_PATH
    if not (mdir and os.path.isdir(mdir) and source_image):
        return None
    script = _entrypoint(mdir, ["inference.py", "demo.py", "run.py"])
    driving = os.path.join(mdir, "assets", "driving.mp4")
    if not script or not os.path.isfile(driving):
        log.info("[pipeline] LivePortrait present but no entrypoint/driving asset — skipping")
        return None
    out = os.path.join(work_dir, "animated.mp4")
    try:
        _run([cfg.MODELS_PYTHON, script, "-s", source_image, "-d", driving,
              "-o", out], "liveportrait", cwd=mdir, timeout=1800)
        return out if os.path.isfile(out) else None
    except Exception as e:
        log.warning("[pipeline] LivePortrait failed (%s) — falling back", e)
        return None


def _patch_wav2lip_source(mdir: str) -> None:
    """
    Make the (2019-era) Wav2Lip source run on modern librosa / numpy WITHOUT any
    manual edits. Idempotent: safe to run before every invocation.

      • librosa.filters.mel(sr, n_fft, ...)  -> keyword args  (librosa >= 0.10)
      • deprecated numpy aliases np.float/np.int/np.bool/np.object/np.complex
        -> python builtins (numpy >= 1.24)
    """
    import re
    # 1. audio.py mel() positional -> keyword
    audio_py = os.path.join(mdir, "audio.py")
    try:
        if os.path.isfile(audio_py):
            src = open(audio_py, encoding="utf-8").read()
            new = re.sub(
                r"librosa\.filters\.mel\(\s*hp\.sample_rate\s*,\s*hp\.n_fft\s*,",
                "librosa.filters.mel(sr=hp.sample_rate, n_fft=hp.n_fft,",
                src,
            )
            if new != src:
                open(audio_py, "w", encoding="utf-8").write(new)
                log.info("[pipeline] patched Wav2Lip audio.py mel() for modern librosa")
    except Exception as e:
        log.warning("[pipeline] could not patch audio.py (%s)", e)

    # 2. numpy deprecated aliases across all .py files
    alias = re.compile(r"\bnp\.(float|int|bool|object|complex)\b(?!\d|_|\s*\()")
    repl = {"float": "float", "int": "int", "bool": "bool",
            "object": "object", "complex": "complex"}
    for root, _d, files in os.walk(mdir):
        if ".git" in root:
            continue
        for f in files:
            if not f.endswith(".py"):
                continue
            p = os.path.join(root, f)
            try:
                s = open(p, encoding="utf-8").read()
                n = alias.sub(lambda m: repl[m.group(1)], s)
                if n != s:
                    open(p, "w", encoding="utf-8").write(n)
                    log.info("[pipeline] patched numpy aliases in %s", p)
            except Exception:
                pass


def run_wav2lip(face_video_or_image: str, audio_path: str, work_dir: str) -> Optional[str]:
    """Wav2Lip: apply accurate lip-sync of `audio` onto `face` video/image."""
    mdir = cfg.WAV2LIP_PATH
    if not (mdir and os.path.isdir(mdir)):
        return None
    script = _entrypoint(mdir, ["inference.py", "wav2lip_inference.py"])
    ckpt = _entrypoint(os.path.join(mdir, "checkpoints"),
                       ["wav2lip_gan.pth", "wav2lip.pth"]) or \
           _entrypoint(mdir, ["wav2lip_gan.pth", "wav2lip.pth"])
    if not script or not ckpt:
        log.info("[pipeline] Wav2Lip present but missing entrypoint/checkpoint — skipping")
        return None
    if cfg.WAV2LIP_AUTOPATCH:
        _patch_wav2lip_source(mdir)
    out = os.path.join(work_dir, "lipsync.mp4")
    try:
        # --nosmooth keeps a still portrait crisp; --resize_factor 1 = full res.
        # Wav2Lip muxes the provided --audio into its output, so the result has voice.
        _run([cfg.MODELS_PYTHON, script, "--checkpoint_path", ckpt,
              "--face", face_video_or_image, "--audio", audio_path,
              "--outfile", out, "--nosmooth", "--resize_factor", "1",
              "--pads", "0", "10", "0", "0"], "wav2lip", cwd=mdir, timeout=3600)
        return out if os.path.isfile(out) else None
    except Exception as e:
        log.warning("[pipeline] Wav2Lip failed (%s) — falling back", e)
        return None


def run_musetalk(face_video_or_image: str, audio_path: str, work_dir: str) -> Optional[str]:
    """
    MuseTalk lip-sync (higher quality than Wav2Lip). Requires an NVIDIA CUDA GPU
    and the MuseTalk repo + weights at MUSETALK_PATH. Skips gracefully otherwise
    (e.g. on macOS, where CUDA is unavailable).
    """
    mdir = cfg.MUSETALK_PATH
    if not (mdir and os.path.isdir(mdir)):
        return None
    # MuseTalk needs CUDA; bail early on platforms without it.
    try:
        import torch  # noqa
        has_cuda = bool(getattr(torch, "cuda", None) and torch.cuda.is_available())
    except Exception:
        has_cuda = False
    if not has_cuda:
        log.info("[pipeline] MuseTalk present but no CUDA GPU — skipping (use Wav2Lip)")
        return None
    script = _entrypoint(mdir, ["inference.py"]) or \
        _entrypoint(os.path.join(mdir, "scripts"), ["inference.py"])
    if not script:
        log.info("[pipeline] MuseTalk present but no entrypoint — skipping")
        return None
    out = os.path.join(work_dir, "musetalk.mp4")
    try:
        _run([cfg.MODELS_PYTHON, script, "--video_path", face_video_or_image,
              "--audio_path", audio_path, "--result_dir", work_dir],
             "musetalk", cwd=mdir, timeout=3600)
        if os.path.isfile(out):
            return out
        for r, _d, files in os.walk(work_dir):
            for f in files:
                if f.lower().endswith(".mp4") and "muse" in f.lower():
                    return os.path.join(r, f)
        return None
    except Exception as e:
        log.warning("[pipeline] MuseTalk failed (%s) — falling back", e)
        return None


def run_codeformer(in_video: str, work_dir: str) -> Optional[str]:
    """CodeFormer / GFPGAN: per-frame face restoration & enhancement."""
    mdir = cfg.CODEFORMER_PATH
    if not (mdir and os.path.isdir(mdir)):
        return None
    script = _entrypoint(mdir, ["inference_codeformer.py", "inference.py"])
    if not script:
        log.info("[pipeline] CodeFormer present but missing entrypoint — skipping")
        return None
    out_dir = os.path.join(work_dir, "cf_out")
    os.makedirs(out_dir, exist_ok=True)
    try:
        _run([cfg.MODELS_PYTHON, script, "-i", in_video, "-o", out_dir,
              "--bg_upsampler", "None", "-w", "0.7"], "codeformer",
             cwd=mdir, timeout=1800)
        for root, _d, files in os.walk(out_dir):
            for f in files:
                if f.lower().endswith(".mp4"):
                    return os.path.join(root, f)
        return None
    except Exception as e:
        log.warning("[pipeline] CodeFormer failed (%s) — falling back", e)
        return None


# ══════════════════════════════════════════════════════════════════════════════
# STAGE 6 — Professional watermark (logo overlay)
# ══════════════════════════════════════════════════════════════════════════════
_POSITIONS = {
    "top-left":     "20:20",
    "top-right":    "W-w-20:20",
    "bottom-left":  "20:H-h-20",
    "bottom-right": "W-w-20:H-h-20",
    "center":       "(W-w)/2:(H-h)/2",
}


# drawtext x/y expressions per position (tw/th = text size, w/h = video size)
_TEXT_POS = {
    "top-left":     "x=m:y=m",
    "top-right":    "x=w-tw-m:y=m",
    "bottom-left":  "x=m:y=h-th-m",
    "bottom-right": "x=w-tw-m:y=h-th-m",
    "center":       "x=(w-tw)/2:y=(h-th)/2",
}


def stage_watermark(in_path: str, out_path: str, *,
                    opacity: float = 0.45, position: str = "bottom-right",
                    scale: float = 0.16, text: str = None) -> None:
    """
    Burn a clean semi-transparent TEXT watermark (HeyGen-style), e.g. "GOLS",
    into every video. Font size scales with the video height; a soft shadow keeps
    it readable on any background. Re-encodes video, copies audio.
    """
    ffmpeg = _find_ffmpeg()
    op = max(0.05, min(1.0, float(opacity)))
    wm_text = (text or cfg.WATERMARK_TEXT or "GOLS").strip()
    # Sanitise for drawtext (avoid breaking the filter syntax)
    safe = wm_text.replace("\\", "").replace(":", "").replace("'", "")
    # margin + font size relative to height
    margin = "h*0.03"
    fontsize = "h*0.045"
    pos = _TEXT_POS.get(position, _TEXT_POS["bottom-right"]).replace("m", margin)

    drawtext = (
        f"drawtext=text='{safe}':"
        f"fontsize={fontsize}:fontcolor=white@{op}:"
        f"shadowcolor=black@{min(1.0, op + 0.15):.2f}:shadowx=2:shadowy=2:"
        f"box=0:{pos}"
    )
    cmd = [ffmpeg, "-y", "-i", in_path, "-vf", f"{drawtext},format=yuv420p",
           "-c:a", "copy", "-c:v", "libx264", "-preset", "fast", "-crf", "20",
           "-movflags", "+faststart", out_path]
    _run(cmd, "watermark-text")
    log.info("[pipeline] Text watermark '%s' applied (op=%.2f, %s) -> %s",
             safe, op, position, out_path)


# ══════════════════════════════════════════════════════════════════════════════
# STAGE 7 — Thumbnail
# ══════════════════════════════════════════════════════════════════════════════
def stage_thumbnail(video_path: str, out_path: str) -> bool:
    ffmpeg = _find_ffmpeg()
    try:
        _run([ffmpeg, "-y", "-ss", "1", "-i", video_path, "-frames:v", "1",
              "-vf", "scale=640:-2", "-q:v", "3", out_path], "thumbnail")
        return os.path.isfile(out_path)
    except Exception as e:
        log.warning("[pipeline] thumbnail generation failed (%s)", e)
        return False


# ══════════════════════════════════════════════════════════════════════════════
# MAIN PIPELINE
# ══════════════════════════════════════════════════════════════════════════════
def run_pipeline(project_id: str) -> None:
    log.info("[pipeline] START %s", project_id)
    work_dir = None
    try:
        ffmpeg_exe = _find_ffmpeg()
        log.info("[pipeline] Using ffmpeg: %s", ffmpeg_exe)

        project = projects_store.get(project_id)
        if not project:
            log.error("[pipeline] Project %s not found", project_id)
            return

        avatar = avatars_store.get(project.get("avatar_id", "")) or {}
        script = project.get("script", "Hello, this is your AI avatar.")
        lang = project.get("language", "en")
        res = project.get("resolution", "1080p")
        gender = project.get("voice_gender") or avatar.get("voice_gender") or "female"
        emotion = project.get("emotion", "natural")

        slist = settings_store.list(per_page=1)["items"]
        settings = slist[0] if slist else {}
        use_wm = settings.get("watermark", False)
        wm_opacity = settings.get("watermark_opacity", cfg.WATERMARK_OPACITY)
        wm_position = settings.get("watermark_position", cfg.WATERMARK_POSITION)

        job_id = str(uuid.uuid4())[:8]
        work_dir = os.path.join(cfg.OUTPUT_DIR, f"_job_{job_id}")
        os.makedirs(work_dir, exist_ok=True)

        audio_path = os.path.join(work_dir, "voice.mp3")
        base_mp4 = os.path.join(work_dir, "base.mp4")
        wm_mp4 = os.path.join(work_dir, "watermarked.mp4")

        source_url = avatar.get("source_url")
        source_path = None
        if source_url:
            candidate = os.path.join(cfg.UPLOAD_DIR, os.path.basename(source_url))
            if os.path.exists(candidate):
                source_path = candidate
            else:
                log.warning("[pipeline] Source file not found: %s", candidate)

        # ── STAGE 1: TTS ──────────────────────────────────────────────────────
        _set(project_id, status="processing", stage="tts",
             stage_label="Generating voice...", progress=10, error=None)
        tts_engine = stage_tts(script, lang, audio_path, gender=gender, emotion=emotion)
        voice_warning = None
        if tts_engine == "silence":
            voice_warning = ("No voice was generated (silent audio). Install OpenAI "
                             "(pip install openai httpx) with OPENAI_API_KEY set, or "
                             "ensure internet access for gTTS, then regenerate.")
        _set(project_id, stage="tts_done", progress=30, voice_engine_used=tts_engine,
             voice_warning=voice_warning)

        # ── STAGE 2: Base video / liveness ────────────────────────────────────
        _set(project_id, stage="video", stage_label="Creating base video...", progress=40)
        duration = stage_make_video(source_path, audio_path, base_mp4, res,
                                    liveness=cfg.ENABLE_LIVENESS)
        _set(project_id, stage="video_done", progress=58)
        current = base_mp4

        # ── STAGE 3: Face animation (LivePortrait/SadTalker) ──────────────────
        if source_path and _is_image(source_path):
            _set(project_id, stage="animation",
                 stage_label="Face animation (LivePortrait)...", progress=64)
            animated = run_liveportrait(source_path, work_dir)
            if animated:
                muxed = os.path.join(work_dir, "animated_av.mp4")
                try:
                    _run([ffmpeg_exe, "-y", "-i", animated, "-i", audio_path,
                          "-c:v", "copy", "-c:a", "aac", "-shortest",
                          "-movflags", "+faststart", muxed], "mux-animated")
                    current = muxed
                except Exception:
                    current = animated
                log.info("[pipeline] LivePortrait animation applied")

        # ── STAGE 4: Lip sync ─────────────────────────────────────────────────
        # Engine selection (per-avatar): "did" (hosted, realistic) | "musetalk"
        # (GPU) | "wav2lip" (CPU/Mac). Hosted D-ID animates directly from the
        # still photo + our voice; local engines fall back gracefully.
        lip_engine = (avatar.get("lip_engine") or "wav2lip").lower()
        lip_face = current if current != base_mp4 else (source_path or base_mp4)
        _set(project_id, stage="lipsync",
             stage_label=f"Lip sync ({lip_engine})...", progress=74)
        lipsync = None

        if lip_engine == "did":
            from app.services import did_service
            did_out = os.path.join(work_dir, "did.mp4")
            if source_path and _is_image(source_path):
                ok, detail = did_service.generate_talk(
                    source_path, audio_path, did_out,
                    script_text=script, gender=gender, language=lang)
                if ok:
                    lipsync = did_out
                else:
                    # Surface the real reason to the UI instead of silently degrading.
                    log.warning("[pipeline] D-ID failed: %s", detail)
                    _set(project_id, did_error=f"D-ID: {detail}")
                    lipsync = run_wav2lip(lip_face, audio_path, work_dir)
            else:
                _set(project_id, did_error="D-ID needs an image source (not a video).")
                lipsync = run_wav2lip(lip_face, audio_path, work_dir)
        elif lip_engine == "musetalk":
            lipsync = run_musetalk(lip_face, audio_path, work_dir) or \
                      run_wav2lip(lip_face, audio_path, work_dir)
        else:
            lipsync = run_wav2lip(lip_face, audio_path, work_dir)

        if lipsync:
            current = lipsync
            log.info("[pipeline] Lip-sync applied via %s", lip_engine)
        else:
            log.info("[pipeline] No lip-sync engine produced output — using living portrait")

        # ── STAGE 5: Face enhancement (CodeFormer) ────────────────────────────
        if avatar.get("enhance", True):
            _set(project_id, stage="enhance",
                 stage_label="Face enhancement (CodeFormer)...", progress=84)
            enhanced = run_codeformer(current, work_dir)
            if enhanced:
                cf_av = os.path.join(work_dir, "enhanced_av.mp4")
                try:
                    _run([ffmpeg_exe, "-y", "-i", enhanced, "-i", audio_path,
                          "-c:v", "copy", "-c:a", "aac", "-shortest",
                          "-movflags", "+faststart", cf_av], "mux-enhanced")
                    current = cf_av
                except Exception:
                    current = enhanced
                log.info("[pipeline] CodeFormer enhancement applied")

        # ── STAGE 6: Watermark (non-fatal — never lose the video over a WM) ────
        _set(project_id, stage="watermark", stage_label="Finalising video...", progress=92)
        final_source = current
        if use_wm:
            try:
                stage_watermark(current, wm_mp4, opacity=wm_opacity,
                                position=wm_position)
                if os.path.isfile(wm_mp4) and os.path.getsize(wm_mp4) > 1000:
                    final_source = wm_mp4
                else:
                    log.warning("[pipeline] watermark produced no file — using un-watermarked")
            except Exception as e:
                log.warning("[pipeline] watermark failed (%s) — using un-watermarked video", e)

        # ── STAGE 7: Save + thumbnail ─────────────────────────────────────────
        filename = f"avatar_{job_id}.mp4"
        final_out = os.path.join(cfg.OUTPUT_DIR, filename)
        shutil.copy2(final_source, final_out)

        thumb_name = f"avatar_{job_id}.jpg"
        thumb_path = os.path.join(cfg.OUTPUT_DIR, thumb_name)
        has_thumb = stage_thumbnail(final_out, thumb_path)

        file_size_mb = round(os.path.getsize(final_out) / 1_048_576, 2)
        log.info("[pipeline] Final video saved -> %s (%.2f MB)", final_out, file_size_mb)

        video = videos_store.create({
            "title": project.get("title", "Avatar Video"),
            "project_id": project_id,
            "avatar_id": project.get("avatar_id"),
            "filename": filename,
            "thumbnail": thumb_name if has_thumb else None,
            "status": "completed",
            "duration_sec": round(duration, 1),
            "resolution": res,
            "language": lang,
            "voice_gender": gender,
            "voice_engine_used": tts_engine,
            "voice_warning": voice_warning,
            "file_size_mb": file_size_mb,
            "watermark": use_wm,
            "download_url": "",
        })
        videos_store.update(video["id"], {
            "download_url": f"/api/videos/{video['id']}/download",
            "stream_url": f"/api/videos/{video['id']}/stream",
            "thumbnail_url": f"/api/videos/{video['id']}/thumbnail" if has_thumb else None,
        })

        _set(project_id, status="completed", stage="done", stage_label="Complete",
             progress=100, output_url=f"/api/videos/{video['id']}/download",
             stream_url=f"/api/videos/{video['id']}/stream",
             video_id=video["id"], duration_sec=round(duration, 1),
             completed_at=_now())
        log.info("[pipeline] DONE %s -> video %s", project_id, video["id"])

    except Exception as exc:
        log.error("[pipeline] EXCEPTION:\n%s", traceback.format_exc())
        _err(project_id, getattr(exc, "stage", "pipeline"), exc)
    finally:
        if work_dir and os.path.exists(work_dir):
            try:
                shutil.rmtree(work_dir)
            except Exception:
                pass
