#!/usr/bin/env python3
"""
tryinfer.com (Infer) Seedance 2.0 reference-to-video, driven through your ALREADY-LOGGED-IN Chrome tab
over the DevTools Protocol — the same in-page-fetch trick as higgsfield_client.py. Every /api call runs
as fetch() INSIDE the real tab, so tryinfer's Clerk SESSION COOKIE rides along automatically (its API is
cookie-authed — no bearer, no token minting).

We pass reference_image_urls as ARBITRARY PUBLIC URLs (verified: tryinfer fetches them), so there is NO
upload step — the backend publishes each board to a public R2 URL and hands us the URL.

Flow (from bridge/tryinfer_capture.json):
  POST /api/create/generation-tasks         {args:{model,capability,input:{reference_image_urls,prompt,
                                             duration_seconds,aspect_ratio,resolution,audio}}, group:{…}}
                                             -> {id, status:"PENDING"}
  GET  /api/create/generation-tasks/{id}     -> status PENDING|RUNNING|SUCCEEDED|FAILED
  GET  /api/create/generation-tasks/{id}/result -> {output:{video_url, width, height, …}}

Two run modes:
  • --emit-json : machine mode for the desktop bridge — NDJSON events on stdout, human logs on stderr.
  • (default)   : standalone test — human logs on stderr, final JSON on stdout.

Deps (same Windows Python as the Higgsfield bridge): websocket-client
Usage (standalone test):
  python tryinfer_client.py --prompt "a slow cinematic push-in" [--image-url URL] [--duration 5] \
      [--aspect 1:1] [--resolution 1080p] [--no-audio]
"""
import argparse
import json
import re
import sys
import time
import urllib.request
import uuid

try:
    import websocket  # websocket-client
except ImportError:
    sys.exit("Missing dep. Install with:  pip install websocket-client")


def log(*a):
    print(*a, file=sys.stderr, flush=True)


# NDJSON progress protocol (only active in --emit-json). The backend's TryinferProvider parses these to
# drive the node's live status: {"event":"progress","phase":…}, {"event":"done","taskId","videoUrl"}, error.
_JSON_OUT = None


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


def http_json(url):
    with urllib.request.urlopen(url, timeout=10) as r:
        return json.loads(r.read().decode())


# ---- minimal synchronous CDP client (same shape as higgsfield_client.py) ----
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):
        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 close(self):
        try:
            self.ws.close()
        except Exception:
            pass


