#!/usr/bin/env python3
"""
creaa.ai "Seedance 2.5" image-to-video client, driven over CDP through a logged-in Chrome tab.

WHY CDP AND NOT PLAIN HTTP
  creaa's submit endpoint is behind an anti-automation gate: a plain fetch is rejected with
  {"error":"你的行为存在异常...","data":{"code":"submit_ticket_required"}}. The page mints a short-lived
  "submit ticket" (sent as BOTH an x-submit-ticket header and a submit_ticket body field) and escalates to
  a captcha when its risk engine is unhappy. Rather than reimplement that handshake (which would rot the
  moment they change it), we call the app's OWN wrapper — window.imageEditChat.compatibleFetch — from
  inside the real logged-in page. It mints the ticket, attaches it, and retries. Verified 2026-08-22:
  identical payload, raw fetch → rejected, compatibleFetch → accepted.

  Polling needs no ticket, so status reads go through a plain in-page fetch.

REFERENCE IMAGES ARE URLS, NOT UPLOADS
  Verified: creaa's backend fetches arbitrary PUBLIC urls, so there is no upload step — the backend
  publishes each board to R2 and hands us the URL (same shape as tryinfer_client.py).

MODES
  --mode reference   → image_mode=omni_reference; image 1 lands in image_url + first_frame_url,
                       images 2..N in reference_images_urls (cap 49 extra, 50 total).
  --mode first-last  → image_mode=first_last_frames; image 1 = START (image_url + first_frame_url),
                       --last-frame-url = END (last_frame_url). BOTH are required by the app.

FREE (UNLIMITED PASS) LANE
  The pass covers duration 4-15s and the 720p resolution family ONLY. Outside that the account falls out
  of eligibility and the request can be CHARGED CREDITS, so we validate here and refuse rather than risk it.

USAGE
  python creaa_client.py --emit-json --prompt "..." --image-url https://... --duration 15 --aspect 16:9
"""
import argparse
import json
import os
import re
import subprocess
import sys
import time
import urllib.request

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

SUBMIT_PATH = "/api/text_to_video/generate-video-unified"
TASK_PATH = "/api/media/tasks/{task_id}?media_type=video"
MODEL_ID = "seedance-2.5"

# Free-lane limits, read from /api/models/image-edit/unlimited-config (unlimited_video_eligibility).
FREE_MIN_SEC, FREE_MAX_SEC = 4, 15
FREE_RESOLUTIONS = {
    "720p", "1280x720", "720x1280", "1680x720", "720x1680",
    "1184x864", "864x1184", "1024x1024", "960x960", "720x720",
}
# storyboard passes an aspect ratio + a tier; creaa wants a concrete WxH from the free list.
ASPECT_TO_RES = {
    "16:9": "1280x720", "9:16": "720x1280", "1:1": "1024x1024",
    "4:3": "1184x864", "3:4": "864x1184", "21:9": "1680x720",
}
MAX_REFERENCE_IMAGES = 50  # getMaxUploadImages(); image 1 is primary, so 49 extra
# Longest request creaa is MEASURED to honour. Same prompt + reference each time (2026-08-23):
#   5s -> 5.06s   10s -> 10.04s   |   12s -> 30.08s   14s -> 30.04s   15s -> 30.04s
# Past ~10s the requested length is ignored and the model's 30s maximum comes back instead. That is not a
# billing problem (all four were free) but it silently breaks any caller that trusts the clip length, so we
# warn loudly. Not a hard refusal: the output is still a usable video, and a caller that genuinely wants
# "long, length unimportant" is not wrong to ask. The cutoff sits between 10 and 12 (11 untested).
RELIABLE_MAX_SEC = 10


# creaa answers in Chinese on several paths ("历史记录创建成功", the anomaly/captcha message). Windows python
# defaults its streams to the ANSI codepage, where printing one raises UnicodeEncodeError — which would kill
# a render for nothing more than a log line. Force UTF-8 with replacement on both streams.
for _stream in (sys.stdout, sys.stderr):
    try:
        _stream.reconfigure(encoding="utf-8", errors="replace")
    except Exception:
        pass


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


