"""
Higgsfield unofficial client.

Drives an already-open, logged-in higgsfield.ai tab in a Chrome started with
--remote-debugging-port=9222, over the Chrome DevTools Protocol.

WHY this design (transport = the real browser, not impersonation):
  - The Higgsfield API (fnf.higgsfield.ai) is behind DataDome. A TLS-impersonating
    HTTP client (curl_cffi) works briefly, but DataDome risk-scores it and starts
    returning 403 captcha challenges after sustained use. So instead we run the API
    calls AS in-page fetch() via CDP (BrowserSession): the request originates from
    the real logged-in tab, carrying genuine Chrome TLS + the JS-earned (HttpOnly)
    datadome cookie, so DataDome treats it like the app's own calls. CORS is fine
    because we borrow the higgsfield.ai origin the server already whitelists.
  - The bearer is minted inline per request via window.Clerk.session.getToken().
  - Only the presigned cloudfront PUT (image upload) and the public video download
    stay on curl_cffi — neither touches the DataDome-protected origin.

RUN (from WSL, using Windows Python so localhost:9222 is reachable natively):
    python.exe "$(wslpath -w scripts/higgsfield/higgsfield_client.py)"
or just:  bash scripts/higgsfield/run.sh

Deps (Windows Python):  curl_cffi  websocket-client
"""

import argparse
import json
import os
import sys
import time

import websocket  # websocket-client
from curl_cffi import requests as cffi_requests

# Windows consoles default to cp1252 and choke on the emoji in our logs.
for _s in (sys.stdout, sys.stderr):
    try:
        _s.reconfigure(encoding="utf-8")
    except Exception:
        pass

# --emit-json mode: machine-readable NDJSON progress on the REAL stdout, while all
# the human print()s get redirected to stderr (set up in main). The backend
# HiggsfieldProvider parses these events to drive the node's live status.
_JSON_OUT = None


def event(**fields):
    if _JSON_OUT is not None:
        _JSON_OUT.write(json.dumps(fields) + "\n")
        _JSON_OUT.flush()


class ImageBlocked(Exception):
    """Verification returned a deterministic block (e.g. nsfw) — do NOT retry."""
    def __init__(self, status):
        super().__init__(f"blocked: {status}")
        self.status = status


class VerifyTimeout(Exception):
    """ip_check didn't finish within the per-attempt window — retry is worthwhile."""

# ==========================================
# CONFIG
# ==========================================
CDP_HTTP = os.environ.get("CDP_HTTP", "http://localhost:9222")
GEN_BASE = "https://fnf.higgsfield.ai"
IMPERSONATE = "chrome131"       # match a recent real Chrome TLS fingerprint
OUTPUT_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "output")

# Defaults (overridable via CLI)
DEFAULT_PROMPT = "create a video of the character in the character sheet slightly animated movement naturally"
DEFAULT_IMAGE_ID = "d6c378a2-5db4-44ed-b3ca-c5cb3c14b968"
DEFAULT_IMAGE_URL = "https://d2ol7oe51mr4n9.cloudfront.net/user_3F4DMevtcYoMWqVsIglYT2yjLIz/d6c378a2-5db4-44ed-b3ca-c5cb3c14b968.png"


# ==========================================
# Minimal synchronous CDP client
# ==========================================
class CDP:
    def __init__(self, ws_url):
        self.ws = websocket.create_connection(ws_url, max_size=None, timeout=30)
        self._id = 0

    def send(self, method, params=None, session_id=None):
        self._id += 1
        msg = {"id": self._id, "method": method, "params": params or {}}
        if session_id:
            msg["sessionId"] = session_id
        self.ws.send(json.dumps(msg))
        return self._id

    def wait_for(self, msg_id, timeout=30):
        """Read frames until we see the response for msg_id; ignore everything else."""
        deadline = time.time() + timeout
        while time.time() < deadline:
            self.ws.settimeout(max(0.1, deadline - time.time()))
            try:
                msg = json.loads(self.ws.recv())
            except websocket.WebSocketTimeoutException:
                break
            if msg.get("id") == msg_id:
                if "error" in msg:
                    raise RuntimeError(f"CDP error for {msg_id}: {msg['error']}")
                return msg.get("result", {})
        raise TimeoutError(f"Timed out waiting for CDP response id={msg_id}")

    def call(self, method, params=None, session_id=None, timeout=30):
        return self.wait_for(self.send(method, params, session_id), timeout)

    def wait_event(self, method, timeout):
        """Block until a given CDP event arrives (or timeout). Returns True/False."""
        deadline = time.time() + timeout
        while time.time() < deadline:
            self.ws.settimeout(max(0.1, deadline - time.time()))
            try:
                msg = json.loads(self.ws.recv())
            except websocket.WebSocketTimeoutException:
                return False
            if msg.get("method") == method:
                return True
        return False

    def close(self):
        try:
            self.ws.close()
        except Exception:
            pass