class PageResponse:
    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):
    """fetch() in the page context. credentials:'include' → the Clerk session cookie rides along
    (tryinfer's API is cookie-authed, so no Authorization bearer is needed, unlike Higgsfield)."""
    parts = [
        "(async()=>{try{",
        f"const o={{method:{json.dumps(method)},headers:{{}},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)


# JS run inside the tab to force Clerk to re-mint its short-lived session token (which updates the __session
# cookie tryinfer's API reads). Works on a HIDDEN/backgrounded tab — Runtime.evaluate ignores visibility — so
# we never foreground. Returns {ok:true} on success, {ok:false,reason|error} if Clerk isn't loaded yet.
_CLERK_REFRESH_EXPR = (
    "(async()=>{try{"
    "const c=window.Clerk;"
    "if(c&&c.session&&typeof c.session.getToken==='function'){"
    "await c.session.getToken({skipCache:true});"
    "return JSON.stringify({ok:true});}"
    "return JSON.stringify({ok:false,reason:'clerk-not-ready'});"
    "}catch(e){return JSON.stringify({ok:false,error:String((e&&e.message)||e)});}})()"
)


# CDP errors that mean the execution context we ran fetch() in is gone — the Studio tab NAVIGATED (it changes
# route when a generation starts/finishes) or reloaded. We recover by re-attaching to the tab's fresh context.
_CTX_DEAD = ("navigated", "context", "closed", "-32000", "detached", "Session with given id")


# ---- worker discovery ----
# Every logged-in tryinfer tab is an interchangeable WORKER. A worker is identified by its
# browserContextId — the isolated session (own cookies), which is what a separate Chrome profile or a temp
# context actually is. We deliberately do NOT identify workers by account/email: the pool is anonymous and
# its membership changes as profiles/contexts come and go.
#
# NOTE the gotcha this exists to avoid: /json/list is the easy endpoint but does NOT carry
# browserContextId, so it cannot tell two tryinfer tabs apart. Target.getTargets does.
def list_workers(cdp_http, match="tryinfer"):
    """[{contextId, targetId, url, title}] — one entry per open tryinfer tab, newest Chrome state."""
    ver = http_json(f"{cdp_http}/json/version")
    cdp = CDP(ver["webSocketDebuggerUrl"])
    try:
        infos = cdp.call("Target.getTargets", timeout=15).get("targetInfos", [])
    finally:
        cdp.close()
    out = []
    for t in infos:
        if t.get("type") == "page" and match in (t.get("url") or ""):
            out.append({"contextId": t.get("browserContextId"), "targetId": t.get("targetId"),
                        "url": t.get("url"), "title": t.get("title")})
    return out


class BrowserSession:
    def __init__(self, cdp_http, match, context=None):
        self.cdp_http = cdp_http
        self.match = match
        # when set, attach ONLY to a tab inside this browser context — that's how the backend pins a job to
        # the worker it leased. Without it we keep the legacy first-match-by-URL behaviour (single worker).
        self.context = context
        ver = http_json(f"{cdp_http}/json/version")
        log(f"Connected to {ver.get('Browser')}")
        self.browser_ws = ver["webSocketDebuggerUrl"]
        self.cdp = None
        self._attach()

    def _resolve_tab(self):
        """(targetId, url) of the tab to drive. Context-pinned when --context was given, else first URL match."""
        if self.context:
            for w in list_workers(self.cdp_http, self.match):
                if w["contextId"] == self.context:
                    return w["targetId"], w["url"]
            # the context is gone (a temp context dies with Chrome, a profile window can be closed) — say so
            # precisely, so the backend can drop that worker from the pool instead of guessing.
            raise RuntimeError(f"WORKER_GONE: no tryinfer tab in browser context {self.context}")
        for t in http_json(f"{self.cdp_http}/json/list"):
            if t.get("type") == "page" and self.match in (t.get("url") or ""):
                return t["id"], t["url"]
        return None, None

    def _attach(self):
        """(Re)resolve the matching tab and attach a fresh Runtime session. Safe to call repeatedly — used
        both at startup and to recover after the Studio tab navigates mid-poll (which kills the context)."""
        if self.cdp:
            try: self.cdp.close()
            except Exception: pass
        tid, url = self._resolve_tab()
        if not tid:
            raise RuntimeError(f"No open tab whose URL contains '{self.match}'. Open tryinfer.com and log in.")
        # A HOST tab left in the background CAN get discarded/frozen/throttled by Chrome (no live renderer),
        # in which case attaching works but Runtime.enable hangs waiting for a renderer → CDP timeout. We wake
        # it WITHOUT foregrounding (activateTarget steals focus, which is annoying in a browser the user is
        # actively using), escalating only as far as needed:
        #   attempt 0 — plain attach, short Runtime.enable timeout so a stalled tab fails fast.
        #   attempt 1 — focus-free un-freeze: Page.setWebLifecycleState('active') resumes a frozen renderer
        #               without bringing the tab forward.
        #   attempt 2 — LAST RESORT ONLY: Target.activateTarget (this is the sole path that steals focus, and
        #               only after a plain attach AND a focus-free un-freeze both failed — a fully discarded tab).
        last = None
        for attempt in range(3):
            self.cdp = CDP(self.browser_ws)
            try:
                if attempt == 2:
                    try: self.cdp.call("Target.activateTarget", {"targetId": tid}, timeout=10)  # foreground: last resort
                    except Exception: pass
                self.session_id = self.cdp.call("Target.attachToTarget", {"targetId": tid, "flatten": True}, timeout=15)["sessionId"]
                if attempt == 1:
                    # attempt 0 stalled → renderer frozen. Resume it without foregrounding.
                    try:
                        self.cdp.call("Page.enable", session_id=self.session_id, timeout=8)
                        self.cdp.call("Page.setWebLifecycleState", {"state": "active"}, session_id=self.session_id, timeout=8)
                    except Exception as e:
                        log(f"   focus-free un-freeze failed ({e})")
                # short timeout on the first pass → fail fast into the wake escalation if the renderer is asleep.
                self.cdp.call("Runtime.enable", session_id=self.session_id, timeout=(6 if attempt == 0 else 20))
                log(f"Routing API calls through tab: {url}")
                return
            except Exception as e:
                last = e
                try: self.cdp.close()
                except Exception: pass
                if attempt < 2:
                    log(f"   CDP attach stalled — {'un-freezing (no focus)' if attempt == 0 else 'foregrounding (last resort)'}… ({e})")
                    time.sleep(2)
        raise RuntimeError(f"couldn't attach to the tryinfer tab — is it awake + logged in? ({last})")

    def request(self, method, url, body=None, timeout=120, retries=4):
        last = None
        refreshed = False  # one focus-free session re-mint per request (a genuine logout still surfaces as 401)
        for attempt in range(retries):
            try:
                res = self.cdp.call(
                    "Runtime.evaluate",
                    {"expression": _fetch_expr(method, url, body), "awaitPromise": True, "returnByValue": True},
                    session_id=self.session_id, timeout=timeout)
            except Exception as e:
                last = f"CDP evaluate error: {e}"
                # The tab stopped answering CDP. Two shapes, ONE recovery: the context died (tab navigated /
                # closed — _CTX_DEAD) OR the renderer froze/was discarded and the evaluate TIMED OUT waiting
                # for a response. Previously we only re-attached on _CTX_DEAD, so a mid-request freeze just
                # hammered the same dead session until it failed. Now we re-attach on ANY evaluate error —
                # _attach focus-free un-freezes a stalled renderer (setWebLifecycleState) — and it also swaps
                # in a fresh CDP socket, clearing any response desynced by the timeout.
                log(f"   tab unresponsive — re-attaching (no focus)… ({e})")
                time.sleep(1)
                try: self._attach()
                except Exception as e2: last = f"re-attach failed: {e2}"
                if attempt < retries - 1:
                    time.sleep(2); continue
                raise RuntimeError(last)
            val = res.get("result", {}).get("value")
            if val is None:
                last = f"evaluate failed: {res.get('exceptionDetails')}"
            else:
                data = json.loads(val)
                status = data.get("status")
                if status == 0:
                    last = f"in-page fetch error: {data.get('error')}"
                elif status in (401, 403) and not refreshed:
                    # Stale Clerk session: a backgrounded tab stops refreshing __session, so tryinfer 401s.
                    # Re-mint the cookie WITHOUT foregrounding, then retry immediately with the fresh session.
                    refreshed = True
                    last = f"auth {status} — re-minting session (no focus)"
                    log(f"   {last}…")
                    try: self.refresh_session(hard=True)
                    except Exception as e2: last = f"session refresh failed: {e2}"
                    continue  # skip the backoff sleep — retry now that the cookie is fresh
                else:
                    return PageResponse(status, data["body"])
            if attempt < retries - 1:
                log(f"   transient request error, retrying… ({last})")
                time.sleep(3)
        raise RuntimeError(last)

    def _clerk_mint(self):
        """Force Clerk to mint a fresh session token (updates __session) via evaluated JS. No foreground.
        Returns True if a token was minted, False if Clerk isn't loaded / errored."""
        try:
            res = self.cdp.call(
                "Runtime.evaluate",
                {"expression": _CLERK_REFRESH_EXPR, "awaitPromise": True, "returnByValue": True},
                session_id=self.session_id, timeout=20)
            val = res.get("result", {}).get("value")
            data = json.loads(val) if val else {}
            if data.get("ok"):
                log("   session re-minted via Clerk (no focus)")
                return True
            log(f"   Clerk refresh unavailable: {data.get('reason') or data.get('error')}")
        except Exception as e:
            log(f"   Clerk refresh errored: {e}")
        return False

    def refresh_session(self, hard=False):
        """Re-mint the Clerk session cookie WITHOUT foregrounding the tab. A backgrounded tab stops running
        Clerk's periodic token refresh, so __session goes stale and tryinfer 401s. Focus-free recovery:
          1. Page.setWebLifecycleState('active') — resume a throttled/frozen renderer (does NOT foreground).
          2. window.Clerk.session.getToken({skipCache:true}) — force a fresh token → fresh cookie.
          3. (hard only) if Clerk isn't reachable, Page.reload + re-attach so Clerk re-inits from scratch.
        Returns True if a fresh token was minted."""
        try:
            self.cdp.call("Page.enable", session_id=self.session_id, timeout=8)
            self.cdp.call("Page.setWebLifecycleState", {"state": "active"}, session_id=self.session_id, timeout=8)
        except Exception:
            pass
        if self._clerk_mint():
            return True
        if hard:
            log("   Clerk not ready — reloading tab to re-init the session (no focus)…")
            try: self.cdp.call("Page.reload", session_id=self.session_id, timeout=15)
            except Exception as e: log(f"   reload failed: {e}")
            time.sleep(5)
            try: self._attach()
            except Exception as e: log(f"   re-attach after reload failed: {e}")
            return self._clerk_mint()
        return False

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


API = "https://tryinfer.com/api"
RATIO_WORD = {"1:1": "square", "16:9": "landscape", "9:16": "portrait", "4:3": "landscape", "3:4": "portrait"}
DONE = {"SUCCEEDED", "COMPLETED"}
FAILED = {"FAILED", "ERROR", "CANCELLED", "CANCELED"}
# Consecutive status responses with NO "status" field we ride out before giving up. A task is only visible
# to the session that created it, so a context/worker mismatch (or a lapsed login) answers with a body that
# has no status at all — but a single blip mid-render shouldn't kill an otherwise healthy job.
MISSING_STATUS_TOLERANCE = 5

# ---- rejection classes ----------------------------------------------------------------------------
# tryinfer REJECTING us is not the same as a render failing, and the difference decides who retries:
#   WORKER_BUSY   → this account is mid-generation; the backend parks the worker and takes another.
#   RATE_LIMITED  → tryinfer/its model provider is pushing back. A different worker does NOT help — only
#                   TIME does — so the backend requeues the whole job on a backoff. Raised for 429s, 5xx,
#                   and for a task that FAILED with a rate-limit message (the provider throttled the render
#                   itself, exactly what the studio UI shows as "The model provider is rate-limiting").
# Anything else stays a hard failure. Retryable statuses are the transport-level ones only: a 4xx that is
# not 429 means the request itself is wrong, and retrying that forever just burns the wind-down clock.
def is_throttled(status):
    return status == 429 or (status is not None and status >= 500)


RATE_LIMIT_RE = re.compile(r"rate[- ]?limit|too many requests|quota|capacity|overloaded|try again", re.I)
MAX_THROTTLED_POLLS = 8      # consecutive rejected status polls before we hand the job back to the backend
POLL_BACKOFF_CAP_S = 120     # ceiling on the in-client poll backoff
# How often we emit a progress event even when nothing changed. The hub's job timer resets on every progress
# EVENT, but on_status only fires on a status CHANGE — and a task sits in PENDING/RUNNING for the whole
# render. So a perfectly healthy long render sent NOTHING for 30 minutes and the hub killed it with
# "tryinfer render timed out". Beating keeps the timer meaning "the client is alive" rather than "the status
# happened to change", which is what it was always documented to mean.
HEARTBEAT_S = 20


def submit(session, prompt, image_urls, duration, aspect, resolution, audio,
           capability="reference-to-video", model="seedance-2.0-pro", last_frame_url=None,
           num_images=1, input_json=None, audio_ref_url=None):
    """POST a generation task. Video + image shapes, all URL-based (no upload). Returns the task id.
      • reference-to-video : reference_image_urls[] + resolution + audio
      • image-to-video     : image_url (start) [+ last_frame_image_url] + audio      (no resolution)
      • edit (Seedream)    : image_url + prompt + num_images                         (UI hides aspect)
      • text-to-image      : prompt + num_images + aspect_ratio                       (no image)
    `input_json` overrides the WHOLE input dict verbatim — for probing undocumented combos."""
    medium = "image" if capability in ("edit", "text-to-image") else "video"
    if input_json is not None:
        inp = json.loads(input_json)  # PROBE: send exactly this
        meta = {"prompt": prompt, "medium": medium}
    elif capability == "image-to-video":
        if not image_urls:
            raise RuntimeError("image-to-video needs a start frame (--image-url)")
        inp = {"image_url": image_urls[0], "prompt": prompt, "duration_seconds": duration, "aspect_ratio": aspect, "audio": audio}
        if last_frame_url:
            inp["last_frame_image_url"] = last_frame_url
        meta = {"prompt": prompt, "medium": "video", "ratio": RATIO_WORD.get(aspect, "square"), "kind": "animate", "sourceUrl": image_urls[0]}
        if last_frame_url:
            meta["lastFrameUrl"] = last_frame_url
    elif capability == "edit":
        if not image_urls:
            raise RuntimeError("edit needs a source image (--image-url)")
        inp = {"image_url": image_urls[0], "prompt": prompt, "num_images": num_images}
        meta = {"prompt": prompt, "medium": "image", "kind": "edit", "sourceUrl": image_urls[0]}
    elif capability == "text-to-image":
        inp = {"prompt": prompt, "num_images": num_images, "aspect_ratio": aspect}
        meta = {"prompt": prompt, "medium": "image", "kind": "generate"}
    else:  # reference-to-video
        inp = {"reference_image_urls": image_urls, "prompt": prompt, "duration_seconds": duration, "aspect_ratio": aspect, "resolution": resolution, "audio": audio}
        # Seedance 2.5 accepts a reference AUDIO clip as a PUBLIC URL (verified in the web-UI capture:
        # input.reference_audio_url, a single URL string — the backend publishes the clip to R2 and hands it
        # here, same as the image refs; there is no upload step). Omitted → plain reference-to-video.
        if audio_ref_url:
            inp["reference_audio_url"] = audio_ref_url
        meta = {"prompt": prompt, "medium": "video", "ratio": RATIO_WORD.get(aspect, "square"), "kind": "reference", "referenceImageUrls": image_urls, "referenceRequestIds": []}
    body = {"args": {"model": model, "capability": capability, "input": inp},
            "group": {"groupId": str(uuid.uuid4()), "position": 0, "meta": meta}}
    r = session.request("POST", f"{API}/create/generation-tasks", body=body)
    if r.status_code not in (200, 201, 202):
        # tryinfer allows ONE generation at a time per account. A 429 here means this worker's account is
        # already busy — often with something started outside this app (their web UI, another tool). That is
        # not a failure of the render: mark it distinctly so the backend can release the worker, cool it
        # down and requeue onto another one, instead of failing the job.
        if r.status_code == 429 and "create_generation_in_progress" in (r.text or ""):
            raise RuntimeError(f"WORKER_BUSY: account already has a generation running ({r.text[:200]})")
        # Any OTHER 429, or a 5xx: tryinfer is throttling or wobbling. Another worker won't help — the whole
        # job needs to come back later — so mark it retryable and let the backend's queue own the backoff.
        if is_throttled(r.status_code):
            raise RuntimeError(f"RATE_LIMITED: submit rejected (HTTP {r.status_code}): {r.text[:300]}")
        raise RuntimeError(f"submit failed ({r.status_code}): {r.text[:800]}")
    tid = r.json().get("id")
    if not tid:
        raise RuntimeError(f"submit returned no task id: {r.text[:400]}")
    return tid


def poll(session, task_id, on_status=None, timeout=1800, interval=3):
    """Poll the task to a terminal state. Fires on_status(status) on each change. Returns final task JSON."""
    deadline = time.time() + timeout
    last = None
    last_beat = 0.0   # 0 → the first poll always beats, so the hub hears from us immediately
    missing = 0
    throttled = 0
    while time.time() < deadline:
        r = session.request("GET", f"{API}/create/generation-tasks/{task_id}")
        # A REJECTED poll is not a failed render — the task is still running on their side. poll() used to
        # ignore status codes entirely, so a 429 body (no "status" key) fell through as a statusless
        # response and killed the job. Back off instead, widening each time, and only give up after
        # MAX_THROTTLED_POLLS — then as RATE_LIMITED so the backend requeues rather than fails.
        if is_throttled(r.status_code):
            throttled += 1
            if throttled > MAX_THROTTLED_POLLS:
                raise RuntimeError(
                    f"RATE_LIMITED: tryinfer rejected the status poll {throttled}x "
                    f"(HTTP {r.status_code}): {(r.text or '')[:200]}")
            wait = min(POLL_BACKOFF_CAP_S, interval * (2 ** throttled))
            log(f"   poll rejected HTTP {r.status_code} — backing off {wait:.0f}s ({throttled}/{MAX_THROTTLED_POLLS})")
            event(event="progress", taskId=task_id,
                  phase=f"waiting · tryinfer rate-limited (retry in {int(wait)}s)")
            time.sleep(wait)
            continue
        throttled = 0
        try:
            j = r.json()
        except Exception:   # an HTML error page / truncated body — same bucket as a statusless response
            j = {"_nonjson": (r.text or "")[:200]}
        st = j.get("status")
        # No "status" key at all — NOT a status we don't recognise, but a response that isn't a task at all
        # (task invisible to this worker context, expired session, an error envelope). This used to fall
        # through to on_status(None) and die as "'NoneType' object has no attribute 'lower'", which threw
        # away the one thing that explains it: the body. Ride out a blip, then fail WITH the body.
        if st is None:
            missing += 1
            if missing > MISSING_STATUS_TOLERANCE:
                raise RuntimeError(
                    f"task {task_id}: status response carried no 'status' {missing}x — the task may not be "
                    f"visible to this worker context, or the tryinfer login lapsed. Last body: "
                    f"{json.dumps(j)[:500]}")
            time.sleep(interval)
            continue
        missing = 0
        changed = st != last
        if changed:
            last = st
        now = time.time()
        if on_status and (changed or now - last_beat >= HEARTBEAT_S):
            last_beat = now
            on_status(st)
        if st in DONE:
            return j
        if st in FAILED:
            err = j.get("error")
            detail = json.dumps(err) if err else "no detail"
            # The task itself was throttled by the model provider (studio shows "The model provider is
            # rate-limiting requests right now"). Retrying LATER is exactly right, so hand it back as
            # retryable instead of failing the node.
            if RATE_LIMIT_RE.search(detail):
                raise RuntimeError(f"RATE_LIMITED: task {st} — {detail[:300]}")
            raise RuntimeError(f"task {st}: {detail}")
        time.sleep(interval)
    raise TimeoutError(f"task {task_id} did not finish within {timeout}s (last status {last})")


def get_result(session, task_id, retries=6, interval=5):
    """GET the result. Returns {'videoUrl':…} for video, or {'images':[url,…]} for image. Raises on
    moderation block / empty output.

    Retries a THROTTLED fetch in-place rather than raising: by here the render has already succeeded, and
    losing a finished clip to a 429 on the last hop would mean re-rendering something we already have.
    """
    last = None
    for attempt in range(retries):
        r = session.request("GET", f"{API}/create/generation-tasks/{task_id}/result")
        if not is_throttled(r.status_code):
            break
        last = f"HTTP {r.status_code}: {(r.text or '')[:200]}"
        wait = min(POLL_BACKOFF_CAP_S, interval * (2 ** attempt))
        log(f"   result fetch rejected {last} — retrying in {wait:.0f}s ({attempt + 1}/{retries})")
        event(event="progress", taskId=task_id,
              phase=f"waiting · tryinfer rate-limited on the result (retry in {int(wait)}s)")
        time.sleep(wait)
    else:
        raise RuntimeError(f"RATE_LIMITED: the finished render's result stayed throttled — {last}")
    j = r.json()
    mod = j.get("moderation_status")
    if mod and mod != "allowed":
        raise RuntimeError(f"moderation blocked: {mod}")
    out = j.get("output") or {}
    if out.get("video_url"):
        return {"videoUrl": out["video_url"], "output": out}
    imgs = [x.get("url") for x in (out.get("images") or []) if isinstance(x, dict) and x.get("url")]
    if imgs:
        return {"images": imgs, "output": out}
    raise RuntimeError(f"no video_url/images in result: {json.dumps(j)[:800]}")


def main():
    ap = argparse.ArgumentParser()
    ap.add_argument("--prompt", default="")
    ap.add_argument("--capability", default="reference-to-video",
                    choices=["reference-to-video", "image-to-video", "edit", "text-to-image"])
    ap.add_argument("--model", default=None, help="defaults by capability: image→seedream-5.0-pro, video→seedance-2.0-pro")
    ap.add_argument("--image-url", action="append", default=[],
                    help="public image URL. reference-to-video: repeat for ordered refs. image-to-video/edit: first = start/source.")
    ap.add_argument("--last-frame-url", default=None, help="image-to-video only: END frame URL (optional)")
    ap.add_argument("--num", type=int, default=1, help="num_images (image capabilities) — probe >1 here")
    ap.add_argument("--input-json", default=None,
                    help="PROBE: raw JSON for args.input verbatim (override the built payload — test undocumented combos)")
    ap.add_argument("--audio-ref-url", default=None,
                    help="reference-to-video (Seedance 2.5): public URL of a reference AUDIO clip → input.reference_audio_url")
    ap.add_argument("--duration", type=int, default=5)
    ap.add_argument("--aspect", default="1:1")
    ap.add_argument("--resolution", default="1080p")
    ap.add_argument("--no-audio", action="store_true")
    ap.add_argument("--resume-task", default=None, help="skip submit; just poll+result an existing task id")
    ap.add_argument("--emit-json", action="store_true", help="NDJSON events on stdout, human logs on stderr")
    ap.add_argument("--match", default="tryinfer", help="substring of the logged-in tab's URL")
    ap.add_argument("--context", default=None,
                    help="pin to the tryinfer tab in this browserContextId (a pool worker). Omit = first URL match.")
    ap.add_argument("--list-workers", action="store_true",
                    help="print [{contextId,targetId,url,title}] for every open tryinfer tab and exit (discovery)")
    ap.add_argument("--host", default="localhost")
    ap.add_argument("--port", type=int, default=9222)
    args = ap.parse_args()

    # discovery mode: no render, no tab attach — just report the pool as Chrome currently sees it.
    if args.list_workers:
        try:
            print(json.dumps({"workers": list_workers(f"http://{args.host}:{args.port}", args.match)}))
        except Exception as e:
            print(json.dumps({"workers": [], "error": str(e)}))
        return

    if args.emit_json:
        global _JSON_OUT
        _JSON_OUT = sys.stdout
        sys.stdout = sys.stderr  # any stray print() can't corrupt the event stream

    # model default is capability-aware: image caps use seedream (seedANCE is the VIDEO model — a common mixup).
    model = args.model or ("seedream-5.0-pro" if args.capability in ("edit", "text-to-image") else "seedance-2.0-pro")
    image_urls = args.image_url or ["https://www.gstatic.com/webp/gallery/1.jpg"]  # default = permissive test image
    sess = None
    try:
        # inside the try so a connect/tab failure ("No open tab … Open tryinfer.com and log in") surfaces as a
        # clean error EVENT the node can show — not a raw Python traceback dumped to the UI.
        sess = BrowserSession(f"http://{args.host}:{args.port}", args.match, context=args.context)
        if args.resume_task:
            task_id = args.resume_task
            log(f"resuming task {task_id}")
            event(event="progress", phase="rendering", taskId=task_id)
        else:
            task_id = submit(sess, args.prompt, image_urls, args.duration, args.aspect, args.resolution, not args.no_audio,
                             capability=args.capability, model=model, last_frame_url=args.last_frame_url,
                             num_images=args.num, input_json=args.input_json, audio_ref_url=args.audio_ref_url)
            log(f"submitted task {task_id}")
            event(event="progress", phase="submitted", taskId=task_id)
        poll(sess, task_id, on_status=lambda st: (log(f"   status={st}"),
                                                  event(event="progress", phase=str(st or "pending").lower(), taskId=task_id)))
        res = get_result(sess, task_id)
    except Exception as e:
        event(event="error", reason="tryinfer", detail=str(e))
        log(f"\n❌ {e}")
        sys.exit(1)
    finally:
        if sess:
            sess.close()

    if res.get("images"):
        event(event="done", taskId=task_id, images=res["images"])
        dims = {x.get("url"): (x.get("width"), x.get("height")) for x in ((res.get("output") or {}).get("images") or []) if isinstance(x, dict)}
        log(f"\n✅ {len(res['images'])} image(s):")
        for u in res["images"]:
            w, h = dims.get(u, (None, None))
            log(f"   {w}x{h}  {u}")
        if not args.emit_json:
            print(json.dumps({"ok": True, "images": res["images"]}))
    else:
        out = res.get("output") or {}
        event(event="done", taskId=task_id, videoUrl=res["videoUrl"],
              width=out.get("width"), height=out.get("height"), duration=out.get("duration_seconds"))
        log(f"\n✅ {out.get('width')}x{out.get('height')}  {out.get('duration_seconds')}s")
        if not args.emit_json:
            print(json.dumps({"ok": True, "video_url": res["videoUrl"]}))


if __name__ == "__main__":
    main()