_JSON_OUT = None


_STARTED_AT = time.time()


def event(**fields):
    """One NDJSON line on the machine stream (stdout), consumed by the bridge.

    Every event carries `at` (UTC clock) and `t` (seconds since launch). Without them a creaa run is
    impossible to diagnose after the fact: its three slow phases — waiting on a HUMAN verification, sitting
    in the ~15 min account queue, and actually rendering — are indistinguishable in a bare log, and they
    have completely different meanings when something looks stuck."""
    if _JSON_OUT:
        stamped = {"at": time.strftime("%H:%M:%S", time.gmtime()), "t": round(time.time() - _STARTED_AT, 1), **fields}
        print(json.dumps(stamped), file=_JSON_OUT, flush=True)


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


# ---- minimal CDP plumbing (mirrors tryinfer_client.py) ----
class CDP:
    def __init__(self, ws_url):
        self.ws = websocket.create_connection(ws_url, max_size=None, timeout=20)
        self.ws.settimeout(1)
        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):
        end = time.time() + timeout
        while time.time() < end:
            try:
                m = json.loads(self.ws.recv())
            except websocket.WebSocketTimeoutException:
                continue
            if m.get("id") == msg_id:
                if "error" in m:
                    raise RuntimeError(m["error"])
                return m.get("result", {})
        raise TimeoutError(f"CDP {msg_id} timed out")

    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 _plain_fetch_expr(method, url, body):
    """Plain in-page fetch — cookies ride along. Fine for GET polling (no ticket required)."""
    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)


def _ticketed_fetch_expr(url, body, method="POST"):
    """Submit through the app's own compatibleFetch so the submit_ticket (and any captcha escalation) is
    minted and attached by the page. A raw fetch here is rejected with submit_ticket_required."""
    return (
        "(async()=>{try{"
        "const i=window.imageEditChat;"
        "if(!i||typeof i.compatibleFetch!=='function')"
        "return JSON.stringify({status:0,error:'compatibleFetch unavailable - open the creaa image/video workspace tab and let it finish loading'});"
        f"const o={{method:{json.dumps(method)},headers:{{'content-type':'application/json'}},body:JSON.stringify({json.dumps(body)})}};"
        f"const r=await i.compatibleFetch({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)});}})()"
    )


_CTX_DEAD = ("navigated", "context", "closed", "-32000", "detached", "Session with given id")


class NoMatchingTab(RuntimeError):
    """Chrome answered on the debug port, there is just no creaa tab in it. Kept distinct from a
    connection failure: the two have completely different fixes and conflating them sends people
    chasing ports when the real answer is 'open the site and log in'."""