def find_higgsfield_tab():
    """Return the targetId of an already-open, healthy higgsfield tab, or None.

    We do NOT spawn a fresh tab: loading higgsfield in a CDP-created tab trips
    DataDome and crashes the renderer ('Render process gone'). An existing tab
    is already past DataDome and authenticated, so we attach to it instead.
    """
    for t in cffi_requests.get(f"{CDP_HTTP}/json/list", timeout=10).json():
        if t.get("type") == "page" and "higgsfield.ai" in (t.get("url") or ""):
            return t["id"], t["url"]
    return None, None


def get_fresh_tokens():
    ver = cffi_requests.get(f"{CDP_HTTP}/json/version", timeout=10).json()
    print(f"🔌 Connected to {ver['Browser']}")

    target_id, url = find_higgsfield_tab()
    if not target_id:
        raise RuntimeError(
            "No open higgsfield.ai tab found.\n"
            "    Open https://higgsfield.ai/ai/video in the debug Chrome and log in,\n"
            "    then re-run. (A freshly-spawned tab gets killed by DataDome.)")
    print(f"📎 Attaching to existing tab: {url}")

    cdp = CDP(ver["webSocketDebuggerUrl"])
    try:
        session_id = cdp.call("Target.attachToTarget",
                              {"targetId": target_id, "flatten": True})["sessionId"]
        cdp.call("Runtime.enable", session_id=session_id)

        # Mint a fresh bearer directly from Clerk — deterministic, no UI clicks.
        mint = ("(async()=>{try{return (window.Clerk&&window.Clerk.session)"
                "?await window.Clerk.session.getToken():null;}catch(e){return null;}})()")
        res = cdp.call("Runtime.evaluate",
                       {"expression": mint, "awaitPromise": True, "returnByValue": True},
                       session_id=session_id)
        jwt = res.get("result", {}).get("value")
        if not jwt:
            raise RuntimeError("Clerk token unavailable (is this tab logged in?)")
        print("✅ Bearer minted via Clerk.")

        # Pull FULL cookie jar (incl. HttpOnly datadome) via CDP.
        cookies = cdp.call("Network.getCookies",
                           {"urls": [f"{GEN_BASE}/", "https://higgsfield.ai/"]},
                           session_id=session_id).get("cookies", [])
        seen = {c["name"]: c["value"] for c in cookies}
        cookie_header = "; ".join(f"{k}={v}" for k, v in seen.items())
        datadome = seen.get("datadome")
        print(f"🍪 {len(seen)} cookies (datadome={'yes' if datadome else 'no'})")

        return f"Bearer {jwt}", datadome, cookie_header
    finally:
        # Detach only — never close the user's real tab.
        try:
            cdp.call("Target.detachFromTarget", {"sessionId": session_id}, timeout=5)
        except Exception:
            pass
        cdp.close()


class PageResponse:
    """Minimal requests-like wrapper over an in-page fetch result."""
    def __init__(self, status, text):
        self.status_code = status
        self.text = text

    def json(self):
        return json.loads(self.text)


def _fetch_expr(method, url, body):
    """Build the JS that runs fetch() in the page context and returns {status, body}.
    Auth bearer is minted inline via Clerk; the datadome cookie + real Chrome TLS
    ride along automatically because the request originates from the real tab."""
    parts = [
        "(async()=>{try{",
        # Clerk can be briefly undefined right after a tab nav/reload — wait for it.
        "let _n=0;while(!(window.Clerk&&window.Clerk.session)&&_n++<40){await new Promise(r=>setTimeout(r,250));}",
        "const t=await window.Clerk.session.getToken();",
        f"const o={{method:{json.dumps(method)},headers:{{authorization:'Bearer '+t}},credentials:'include'}};",
    ]
    if body is not None:
        parts.append("o.headers['content-type']='application/json';")
        parts.append(f"o.body=JSON.stringify({json.dumps(body)});")
    parts += [
        f"const r=await fetch({json.dumps(url)},o);",
        "const x=await r.text();",
        "return JSON.stringify({status:r.status,body:x});",
        "}catch(e){return JSON.stringify({status:0,error:String((e&&e.message)||e)});}})()",
    ]
    return "".join(parts)