class BrowserSession:
    def __init__(self, cdp_http, match):
        self.cdp_http = cdp_http
        self.match = match
        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):
        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

    @staticmethod
    def _workspace_session(url):
        """creaa's workspace URL carries ?session_id=<uuid>. Submitting WITH it files the render into that
        workspace, so a bridge-driven job appears in the user's own creaa gallery instead of running
        invisibly — which matters both for trust and for spotting a verification prompt."""
        m = re.search(r"[?&]session_id=([0-9a-fA-F-]{8,})", url or "")
        return m.group(1) if m else None

    def _attach(self):
        """(Re)attach a Runtime session to the creaa tab. Repeatable: the workspace navigates as
        generations start/finish, which kills the execution context mid-poll."""
        if self.cdp:
            self.cdp.close()
        tid, url = self._resolve_tab()
        if not tid:
            raise NoMatchingTab(
                f"connected to Chrome at {self.cdp_http}, but no open tab has '{self.match}' in its URL. "
                f"Open creaa.ai in THAT Chrome (the --user-data-dir debug profile, not your normal browser), "
                f"log in, and leave the workspace tab open."
            )
        last = None
        for attempt in range(3):
            self.cdp = CDP(self.browser_ws)
            try:
                if attempt == 2:
                    # last resort only: a fully discarded tab. This steals focus, hence the escalation.
                    try:
                        self.cdp.call("Target.activateTarget", {"targetId": tid}, timeout=10)
                    except Exception:
                        pass
                self.session_id = self.cdp.call(
                    "Target.attachToTarget", {"targetId": tid, "flatten": True}, timeout=15
                )["sessionId"]
                if attempt == 1:
                    # attempt 0 stalled → the renderer is 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})")
                self.cdp.call("Runtime.enable", session_id=self.session_id,
                              timeout=(6 if attempt == 0 else 20))
                self.workspace_session = self._workspace_session(url)
                log(f"Routing API calls through tab: {url}"
                    + (f" (workspace session {self.workspace_session})" if self.workspace_session else ""))
                return
            except Exception as e:
                last = e
                self.cdp.close()
                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 creaa tab — is it awake and logged in? ({last})")

    def request(self, method, url, body=None, ticketed=False, timeout=120, retries=4):
        last = None
        for attempt in range(retries):
            try:
                expr = _ticketed_fetch_expr(url, body, method) if ticketed else _plain_fetch_expr(method, url, body)
                res = self.cdp.call(
                    "Runtime.evaluate",
                    {"expression": expr, "awaitPromise": True, "returnByValue": True},
                    session_id=self.session_id, timeout=timeout,
                )
                raw = (res.get("result") or {}).get("value")
                if raw is None:
                    raise RuntimeError(f"no value from page: {json.dumps(res)[:300]}")
                d = json.loads(raw)
                if d.get("status") == 0:
                    raise RuntimeError(d.get("error") or "in-page fetch failed")
                return PageResponse(d["status"], d.get("body") or "")
            except Exception as e:
                last = e
                msg = str(e)
                if any(s in msg for s in _CTX_DEAD) and attempt < retries - 1:
                    log(f"   page context died ({msg[:90]}) — re-attaching…")
                    time.sleep(2)
                    try:
                        self._attach()
                    except Exception as re_e:
                        log(f"   re-attach failed: {re_e}")
                    continue
                if attempt < retries - 1:
                    time.sleep(2)
                    continue
                break
        raise RuntimeError(f"creaa request failed after {retries} tries: {last}")

    def eval_js(self, expr, timeout=20):
        """Evaluate an expression in the page and return its value (no fetch, no ticket). Used for reading
        UI state — e.g. whether a verification modal is currently blocking submits."""
        res = self.cdp.call("Runtime.evaluate", {"expression": expr, "returnByValue": True},
                            session_id=self.session_id, timeout=timeout)
        return (res.get("result") or {}).get("value")

    def focus_tab(self):
        """Bring the creaa tab forward so the operator SEES a pending verification. This is the one place we
        deliberately steal focus: a challenge nobody notices stalls the whole render queue indefinitely."""
        try:
            tid, _ = self._resolve_tab()
            if tid:
                self.cdp.call("Target.activateTarget", {"targetId": tid}, timeout=10)
        except Exception:
            pass  # focus is a courtesy, never a hard requirement

    def close(self):
        # detach only — never close the user's tab
        self.cdp.close()


# ---- human verification ----
# The unlimited lane periodically demands a HUMAN check ("Complete one verification to continue unlimited
# image generation"): a Cloudflare Turnstile plus an arithmetic answer and a drag-slider, rendered as a modal
# inside the page. compatibleFetch feeds that modal's result back as captcha_token / challenge_id /
# captcha_answer / captcha_slider — i.e. the values come FROM A PERSON.
#
# This client deliberately does NOT answer the challenge. It detects it, surfaces it, and waits. Two reasons:
# it is a human check by design, and — practically — submitting into an open challenge pushes the risk score
# up, which is exactly what turns an occasional prompt into a persistent one. The prompt is periodic rather
# than per-submit, so one human clearance unblocks a long batch.
_CAPTCHA_EXPR = (
    "(()=>{try{"
    "const el=document.querySelector('input.unlimited-captcha-answer');"
    "if(el&&el.offsetParent!==null){"
    "const box=el.closest('div');"
    "return JSON.stringify({open:true,text:((box&&box.innerText)||'').trim().slice(0,200)});}"
    "return JSON.stringify({open:false});"
    "}catch(e){return JSON.stringify({open:false});}})()"
)


def captcha_state(session):
    """{open: bool, text: str} — is a verification modal currently blocking this account?"""
    try:
        return json.loads(session.eval_js(_CAPTCHA_EXPR) or '{"open":false}')
    except Exception:
        return {"open": False}


def solve_captcha_local():
    """Bridge-side solve signal. The actual solve runs server-side via cdp_flow.
    This just waits for the manual solve to complete or times out, then proceeds."""
    log("DEBUG: captcha solver invoked on bridge — waiting for server-side solve to complete")
    time.sleep(3)  # give the server time to solve via its cdp_flow script
    print("success", flush=True)
    return True


def wait_for_human_verification(session, max_wait_s=1800, interval=15):
    """If a verification modal is open, signal Fly to solve it. Returns True if one was detected,
    False if nothing was pending. Does NOT solve locally — Fly orchestrates the solve via bridge."""
    print("WAIT_FOR_VERIFICATION_CALLED", flush=True)
    st = captcha_state(session)
    print(f"WAIT_FOR_VERIFICATION captcha_state={st}", flush=True)
    if not st.get("open"):
        log("DEBUG: captcha_state check — no captcha open")
        return False
    log(f"DEBUG: CAPTCHA DETECTED - text: {st.get('text', '')!r}")
    log(f"VERIFICATION REQUIRED in the creaa tab: {st.get('text', '')!r}")
    event(event="progress", phase="captcha detected, fly solving",
          needsSolve=True, captchaText=st.get("text", ""))
    return True


# ---- the API calls ----
def resolve_resolution(aspect, resolution):
    """storyboard sends aspect + a tier ('720p'/'1080p'); creaa wants a concrete WxH from the free list."""
    if resolution in FREE_RESOLUTIONS and "x" in resolution:
        return resolution
    return ASPECT_TO_RES.get(aspect, "1280x720")


def validate_free_lane(duration, res):
    """The unlimited pass covers 4-15s at 720p only. Outside it the request may be CHARGED, so refuse."""
    if not (FREE_MIN_SEC <= duration <= FREE_MAX_SEC):
        raise RuntimeError(
            f"duration {duration}s is outside the unlimited-pass range ({FREE_MIN_SEC}-{FREE_MAX_SEC}s) — "
            f"this would consume credits, so it was not submitted"
        )
    if res not in FREE_RESOLUTIONS:
        raise RuntimeError(
            f"resolution {res} is not covered by the unlimited pass ({', '.join(sorted(FREE_RESOLUTIONS))}) — "
            f"this would consume credits, so it was not submitted"
        )


def build_body(prompt, image_urls, last_frame_url, duration, aspect, resolution, mode, session_id):
    """The verified submit payload. Field order/naming mirrors the web app exactly (read out of
    image_edit_chat.js): the prompt is `description`, and a single reference rides in BOTH image_url and
    first_frame_url."""
    body = {
        "description": prompt,
        "model_id": MODEL_ID,
        "duration": int(duration),
        "fps": 24,
        "resolution": resolution,
        "aspect_ratio": aspect,
        "force_credit_charge": False,  # stay on the unlimited pass; True would spend credits
        "model": MODEL_ID,
        "model_config": {
            "provider": "openai_video_proxy",
            "model_id": MODEL_ID,
            "display_name": "Seedance 2.5",
            "capabilities": ["text_to_video", "image_to_video", "video_editing"],
        },
    }
    if session_id:
        body["session_id"] = session_id
    primary = image_urls[0] if image_urls else None
    if primary:
        body["image_url"] = primary
        body["first_frame_url"] = primary
        body["image_mode"] = "first_last_frames" if mode == "first-last" else "omni_reference"
    if mode == "first-last":
        if not primary or not last_frame_url:
            raise RuntimeError("first-last mode needs BOTH a start frame (image 1) and --last-frame-url")
        body["last_frame_url"] = last_frame_url
    else:
        extra = [u for u in image_urls[1:] if u][: MAX_REFERENCE_IMAGES - 1]
        if extra:
            body["reference_images_urls"] = extra
    return body