class BrowserSession:
    """Routes every higgsfield API call through the real logged-in Chrome tab via
    CDP `fetch` — so DataDome sees genuine browser requests (real TLS + the JS-earned
    datadome cookie) and never challenges them. Replaces the curl_cffi/TLS-impersonation
    + token/cookie-harvesting transport, which DataDome eventually risk-scores and 403s.

    The cloudfront presigned PUT (no DataDome) and the public video download stay on
    curl_cffi — they don't touch the protected fnf.higgsfield.ai origin."""

    def __init__(self):
        ver = cffi_requests.get(f"{CDP_HTTP}/json/version", timeout=10).json()
        print(f"🔌 Connected to {ver['Browser']}")
        tid, url = find_higgsfield_tab()
        if not tid:
            raise RuntimeError(
                "No open higgsfield.ai tab found.\n"
                "    Open https://higgsfield.ai/ai/video in the debug Chrome and log in.")
        print(f"📎 Routing API calls through tab: {url}")
        self.cdp = CDP(ver["webSocketDebuggerUrl"])
        self.session_id = self.cdp.call(
            "Target.attachToTarget", {"targetId": tid, "flatten": True})["sessionId"]
        self.cdp.call("Runtime.enable", session_id=self.session_id)

    def request(self, method, url, body=None, timeout=120, retries=3):
        # Retry transient blips (tab reload, Clerk not ready, CDP hiccup) so a
        # single bad evaluate doesn't kill a long poll loop.
        last = None
        for attempt in range(retries):
            res = self.cdp.call(
                "Runtime.evaluate",
                {"expression": _fetch_expr(method, url, body),
                 "awaitPromise": True, "returnByValue": True},
                session_id=self.session_id, timeout=timeout)
            val = res.get("result", {}).get("value")
            if val is None:
                last = f"evaluate failed: {res.get('exceptionDetails')}"
            else:
                data = json.loads(val)
                if data.get("status") == 0:
                    last = f"in-page fetch error: {data.get('error')}"
                else:
                    return PageResponse(data["status"], data["body"])
            if attempt < retries - 1:
                print(f"   ↻ transient request error, retrying… ({last})")
                time.sleep(3)
        raise RuntimeError(last)

    def close(self):
        try:
            self.cdp.call("Target.detachFromTarget", {"sessionId": self.session_id}, timeout=5)
        except Exception:
            pass
        self.cdp.close()


MIME_BY_EXT = {".jpg": "image/jpeg", ".jpeg": "image/jpeg", ".png": "image/png",
               ".webp": "image/webp"}
# A FINAL verdict in this set means the image is blocked. It is only valid once the NSFW/IP check has
# FINISHED (ip_check_finished=True) — reading `status` mid-check returns a not-yet-settled value, which is
# what falsely flagged images Higgsfield ultimately verifies as clean. See the verify loop below.
BLOCK_STATUSES = {"nsfw", "rejected", "failed", "blocked"}


VERIFY_TIMEOUT_S = 90  # per-attempt ceiling on the NSFW/IP check before retrying