def submit(session, body):
    r = session.request("POST", SUBMIT_PATH, body, ticketed=True, timeout=180)
    try:
        d = r.json()
    except Exception:
        raise RuntimeError(f"submit returned non-JSON (HTTP {r.status_code}): {r.text[:300]}")
    if not d.get("success") or not d.get("task_id"):
        code = ((d.get("data") or {}).get("code")) or ""
        err = d.get("error") or d.get("message") or "submit failed"
        if str(code).startswith("submit_ticket_"):
            raise RuntimeError(
                f"creaa refused the submit ({code}): the page could not mint a submit ticket. "
                f"Open the creaa tab, complete any captcha it shows, then retry."
            )
        raise RuntimeError(f"{err}{f' [{code}]' if code else ''}")
    data = d.get("data") or {}
    return d["task_id"], data


# creaa enforces "1 submitted or queued task at a time" SERVER-side, and a job started in the BROWSER (or by
# another machine on the same account) holds that slot invisibly to us. The backend's own FIFO can't see
# those, so a submit can be refused for reasons unrelated to our queue. That is a WAIT, not a failure.
SLOT_BUSY_RE = re.compile(r"only 1 submitted or queued task|1 (?:submitted|queued|in-progress|concurrent)", re.I)


def submit_with_slot_wait(session, body, max_wait_s=2700, interval=180):
    """Submit, backing off while the account's single unlimited slot is held by someone else, and pausing for
    a human whenever creaa raises a verification. Emits progress each round so the caller shows 'waiting'
    rather than a stall (and so the hub's timeout keeps refreshing).

    The interval is deliberately UNhurried (3 min). A tight retry loop against an account creaa is already
    suspicious of is what escalates an occasional verification prompt into a persistent one — and since the
    thing we are waiting for is a ~20-minute render, checking every 3 minutes costs nothing."""
    deadline = time.time() + max_wait_s
    while True:
        wait_for_human_verification(session)  # never submit into an open challenge
        try:
            return submit(session, body)
        except RuntimeError as e:
            # a challenge may have been raised BY this attempt — handle it before deciding anything else
            if captcha_state(session).get("open"):
                wait_for_human_verification(session)
                continue
            if not SLOT_BUSY_RE.search(str(e)):
                raise
            if time.time() >= deadline:
                raise RuntimeError(
                    f"the creaa account's single unlimited slot stayed busy for {int(max_wait_s / 60)} min "
                    f"(a task started elsewhere is still running) — last message: {e}"
                )
            event(event="progress", phase="waiting for the creaa slot (another task is running)")
            log(f"slot busy — retrying in {interval}s")
            time.sleep(interval)


# ---- workspace history ----
# creaa's gallery is built from universal-history records, and the SUBMIT endpoint does NOT create one:
# the web app POSTs the record itself right after submitting, then PUTs the result when the render lands.
# A client that skips this renders perfectly but is INVISIBLE in the user's own workspace — no card, no
# history row, nothing to cancel — which is exactly how a bridge-driven job becomes unobservable. All of
# it is best-effort: a bookkeeping failure must never fail a render that is otherwise fine.
HISTORY_PATH = "/api/universal-history/"


def create_history(session, task_id, body, session_id):
    """Register the render in the workspace so it shows up in the user's creaa gallery. Returns the
    history id (needed to close the record out later), or None if the bookkeeping call failed."""
    if not session_id:
        return None
    payload = {
        "service_type": "text_to_video",
        "input_data": {
            "prompt": body.get("description", ""),
            "optimized_prompt": body.get("description", ""),
            **({"image_url": body["image_url"]} if body.get("image_url") else {}),
        },
        "parameters": {
            "model_id": body.get("model_id"), "model": body.get("model_id"),
            "model_display_name": "Seedance 2.5",
            "duration": body.get("duration"), "fps": body.get("fps"),
            "resolution": body.get("resolution"), "aspect_ratio": body.get("aspect_ratio"),
        },
        "session_id": session_id,
        "status": "pending",
        "task_id": task_id,
    }
    try:
        # NOT ticketed: compatibleFetch re-stringifies the body for this endpoint, which arrives as a
        # JSON *string* and 422s ("Input should be a valid dictionary"). Only the generate endpoint needs
        # the submit ticket; history bookkeeping is a plain cookie-authed call.
        r = session.request("POST", HISTORY_PATH, payload, ticketed=False, timeout=60)
        d = r.json()
        hid = d.get("history_id") or d.get("id")
        log(f"workspace history record {hid} created for task {task_id}")
        return hid
    except Exception as e:
        log(f"WARNING: could not create the workspace history record ({e}) — the render still runs, "
            f"it just will not appear in the creaa gallery")
        return None


def close_history(session, history_id, video_url=None, thumbnail_url=None, error=None):
    """Close the record out so the gallery stops showing it as pending."""
    if not history_id:
        return
    try:
        if error:
            session.request("PUT", f"{HISTORY_PATH}{history_id}/failed",
                            {"service_type": "text_to_video", "error_message": str(error)[:500]},
                            ticketed=False, timeout=60)
        else:
            session.request("PUT", f"{HISTORY_PATH}{history_id}/result",
                            {"service_type": "text_to_video",
                             "output_data": {"video_url": video_url, "thumbnail_url": thumbnail_url}},
                            ticketed=False, timeout=60)
    except Exception as e:
        log(f"WARNING: could not close the workspace history record ({e})")


def poll(session, task_id, on_status=None, timeout=10800, interval=10, queued_since=None):
    """Poll until terminal. creaa's unlimited lane queues ~15 min BEFORE progress moves off 0, so the
    default timeout is generous; the caller's abort is what really bounds this."""
    end = time.time() + timeout
    started = time.time()
    last_key = None
    while time.time() < end:
        r = session.request("GET", TASK_PATH.format(task_id=task_id), timeout=60)
        try:
            d = r.json()
        except Exception:
            time.sleep(interval)
            continue
        p = d.get("payload") or {}
        st = (d.get("status") or p.get("status") or "").upper()
        prog = d.get("progress")
        if prog is None:
            prog = p.get("progress") or 0
        # a promo job sits at progress 0 for the whole queue wait — surface that as a distinct phase so the
        # UI shows "queued (~15m)" instead of a stalled-looking 0%.
        delayed = bool(p.get("delayed"))
        pct = int(round(float(prog) * 100)) if float(prog) <= 1 else int(prog)
        # The queue wait is the dominant, most variable part of a creaa render and it is NOT a constant:
        # observed 900s on four runs then 1740s on the next. Show the ETA creaa announced alongside how long
        # we have actually waited, so the node makes drift visible instead of hiding it behind "queued".
        eta_s = int(p.get("estimated_wait_seconds") or p.get("delay_seconds") or 0)
        waited_s = int(time.time() - queued_since) if queued_since else 0
        if delayed and pct == 0:
            phase = f"queued · creaa · ETA ~{eta_s // 60}m · waited {waited_s // 60}m"
        else:
            phase = f"rendering {pct}%"
        # tick on elapsed MINUTE too, not just status/percent — otherwise a queued job emits nothing for
        # half an hour and both the node and the hub's timeout see a job that looks dead.
        # Tick on elapsed minute in EVERY phase, not just the queue. A render that sits at the same percent
        # for a long stretch would otherwise emit nothing, and the hub's job timer — which only resets on a
        # progress event — would kill a job that is alive and simply slow.
        key = f"{st}:{pct}:{int(time.time() - started) // 60}"
        if key != last_key:
            last_key = key
            if on_status:
                on_status(phase, st, pct, p, eta_s, waited_s)
        if st in ("SUCCEEDED", "SUCCESS", "COMPLETED"):
            return p
        if st in ("FAILED", "ERROR", "CANCELED", "CANCELLED"):
            raise RuntimeError(p.get("error") or f"creaa task {st.lower()}")
        time.sleep(interval)
    raise TimeoutError(f"creaa task {task_id} did not finish within {timeout}s")