def upload_image(session, path, surface="seedance_2", verify_timeout=VERIFY_TIMEOUT_S):
    """Upload a local image to Higgsfield and wait for the NSFW/IP check (one attempt).

    Reproduces the UI flow:
      1. POST /media/batch       -> reserve id + presigned upload_url
      2. PUT  <upload_url>       -> raw bytes to cloudfront (presigned, no auth)
      3. POST /media/{id}/upload -> trigger NSFW + IP checks
      4. GET  /media/{id}        -> poll until ip_check_finished
    Returns the final media object, or raises ImageBlocked / VerifyTimeout.
    """
    ext = os.path.splitext(path)[1].lower()
    mime = MIME_BY_EXT.get(ext)
    if not mime:
        raise RuntimeError(f"Unsupported image type '{ext}' (use {list(MIME_BY_EXT)})")
    with open(path, "rb") as f:
        body = f.read()
    print(f"⬆️  Uploading {os.path.basename(path)} ({len(body)} bytes, {mime})…")
    event(event="progress", phase="uploading")

    # 1. reserve
    r = session.request("POST", f"{GEN_BASE}/media/batch",
                        body={"mimetypes": [mime], "source": "user_upload",
                              "surface": surface, "force_ip_check": True})
    if r.status_code != 200:
        raise RuntimeError(f"/media/batch failed ({r.status_code}): {r.text[:500]}")
    m = r.json()[0]
    mid, url, upload_url = m["id"], m["url"], m["upload_url"]
    ctype = m.get("content_type", mime)

    # 2. presigned PUT — bare request, NOT session (no auth/cookies; only the
    #    signed Content-Type + host headers may be sent).
    put = cffi_requests.put(upload_url, data=body, headers={"content-type": ctype},
                            impersonate=IMPERSONATE)
    if put.status_code not in (200, 204):
        raise RuntimeError(f"presigned PUT failed ({put.status_code}): {put.text[:500]}")

    # 3. notify -> kicks off checks
    session.request("POST", f"{GEN_BASE}/media/{mid}/upload",
                    body={"filename": os.path.basename(path), "force_nsfw_check": True,
                          "force_ip_check": True, "surface": surface})

    # 4. poll until verification finishes. The poll response omits `url`, so
    #    carry the cloudfront url from the batch step into the returned object.
    print("🔎 Verifying (NSFW / IP check)…")
    event(event="progress", phase="verifying")
    deadline = time.time() + verify_timeout
    while True:
        media = session.request("GET", f"{GEN_BASE}/media/{mid}").json()
        status = media.get("status")
        # ONLY judge the verdict once the check has FINISHED. A mid-check `status` is not the result —
        # judging it early is exactly what flagged images Higgsfield then verifies as clean. So wait for
        # ip_check_finished, THEN read the final status: block if it's a real block verdict, else proceed.
        if media.get("ip_check_finished"):
            if status in BLOCK_STATUSES:
                print(f"⛔ id={mid} FINAL verdict={status} | media={json.dumps(media)}")  # dump for diagnosis
                raise ImageBlocked(status)
            media["url"] = url
            print(f"✅ Verified: id={mid} status={status} "
                  f"face={media.get('is_face_detected')}")
            return media
        if time.time() > deadline:
            raise VerifyTimeout(f"ip_check not finished in {verify_timeout}s")
        time.sleep(2)


def upload_with_retry(session, path, surface="seedance_2", attempts=3):
    """Upload + verify, retrying transient verify timeouts up to `attempts` times.
    A deterministic block (nsfw) is NOT retried — same bytes give the same verdict."""
    for attempt in range(1, attempts + 1):
        print(f"📤 Upload attempt {attempt}/{attempts}…")
        event(event="progress", phase="uploading", attempt=attempt, max=attempts)
        try:
            return upload_image(session, path, surface=surface)
        except ImageBlocked as b:
            event(event="error", reason="nsfw", status=b.status, fatal=True)
            raise SystemExit(f"Image blocked by verification (status={b.status})")
        except VerifyTimeout as t:
            print(f"   ⏱ verify timed out ({t}); retrying…")
            event(event="progress", phase="verify-retry", attempt=attempt, max=attempts)
    event(event="error", reason="verify-timeout", fatal=True)
    raise SystemExit(f"Verification did not finish after {attempts} attempts")


def media_data(media):
    """Build the generation-payload `data` block from a media object."""
    finished = bool(media.get("ip_check_finished", True))
    face = bool(media.get("is_face_detected", False))
    status = media.get("status", "uploaded")
    return {
        "id": media["id"], "type": "media_input", "url": media["url"],
        "status": status, "ip_check_finished": finished, "is_face_detected": face,
        "ipCheckFinished": finished, "isFaceDetected": face, "ipStatus": status,
    }


# (resolution, aspect) -> (width, height) for the seedance payload.
def _dims(resolution, aspect):
    short = 480 if str(resolution).startswith("480") else 720
    long = round(short * 16 / 9)
    return (long, short) if aspect == "16:9" else (short, long)


GEN_SLOT_BACKOFF_S = 20  # poll interval while the single job slot is busy