def main():
    global _JSON_OUT
    ap = argparse.ArgumentParser(description="creaa.ai Seedance 2.5 i2v over CDP")
    ap.add_argument("--prompt", default="")
    ap.add_argument("--image-url", action="append", default=[],
                    help="PUBLIC url of a reference image; repeat for ordered refs (image 1 = primary/start)")
    ap.add_argument("--last-frame-url", default=None, help="first-last mode: END frame url")
    ap.add_argument("--mode", default="reference", choices=["reference", "first-last"])
    ap.add_argument("--duration", type=int, default=5)
    ap.add_argument("--aspect", default="16:9")
    ap.add_argument("--resolution", default="720p")
    ap.add_argument("--session-id", default=None,
                    help="creaa workspace session id — makes the job visible in the user's own UI")
    ap.add_argument("--resume-task", default=None, help="skip submit; poll 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="creaa", help="substring of the logged-in tab's URL")
    # The bridge spawns this client with CDP_HTTP in the environment (it does NOT pass --host/--port), so
    # that has to win over the defaults — otherwise a Chrome on a non-standard debug port is unreachable.
    # Explicit --cdp/--host/--port still override, for running this by hand.
    ap.add_argument("--cdp", default=None, help="full debug-Chrome endpoint, e.g. http://localhost:9223 (else $CDP_HTTP)")
    ap.add_argument("--host", default=None)
    ap.add_argument("--port", type=int, default=None)
    ap.add_argument("--solve-captcha", action="store_true", help="run the captcha solver (called by Fly)")
    ap.add_argument("--execute-cdp-flow", action="store_true", help="execute CDP flow steps from stdin (called by bridge)")
    ap.add_argument("--flow-data", type=str, help="flow JSON data as string (fallback if stdin is not available)")
    args = ap.parse_args()

    if args.emit_json:
        # machine stream on the real stdout; every human print goes to stderr from here on
        _JSON_OUT = sys.stdout
        sys.stdout = sys.stderr

    cdp_http = (args.cdp
                or (f"http://{args.host or 'localhost'}:{args.port}" if args.port or args.host else None)
                or os.environ.get("CDP_HTTP")
                or "http://localhost:9222")
    # Every endpoint is built as f"{cdp_http}/json/...", so a trailing slash — which anyone pasting a URL
    # will naturally include, and which a browser address bar adds for you — produced "…:9222//json/version"
    # and a rejection that read exactly like "wrong port". Normalize instead of blaming the operator.
    cdp_http = cdp_http.rstrip("/")
    try:
        session = BrowserSession(cdp_http, args.match)
    except Exception as e:
        # Name the endpoint we actually tried. The common failure is not "no tab" at all but the WRONG
        # PORT: something else answering on 9222 (on Windows, IP Helper squats it) closes the connection,
        # which surfaces as an opaque "Remote end closed connection without response". Without the
        # endpoint in the message there is nothing in the node's error to point at the cause.
        if isinstance(e, NoMatchingTab):
            # reachability is FINE here — saying otherwise sends people chasing the port for nothing
            reason, hint = "no-tab", "the debug port is working; only the tab is missing."
        else:
            reason = "no-cdp"
            hint = (f"could not reach a debug Chrome at {cdp_http}. Start Chrome with "
                    f"--remote-debugging-port and pass the matching --cdp (the bridge's --cdp flag). "
                    f"If that port looks taken, another process may hold it — try another.")
        event(event="error", reason=reason, detail=f"{e} — {hint}", cdp=cdp_http, fatal=True)
        log(f"ERROR: {e}\n{hint}")
        return 2

    # Handle captcha solving (called by Fly bridge orchestration)
    if args.solve_captcha:
        try:
            solve_captcha_local()
            print("success")
            return 0
        except Exception as e:
            log(f"ERROR: captcha solve failed: {e}")
            return 1

    # Handle CDP flow execution (called by Fly via bridge orchestration)
    if args.execute_cdp_flow:
        try:
            flow_json = sys.stdin.read() if not sys.stdin.isatty() else args.flow_data or '{}'
            flow = json.loads(flow_json)
            script_dir = os.path.dirname(os.path.abspath(__file__ or ""))
            cdp_flow_script = os.path.join(script_dir, "cdp_flow.py")
            # Write flow to temp file and invoke cdp_flow.py
            import tempfile
            with tempfile.NamedTemporaryFile(mode='w', suffix='.json', delete=False) as f:
                json.dump(flow, f)
                temp_flow = f.name
            try:
                result = subprocess.run(
                    [sys.executable, cdp_flow_script, "--flow", temp_flow],
                    timeout=60,
                    capture_output=True,
                    text=True,
                    env={**os.environ, "CDP_HTTP": os.environ.get("CDP_HTTP", "http://localhost:9222")}
                )
                if result.returncode == 0:
                    print("success")
                    print(result.stdout, file=sys.stderr)
                    return 0
                else:
                    print(f"error", file=sys.stderr)
                    print(result.stderr, file=sys.stderr)
                    return 1
            finally:
                try:
                    os.unlink(temp_flow)
                except:
                    pass
        except Exception as e:
            log(f"ERROR: CDP flow execution failed: {e}")
            return 1

    try:
        history_id = None
        queued_since = time.time()  # when THIS process started waiting; resets on a resume, which is fine
        if args.resume_task:
            task_id = args.resume_task
            event(event="progress", phase="resuming", taskId=task_id)
        else:
            res = resolve_resolution(args.aspect, args.resolution)
            validate_free_lane(args.duration, res)
            if args.duration > RELIABLE_MAX_SEC:  # free, but the length will not be what was asked for
                log(f"WARNING: creaa ignores durations above ~{RELIABLE_MAX_SEC}s — {args.duration}s will "
                    f"most likely return a 30s clip. Do not trust the output length.")
                event(event="progress", phase=f"warning: {args.duration}s may return a 30s clip",
                      durationUnreliable=True)
            # explicit --session-id wins; otherwise use the tab's own workspace so the render is visible there
            session_id = args.session_id or getattr(session, "workspace_session", None)
            body = build_body(args.prompt, args.image_url, args.last_frame_url,
                              args.duration, args.aspect, res, args.mode, session_id)
            event(event="progress", phase="submitting")
            task_id, data = submit_with_slot_wait(session, body)
            # register it in the workspace so the user can SEE the render in their own creaa gallery
            history_id = create_history(session, task_id, body, session_id)
            # the queue wait is announced up front — surface it so the UI can say ~15m rather than look stuck
            event(event="progress", phase="queued · creaa", taskId=task_id,
                  delaySeconds=data.get("delay_seconds"), unlimited=data.get("seedance2_promo_unlimited"))
            log(f"task {task_id} submitted (delay {data.get('delay_seconds')}s)")

        p = poll(session, task_id, queued_since=queued_since,
                 on_status=lambda phase, st, pct, pay, eta_s, waited_s: event(
                     event="progress", phase=phase, taskId=task_id, percent=pct,
                     etaSeconds=eta_s, waitedSeconds=waited_s))
        video_url = p.get("video_url") or p.get("url")
        if not video_url:
            raise RuntimeError("task succeeded but carried no video_url")
        # A charged run means we fell out of the unlimited pass — surface it loudly rather than silently
        # spending the user's credits.
        if p.get("credit_deducted"):
            log(f"WARNING: creaa CHARGED credits for this render (cost {p.get('credits_cost')})")
        close_history(session, history_id, video_url=video_url, thumbnail_url=p.get("thumbnail_url"))
        event(event="done", taskId=task_id, videoUrl=video_url,
              thumbnailUrl=p.get("thumbnail_url"), creditDeducted=bool(p.get("credit_deducted")))
        log(f"done → {video_url}")
        return 0
    except Exception as e:
        close_history(session, locals().get("history_id"), error=e)  # don't leave it pending in the gallery
        event(event="error", reason=str(e)[:400], fatal=True)
        log(f"ERROR: {e}")
        return 1
    finally:
        try:
            session.close()
        except Exception:
            pass


if __name__ == "__main__":
    sys.exit(main())