def generate(session, prompt, medias, duration=8, resolution="720p", aspect="16:9",
             gen_wait=480):
    """`medias` is an ORDERED list of verified media objects (1+ reference images).
    Order is preserved into the payload's medias[] — that's the @imageN / ImgN order."""
    if isinstance(medias, dict):  # tolerate a single media object
        medias = [medias]
    url = f"{GEN_BASE}/jobs/v2/seedance_unlimited"
    width, height = _dims(resolution, aspect)
    payload = {
        "params": {
            "prompt": prompt,
            "duration": int(duration),
            "aspect_ratio": aspect,
            "resolution": resolution,
            "bitrate_mode": "high",
            "generate_audio": True,
            "width": width,
            "height": height,
            "model": "seedance_unlimited",
            "medias": [{"role": "image", "data": media_data(m)} for m in medias],
        },
        "use_unlim": True,
        "use_free_gens": False,
    }
    print("\n🎬 Starting generation…")
    # Higgsfield allows 1 concurrent job. The backend FIFO queue serializes OUR
    # jobs, but a job started directly in the browser UI also holds that slot —
    # so on a 429 we wait it out (up to gen_wait) instead of failing.
    waited = 0
    while True:
        r = session.request("POST", url, body=payload)
        if r.status_code == 200:
            break
        if r.status_code == 429 and waited < gen_wait:
            print(f"   ⏳ job slot busy (429); waiting {GEN_SLOT_BACKOFF_S}s ({waited}/{gen_wait})…")
            event(event="progress", phase=f"Waiting for slot ({waited}s)")
            time.sleep(GEN_SLOT_BACKOFF_S)
            waited += GEN_SLOT_BACKOFF_S
            continue
        print(f"❌ Generation request failed ({r.status_code}):\n{r.text[:1000]}")
        reason = "rate-limited" if r.status_code == 429 else "generate-failed"
        event(event="error", reason=reason, status=r.status_code, detail=r.text[:300], fatal=True)
        sys.exit(1)
    job_id = r.json()["job_sets"][0]["jobs"][0]["id"]
    print(f"✅ Job ID: {job_id}")
    event(event="progress", phase="rendering", jobId=job_id)
    return job_id


def poll_and_download(session, job_id, download=True, own_rights=False):
    """Poll the job to completion. Returns (video_url, local_path|None).
    Emits a 'done' event (with the cloudfront url) for --emit-json consumers."""
    status_url = f"{GEN_BASE}/jobs/{job_id}/status"
    details_url = f"{GEN_BASE}/jobs/{job_id}"
    confirm_url = f"{GEN_BASE}/jobs/{job_id}/seedance/confirm-rights"
    rights_confirmed = False
    print("⏳ Rendering…")
    while True:
        r = session.request("GET", status_url)
        status = r.json().get("status")
        print(f"   status: {status}")
        # Rights gate: the clip RENDERED, but Higgsfield flagged copyright/likeness content and holds it behind
        # "I own rights to this content" — status stays 'ip_detected' (it NEVER becomes 'completed' on its own,
        # so without this we'd just poll until timeout). With explicit opt-in we fire the same PUT the UI button
        # does, which flips the job to 'completed'; otherwise we fail clearly so the user knows why.
        if status == "ip_detected":
            if not own_rights:
                print("⛔ rights verification required — output flagged copyright/likeness. Pass --own-rights to auto-confirm.")
                event(event="error", reason="rights-required", jobId=job_id, fatal=True)
                return None, None
            if not rights_confirmed:
                print("⚖️  asserting content rights (I own rights to this content)…")
                session.request("PUT", confirm_url)  # exact call the UI button fires → flips status to completed
                rights_confirmed = True
                event(event="progress", phase="rights-confirmed")
            time.sleep(2)
            continue
        if status == "completed":
            details = session.request("GET", details_url).json()
            video_url = details["results"]["raw"]["url"]
            print(f"🎬 {video_url}")
            path = None
            if download:
                os.makedirs(OUTPUT_DIR, exist_ok=True)
                path = os.path.join(OUTPUT_DIR, f"higgsfield_{job_id}.mp4")
                data = cffi_requests.get(video_url, impersonate=IMPERSONATE).content
                with open(path, "wb") as f:
                    f.write(data)
                print(f"🎉 Saved: {path}")
            event(event="done", jobId=job_id, videoUrl=video_url, file=path)
            return video_url, path
        if status in ("failed", "error", "cancelled", "nsfw", "moderated", "rejected"):
            # a moderation flag can land even though the render PRODUCED a clip — grab it if so rather than
            # discarding a finished video over a soft flag. Only fail when there's genuinely no output.
            details = session.request("GET", details_url).json()
            video_url = (((details.get("results") or {}).get("raw")) or {}).get("url")
            if video_url:
                print(f"🎬 (flagged {status}, but a clip rendered) {video_url}")
                path = None
                if download:
                    os.makedirs(OUTPUT_DIR, exist_ok=True)
                    path = os.path.join(OUTPUT_DIR, f"higgsfield_{job_id}.mp4")
                    data = cffi_requests.get(video_url, impersonate=IMPERSONATE).content
                    with open(path, "wb") as f:
                        f.write(data)
                    print(f"🎉 Saved: {path}")
                event(event="done", jobId=job_id, videoUrl=video_url, file=path)
                return video_url, path
            print(f"❌ Generation {status}")
            event(event="error", reason=f"render-{status}", jobId=job_id, fatal=True)
            return None, None
        time.sleep(5)


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--prompt", default=DEFAULT_PROMPT)
    ap.add_argument("--image-file", action="append", default=[],
                    help="Local image to upload + verify, then use. Repeat for ordered refs.")
    ap.add_argument("--image-id", default=DEFAULT_IMAGE_ID,
                    help="Use an already-uploaded media id (ignored if --image-file).")
    ap.add_argument("--image-url", default=DEFAULT_IMAGE_URL)
    ap.add_argument("--surface", default="seedance_2")
    ap.add_argument("--duration", type=int, default=8)
    ap.add_argument("--resolution", default="720p")
    ap.add_argument("--aspect", default="16:9")
    ap.add_argument("--retries", type=int, default=3, help="Upload+verify attempts.")
    ap.add_argument("--gen-wait", type=int, default=480,
                    help="Max seconds to wait out a busy job slot (429) before failing.")
    ap.add_argument("--resume-job", default=None,
                    help="Skip upload/generate; just poll + download an existing Higgsfield job id.")
    ap.add_argument("--media-json", default=None,
                    help="JSON list of already-verified media objects; skip upload, go to generate.")
    ap.add_argument("--no-download", action="store_true",
                    help="Don't save the mp4 locally; just report the video URL.")
    ap.add_argument("--own-rights", action="store_true",
                    help="Auto-confirm Higgsfield's copyright/likeness rights gate on a rendered output "
                         "(asserts you own the rights to the content). Without it, a gated output fails.")
    ap.add_argument("--emit-json", action="store_true",
                    help="Emit NDJSON progress on stdout (human logs go to stderr).")
    ap.add_argument("--capture-only", action="store_true",
                    help="Only harvest + print tokens, don't generate.")
    ap.add_argument("--upload-only", action="store_true",
                    help="Upload + verify --image-file, print the media object, stop.")
    args = ap.parse_args()

    # In --emit-json mode, machine events go to the real stdout; every human print()
    # is redirected to stderr so the two streams never interleave.
    if args.emit_json:
        global _JSON_OUT
        _JSON_OUT = sys.stdout
        sys.stdout = sys.stderr

    if args.capture_only:
        auth, datadome, cookies = get_fresh_tokens()
        print("\n--- CAPTURED ---")
        print("auth:", auth[:40], "…")
        print("datadome:", datadome)
        print("cookie len:", len(cookies))
        return

    session = BrowserSession()

    # Resume: a job was already submitted (id persisted); just poll it to completion.
    if args.resume_job:
        print(f"↩️  Resuming job {args.resume_job}")
        event(event="progress", phase="rendering", jobId=args.resume_job)
        poll_and_download(session, args.resume_job, download=not args.no_download, own_rights=args.own_rights)
        return

    if args.media_json:
        # Resume after a restart: media already uploaded + verified — skip straight to generate.
        medias = json.loads(args.media_json)
        print(f"↩️  Using {len(medias)} pre-verified media (skip upload)")
    elif args.image_file:
        # Upload + verify each ref IN ORDER → ordered medias[] (= @imageN / ImgN order).
        medias = []
        for n, f in enumerate(args.image_file, 1):
            print(f"📎 Reference {n}/{len(args.image_file)}: {os.path.basename(f)}")
            medias.append(upload_with_retry(session, f, surface=args.surface,
                                            attempts=args.retries))
        # Checkpoint: refs are verified → let the backend persist them so a restart
        # mid-render can resume without re-uploading.
        event(event="verified", medias=medias)
    else:
        medias = [{"id": args.image_id, "url": args.image_url, "status": "uploaded",
                   "ip_check_finished": True, "is_face_detected": False}]

    if args.upload_only:
        print("\n--- MEDIA ---")
        print(json.dumps(medias, indent=2))
        return

    job_id = generate(session, args.prompt, medias, duration=args.duration,
                      resolution=args.resolution, aspect=args.aspect, gen_wait=args.gen_wait)
    poll_and_download(session, job_id, download=not args.no_download, own_rights=args.own_rights)


if __name__ == "__main__":
    main()
