From 5fd5ec8df3c28ffc53508c53f4f68a07ddcf60d7 Mon Sep 17 00:00:00 2001 From: Hoang Tuan Nguyen Date: Sun, 6 Sep 2026 22:49:49 +0700 Subject: [PATCH 1/6] feat(flow): speak Flow's batchexecute API on flow.google.com MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Flow moved off labs.google in September 2026 and stopped minting the `Bearer ya29.…` the aisandbox-pa REST API needed. The rewritten frontend signs every call with the session cookie plus a per-page `at` token against one batchexecute endpoint, so nothing can be replayed from here: the agent builds the envelope, the extension runs it in the Flow tab. flow_batch.py holds the codec, the request builders and the readers, ported from the flowgen bridge along with the traps that cost it hours — the poll's third uuid is the scene and not the media, "Media not found." is a complaint rather than a verdict, the listing decides where the poll only hints, a media id arrives before the clip is fetchable, and image aspect (1 = square) does not share video aspect's encoding (1 = portrait). flow_client keeps its old surface and answers in the old REST shapes, so the worker's parsers, the operation poller and the DB writers never learn which transport ran. The legacy path stays behind USE_BATCH_RPC=0 as a post-mortem tool. Upscale, r2v and start+end-frame chaining have no captured payload and now fail as UNSUPPORTED_ON_BATCH_API instead of reaching for the dead bearer; FLOW_ALLOW_DEGRADED=1 drops the latter two to plain i2v. refresh_project_urls stops being a stub — the media rpc re-signs a stored id, which is exactly what /fk-refresh-urls always wanted. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015JkgMuy89iJ9qc1NkRdPk5 --- agent/api/flow.py | 10 +- agent/api/models.py | 4 + agent/config.py | 27 ++ agent/models.json | 10 + agent/services/flow_batch.py | 514 +++++++++++++++++++++++ agent/services/flow_client.py | 593 +++++++++++++++++++++++++-- tests/unit/test_flow_batch.py | 233 +++++++++++ tests/unit/test_flow_client_batch.py | 395 ++++++++++++++++++ 8 files changed, 1753 insertions(+), 33 deletions(-) create mode 100644 agent/services/flow_batch.py create mode 100644 tests/unit/test_flow_batch.py create mode 100644 tests/unit/test_flow_client_batch.py diff --git a/agent/api/flow.py b/agent/api/flow.py index 1373e113b..b9d8b141c 100644 --- a/agent/api/flow.py +++ b/agent/api/flow.py @@ -2,6 +2,7 @@ from fastapi import APIRouter, HTTPException from pydantic import BaseModel from typing import Optional +from agent.config import USE_BATCH_RPC, FLOW_PROJECT_ID, FLOW_ALLOW_DEGRADED from agent.services.flow_client import get_flow_client router = APIRouter(prefix="/flow", tags=["flow"]) @@ -61,10 +62,17 @@ class EditImageRequest(BaseModel): @router.get("/status") async def extension_status(): - """Check if extension is connected.""" + """Extension health, and which transport it is being asked to speak. + + `flow_key_present` is a legacy-path signal: the batchexecute path has no + bearer token at all, so false is expected there rather than a fault. + """ client = get_flow_client() return { "connected": client.connected, + "transport": "batch" if USE_BATCH_RPC else "legacy_rest", + "flow_project_id": FLOW_PROJECT_ID or None, + "allow_degraded": FLOW_ALLOW_DEGRADED, "flow_key_present": client._flow_key is not None, } diff --git a/agent/api/models.py b/agent/api/models.py index aea93127e..b4fb05bab 100644 --- a/agent/api/models.py +++ b/agent/api/models.py @@ -32,6 +32,7 @@ def _reload_config(data: dict): config.UPSCALE_MODELS.update(data["upscale_models"]) config.IMAGE_MODELS.clear() config.IMAGE_MODELS.update(data["image_models"]) + config.DEFAULT_IMAGE_MODEL = data.get("default_image_model", "NANO_BANANA_PRO") @router.get("") @@ -57,6 +58,9 @@ async def patch_models(body: dict): """ current = _read_models() + if "default_image_model" in body: + current["default_image_model"] = body["default_image_model"] + # Deep merge: only update keys that are provided for section in ("video_models", "image_models", "upscale_models"): if section not in body: diff --git a/agent/config.py b/agent/config.py index d180bdb78..f2c1c567d 100644 --- a/agent/config.py +++ b/agent/config.py @@ -16,10 +16,34 @@ WS_PORT = int(os.environ.get("WS_PORT", "9222")) # ─── Google Flow API ──────────────────────────────────────── +# Legacy REST host. Flow moved to flow.google.com in September 2026 and stopped +# minting the `Bearer ya29.…` this host needs, so these are only reachable with +# USE_BATCH_RPC=0 on a browser profile that still has an old token. GOOGLE_FLOW_API = "https://aisandbox-pa.googleapis.com" GOOGLE_API_KEY = os.environ.get("GOOGLE_API_KEY", "AIzaSyBtrm0o5ab1c-Ec8ZuLcGt3oJAA5VWt3pY") RECAPTCHA_SITE_KEY = os.environ.get("RECAPTCHA_SITE_KEY", "6LdsFiUsAAAAAIjVDZcuLhaHiDn5nnHVXVRQGeMV") +# ─── Flow batchexecute (the current path) ─────────────────── +# Every call is signed in the page with the session cookie plus a per-page `at` +# token, so the extension runs it inside a signed-in flow.google.com tab. Set +# USE_BATCH_RPC=0 only to fall back to the dead REST path for a post-mortem. +USE_BATCH_RPC = os.environ.get("USE_BATCH_RPC", "1") == "1" + +# The Flow project every RPC is scoped to. Project creation went with the old +# labs.google tRPC endpoint, so a project is made once in the Flow UI and its +# uuid pinned here; POST /api/projects falls back to it when no id is given. +FLOW_PROJECT_ID = os.environ.get("FLOW_PROJECT_ID", "") + +# Capabilities whose payloads were never captured off the new UI (4K upscale, +# reference-to-video, start+end-frame chaining) fail loudly by default. With +# this on, the two that have a sane fallback degrade instead: chaining and r2v +# both drop to plain i2v off the start frame. Upscale has no fallback. +FLOW_ALLOW_DEGRADED = os.environ.get("FLOW_ALLOW_DEGRADED", "0") == "1" + +# The tier no longer picks a model — aspect is its own slot and the model names +# are fixed — so it is only carried for the DB column and the dashboard. +DEFAULT_PAYGATE_TIER = os.environ.get("DEFAULT_PAYGATE_TIER", "PAYGATE_TIER_TWO") + # ─── Worker ────────────────────────────────────────────────── POLL_INTERVAL = int(os.environ.get("POLL_INTERVAL", "5")) VIDEO_POLL_INTERVAL = int(os.environ.get("VIDEO_POLL_INTERVAL", "10")) # polling interval for video/upscale status @@ -37,6 +61,9 @@ VIDEO_MODELS = _MODELS["video_models"] UPSCALE_MODELS = _MODELS["upscale_models"] IMAGE_MODELS = _MODELS["image_models"] +# Nickname from image_models. The batch path accepts GEM_PIX_2 (Nano Banana Pro) +# and NARWHAL (Banana 2) and rejects everything else. +DEFAULT_IMAGE_MODEL = _MODELS.get("default_image_model", "NANO_BANANA_PRO") # ─── API Endpoints ─────────────────────────────────────────── ENDPOINTS = { diff --git a/agent/models.json b/agent/models.json index 6bd4242a7..9a45e7351 100644 --- a/agent/models.json +++ b/agent/models.json @@ -36,5 +36,15 @@ "image_models": { "NANO_BANANA_PRO": "GEM_PIX_2", "NANO_BANANA_2": "NARWHAL" + }, + "default_image_model": "NANO_BANANA_PRO", + "batch_video_models": { + "_comment": "Wire names the flow.google.com batchexecute path accepts. video_models above are REST-era keys; resolve_video_model() folds them onto these (anything asking for 'ultra' lands on the ultra model).", + "accepted": [ + "veo_3_1_i2v_lite_low_priority", + "veo_3_1_i2v_lite", + "veo_3_1_i2v_s_fast_ultra" + ], + "default": "veo_3_1_i2v_lite_low_priority" } } diff --git a/agent/services/flow_batch.py b/agent/services/flow_batch.py new file mode 100644 index 000000000..f01116320 --- /dev/null +++ b/agent/services/flow_batch.py @@ -0,0 +1,514 @@ +"""Flow's batchexecute API — the transport Flow moved to in September 2026. + +The old world was a REST call to ``aisandbox-pa.googleapis.com`` carrying a +``Bearer ya29.…`` the extension sniffed off the page. That token no longer +exists: the rewritten frontend on ``flow.google.com`` signs every call with the +session cookie plus a per-page ``at`` token, against a single batchexecute +endpoint. Nothing can be replayed from outside the browser — a generate call +also carries a **single-use** reCAPTCHA token, and a replayed one comes back +``PUBLIC_ERROR_UNUSUAL_ACTIVITY``. + +So this module never touches the network. It builds request envelopes and reads +responses; issuing the request is the job of the Chrome extension, which runs +the payload inside the Flow tab (see ``FlowClient.batch_rpc``). + +The wire format is Google's usual batchexecute: + + f.req = [[[rpcid, "", null, "generic"]]] + +and the response is a ``)]}'`` sentinel followed by length-prefixed chunks of +``["wrb.fr", rpcid, "", …]`` envelopes. + +Ported from postforge/bridges/flowgen/flow_batch.py — see its MIGRATION.md for +the traps behind each of the comments below. +""" +from __future__ import annotations + +import json +import random +import re +import uuid +from dataclasses import dataclass +from typing import Any, Optional + +BATCH_PATH = "/_/AiSandboxAngularFrontend/data/batchexecute" +MEDIA_HOST = "flow-content.google" + +RPC_GEN_IMAGE = "ogiZ0b" +RPC_GEN_VIDEO = "eb1hJf" +RPC_OPERATION = "jwpduf" +RPC_PROJECT_MEDIA = "Zzl0ze" +RPC_MEDIA = "as29s" +RPC_UPLOAD_IMAGE = "maseQ" + +CAPTCHA_IMAGE = "IMAGE_GENERATION" +CAPTCHA_VIDEO = "VIDEO_GENERATION" + +#: The extension substitutes a freshly minted reCAPTCHA token for this marker. +#: It has to be a placeholder rather than a real token because the mint has to +#: happen in the page, moments before the request leaves. +CAPTCHA_SLOT = "__CAPTCHA__" + +#: Wire names this path accepts. Everything else is rejected outright by Flow. +#: ``GEM_PIX_2`` is Nano Banana Pro, ``NARWHAL`` is Banana 2. Flow Kit uses Pro +#: by default (see agent/models.json), which is also what the new path defaults +#: to; a caller that wants Banana 2 has to name it. +IMAGE_MODELS = {"GEM_PIX_2", "NARWHAL"} +IMAGE_MODEL = "GEM_PIX_2" + +#: The nicknames models.json speaks, resolved to wire names. +IMAGE_MODEL_BY_NICKNAME = {"NANO_BANANA_PRO": "GEM_PIX_2", "NANO_BANANA_2": "NARWHAL"} + +#: Image aspect ratios, measured by generating one of each and reading the +#: JPEG header. This slot was mistaken for a variant count at first — 1 means +#: square, which is why a `count=1` request looked like it was working. +ASPECT_SQUARE = 1 # 1024x1024 +ASPECT_PORTRAIT = 2 # 768x1376 (9:16) +ASPECT_LANDSCAPE = 3 # 1376x768 (16:9) +ASPECT_PORTRAIT_4_3 = 4 # 896x1200 (3:4) +ASPECT_LANDSCAPE_4_3 = 5 # 1200x896 (4:3) + +#: The names the REST payload used, so callers can keep speaking them. +ASPECT_BY_NAME = { + "IMAGE_ASPECT_RATIO_SQUARE": ASPECT_SQUARE, + "IMAGE_ASPECT_RATIO_PORTRAIT": ASPECT_PORTRAIT, + "IMAGE_ASPECT_RATIO_LANDSCAPE": ASPECT_LANDSCAPE, + "IMAGE_ASPECT_RATIO_PORTRAIT_FOUR_THREE": ASPECT_PORTRAIT_4_3, + "IMAGE_ASPECT_RATIO_LANDSCAPE_FOUR_THREE": ASPECT_LANDSCAPE_4_3, +} + +#: Video models this path accepts. The REST-era map was keyed by +#: [tier][quality][aspect] and carried `…_portrait` / `…_fl` / `…_relaxed` +#: variants; those are gone — aspect is its own slot now, and the suffixed +#: names are rejected. +VIDEO_MODEL = "veo_3_1_i2v_lite_low_priority" +VIDEO_MODELS = { + "veo_3_1_i2v_lite_low_priority", + "veo_3_1_i2v_lite", + "veo_3_1_i2v_s_fast_ultra", +} + +#: Video aspect, and note it does NOT share the image encoding: here 1 is +#: portrait, where for an image 1 is square. Measured by rendering one of each +#: from the same portrait still — 720x1280 against 1280x720. +VIDEO_ASPECT_PORTRAIT = 1 +VIDEO_ASPECT_LANDSCAPE = 2 + +VIDEO_ASPECT_BY_NAME = { + "VIDEO_ASPECT_RATIO_PORTRAIT": VIDEO_ASPECT_PORTRAIT, + "VIDEO_ASPECT_RATIO_LANDSCAPE": VIDEO_ASPECT_LANDSCAPE, +} + +#: `CAE` is the operation's terminal state. Anything else means still working. +STATUS_DONE = "CAE" + +#: Outcome codes seen in the operation's status block. Code 4 carries a +#: message like "Media not found." — but it is NOT a verdict: jobs that report +#: it still finish, and the finished media shows up in the project listing +#: seconds later. Treat it as something to quote on a timeout, never as a +#: reason to stop waiting. +OUTCOME_OK = 3 +OUTCOME_COMPLAINT = 4 + +#: Surface id the web client stamps on every call. Constant in every capture. +SURFACE_ID = 22 + +#: Crop box on the reference image, verbatim from the UI when nothing was +#: reframed by hand: a hair inside the edges, spanning 128/129 of the frame. +FULL_FRAME_CROP = [None, 0.0038759689922481244, 1, 0.9961240310077519] + +#: A reference image, as the UI sends it: the media id FIRST and a type flag +#: four slots later. Probing never found this — the id sat in the wrong +#: position, the payload was accepted, and the picture quietly ignored it. +REF_TYPE_IMAGE = 1 + + +class RpcError(RuntimeError): + """A batchexecute envelope came back with an error slot instead of data.""" + + def __init__(self, rpcid: str, detail: Any): + super().__init__(f"{rpcid} failed: {detail!r}") + self.rpcid = rpcid + self.detail = detail + + +class FlowBatchError(RuntimeError): + """The call succeeded but the payload did not hold what we came for.""" + + +@dataclass(frozen=True) +class RpcResult: + rpcid: str + data: Any + error: Any = None + + @property + def ok(self) -> bool: + return self.error is None + + +@dataclass(frozen=True) +class GeneratedImage: + media_id: str + url: str + + +@dataclass(frozen=True) +class Operation: + operation_id: str + project_id: Optional[str] + status: Optional[str] + error: Optional[str] = None + + @property + def done(self) -> bool: + return self.status == STATUS_DONE + + @property + def complained(self) -> bool: + """The poll grumbled. Observed to be survivable — check the listing.""" + return self.error is not None + + +@dataclass(frozen=True) +class MediaUrls: + media_id: str + video: Optional[str] = None + image: Optional[str] = None + + +# ── model / aspect resolvers ───────────────────────────────────────────────── + +def resolve_image_model(key: Optional[str]) -> str: + """Nickname or wire name in, wire name out; anything unknown coerces.""" + if isinstance(key, str): + if key in IMAGE_MODEL_BY_NICKNAME: + return IMAGE_MODEL_BY_NICKNAME[key] + if key in IMAGE_MODELS: + return key + return IMAGE_MODEL + + +def resolve_video_model(key: Optional[str]) -> str: + """Map a REST-era model key onto one the batch path accepts. + + The old keys encoded tier, quality, aspect and chaining in the name + (``veo_3_1_i2v_s_fast_ultra_relaxed``, ``…_portrait``, ``…_fl``). Aspect + and chaining are their own slots now and the suffixed names are rejected, + so the tier/quality intent is all that survives: anything that asked for + "ultra" gets the ultra model, anything else lands on the lite default. + """ + if isinstance(key, str): + if key in VIDEO_MODELS: + return key + if "ultra" in key: + return "veo_3_1_i2v_s_fast_ultra" + if "lite_low_priority" in key: + return "veo_3_1_i2v_lite_low_priority" + if "lite" in key: + return "veo_3_1_i2v_lite" + return VIDEO_MODEL + + +def resolve_aspect(aspect: Any) -> int: + """Take either the wire value or the REST-era name.""" + if isinstance(aspect, int): + return aspect + try: + return ASPECT_BY_NAME[aspect] + except KeyError: + raise ValueError( + f"unknown aspect {aspect!r} — use one of {sorted(ASPECT_BY_NAME)} or 1-5" + ) from None + + +def resolve_video_aspect(aspect: Any) -> int: + if isinstance(aspect, int): + if aspect not in (VIDEO_ASPECT_PORTRAIT, VIDEO_ASPECT_LANDSCAPE): + # 3 is a perfectly good IMAGE aspect and a meaningless video one + raise ValueError(f"video aspect must be 1 or 2, got {aspect}") + return aspect + try: + return VIDEO_ASPECT_BY_NAME[aspect] + except KeyError: + raise ValueError( + f"unknown video aspect {aspect!r} — use one of " + f"{sorted(VIDEO_ASPECT_BY_NAME)} or 1-2" + ) from None + + +# ── envelope codec ─────────────────────────────────────────────────────────── + +def build_envelope(rpcid: str, inner: Any) -> str: + """Wrap an inner payload as the ``f.req`` string batchexecute expects.""" + return json.dumps( + [[[rpcid, json.dumps(inner, separators=(",", ":"), ensure_ascii=False), None, "generic"]]], + separators=(",", ":"), + ensure_ascii=False, + ) + + +def parse_envelope(text: str) -> list[RpcResult]: + """Unwrap the `)]}'` sentinel and the length-prefixed chunks. + + The chunk lengths count characters, but a payload can disagree with them by + a byte or two once escapes are involved, so the JSON is decoded by scanning + rather than by trusting the prefix. + """ + if not text: + return [] + body = text.split("\n", 1)[1] if text.startswith(")]}'") else text + decoder = json.JSONDecoder() + results: list[RpcResult] = [] + index = 0 + while index < len(body): + start = body.find("[", index) + if start == -1: + break + try: + chunk, consumed = decoder.raw_decode(body[start:]) + except json.JSONDecodeError: + # step past this `[` and keep scanning — a chunk boundary landing + # mid-token should not cost us the envelopes that follow it + index = start + 1 + continue + index = start + consumed + for entry in chunk if isinstance(chunk, list) else []: + if not isinstance(entry, list) or not entry or entry[0] != "wrb.fr": + continue + rpcid = entry[1] if len(entry) > 1 else "?" + payload = entry[2] if len(entry) > 2 else None + if payload is None: + # index 5 is the error slot; it is `[5]`-style codes, not text + results.append(RpcResult(rpcid, None, entry[5] if len(entry) > 5 else True)) + continue + results.append( + RpcResult(rpcid, json.loads(payload) if isinstance(payload, str) else payload) + ) + return results + + +def first_payload(text: str, rpcid: str) -> Any: + """The payload of the first matching envelope, or raise what went wrong.""" + results = parse_envelope(text) + for result in results: + if result.rpcid != rpcid: + continue + if not result.ok: + raise RpcError(rpcid, result.error) + return result.data + raise FlowBatchError(f"no {rpcid} envelope in response ({len(results)} others)") + + +# ── request builders ───────────────────────────────────────────────────────── + +def _client_uuid() -> str: + """Client-side request ids. The UI sends them upper-case; match it.""" + return str(uuid.uuid4()).upper() + + +def _context(project_id: str) -> list: + """The surface/project/captcha envelope every generate call repeats.""" + return [None, SURFACE_ID, None, None, None, project_id, None, None, None, None, + [CAPTCHA_SLOT, 1]] + + +def _reference(media_id: str) -> list: + return [media_id, None, None, None, REF_TYPE_IMAGE] + + +def image_request(prompt: str, project_id: str, count: int = 1, + aspect: Any = ASPECT_SQUARE, seed: Optional[int] = None, + prompts: Optional[list[str]] = None, + model: str = IMAGE_MODEL, + ref_media_ids: Optional[list[str]] = None) -> str: + """One request item per variant, exactly as the REST payload did it. + + There is no "how many" field: Flow returns one image per item in the list, + so `count` replicates the item under fresh seeds. `ref_media_ids` conditions + the result on images already in the project — this is what keeps a character + the same person from beat to beat. + """ + ratio = resolve_aspect(aspect) + base = seed if seed is not None else random.randint(1, 10**9) + items = [] + for index in range(max(1, count)): + text = prompts[index] if prompts and index < len(prompts) else prompt + refs = [_reference(mid) for mid in (ref_media_ids or [])] or None + items.append([None, None, refs, base + index * 9973, ratio, model, None, + _context(project_id), [[[text]]], None, None, None, + _client_uuid(), _client_uuid()]) + return build_envelope(RPC_GEN_IMAGE, [None, items, 1, _context(project_id), + [_client_uuid()]]) + + +def video_request(prompt: str, project_id: str, source_media_id: str, + crop: Optional[list] = None, + aspect: Any = VIDEO_ASPECT_LANDSCAPE, + model: str = VIDEO_MODEL) -> str: + inner = [ + [[[None, None, [[[prompt]]]], model, resolve_video_aspect(aspect), None, + [None, source_media_id, None, None, None, + FULL_FRAME_CROP if crop is None else crop], + [None, None, None, None, _client_uuid(), _client_uuid()]]], + _context(project_id), + [_client_uuid(), 2], + ] + return build_envelope(RPC_GEN_VIDEO, inner) + + +def upload_request(image_b64: str, project_id: str, mime_type: str = "image/jpeg", + file_name: str = "upload.jpg") -> str: + """Put a local image into the project so it can be used as a reference. + + The bytes ride inside the RPC as plain base64 — no data: prefix, no separate + upload endpoint — and the call carries a captcha like a generate does. + """ + return build_envelope(RPC_UPLOAD_IMAGE, [ + _context(project_id), image_b64, mime_type, 1, None, None, None, None, + file_name, None, _client_uuid(), _client_uuid(), + ]) + + +def operation_request(operation_id: str) -> str: + return build_envelope(RPC_OPERATION, [None, None, [[operation_id]]]) + + +def project_media_request(project_id: str) -> str: + return build_envelope(RPC_PROJECT_MEDIA, [f"projects/{project_id}", None, None, None, [1]]) + + +def media_request(media_id: str) -> str: + return build_envelope(RPC_MEDIA, [media_id]) + + +# ── response readers ───────────────────────────────────────────────────────── + +def _walk_strings(node: Any): + if isinstance(node, str): + yield node + elif isinstance(node, list): + for item in node: + yield from _walk_strings(item) + + +def _walk_lists(node: Any): + if isinstance(node, list): + yield node + for item in node: + yield from _walk_lists(item) + + +def read_images(payload: Any) -> list[GeneratedImage]: + """Signed CDN urls come back inline on the image call — one per variant. + + The media id is read out of the url path rather than from a fixed index: + the url is the thing we actually need, and pairing them at the source keeps + a reshuffled response from mismatching ids to pictures. + """ + images: list[GeneratedImage] = [] + seen: set[str] = set() + for text in _walk_strings(payload): + if MEDIA_HOST + "/image/" not in text: + continue + media_id = text.split("/image/", 1)[1].split("?", 1)[0] + if media_id in seen: + continue + seen.add(media_id) + images.append(GeneratedImage(media_id=media_id, url=text)) + return images + + +def read_uploaded_media_id(payload: Any) -> str: + """`[[mediaId, projectId, operationId, "CAE", …]]` — the id is the handle a + later generate passes as a reference.""" + record = payload[0] if isinstance(payload, list) and payload else None + media_id = record[0] if isinstance(record, list) and record else None + if not isinstance(media_id, str) or not media_id: + raise FlowBatchError("upload response carried no media id") + return media_id + + +def read_operation(payload: Any) -> Operation: + """`[null, 50, [[opId, projectId, sceneId, status, …]]]`. + + Note the third uuid is the **scene**, not the media. Reading it as a media + id is what made every `as29s` lookup answer NOT_FOUND. + """ + records = payload[2] if isinstance(payload, list) and len(payload) > 2 else None + record = records[0] if isinstance(records, list) and records else None + if not isinstance(record, list) or not record: + raise FlowBatchError("operation payload carried no record") + return Operation( + operation_id=record[0], + project_id=record[1] if len(record) > 1 else None, + status=record[3] if len(record) > 3 else None, + error=read_operation_error(record), + ) + + +def read_operation_error(record: list) -> Optional[str]: + """The complaint attached to this operation, if it carries one. + + It hides in the detail block's status slot as + ``[4, [null, "Media not found."], ["Media not found."]]``. Measured + behaviour: an operation can report exactly that and still deliver a + finished 8-second clip, so this is a diagnostic string and nothing more. + """ + detail = record[5] if len(record) > 5 else None + if not isinstance(detail, list) or len(detail) <= 8: + return None + block = detail[8] + if not isinstance(block, list) or not block or block[0] != OUTCOME_COMPLAINT: + return None + for text in _walk_strings(block): + return text + return "operation failed without a message" + + +def find_media_id(payload: Any, operation_id: str) -> Optional[str]: + """Look an operation up in the project listing and take its media id. + + Entries look like + ``[opId, null, null, [title, created, null, null, mediaId, clientUuid, done], projectId]``. + """ + for node in _walk_lists(payload): + if len(node) < 4 or node[0] != operation_id: + continue + detail = node[3] + if isinstance(detail, list) and len(detail) > 4 and isinstance(detail[4], str): + return detail[4] + return None + + +#: The media slot in a listing entry, matched straight off the wire: a title, +#: a timestamp pair, two nulls, then the media id. Escaped or not, both forms +#: appear depending on whether the text has been through a JSON decode. +_MEDIA_SLOT = re.compile(r'null,null,\\?"([0-9a-fA-F-]{36})\\?"') + + +def find_media_id_in_text(text: str, operation_id: str) -> Optional[str]: + """Same lookup as :func:`find_media_id`, but on an unparsed listing. + + The project listing has no page size that shrinks it and grows with every + generation, so it will outrun whatever response cap is in place — and a + truncated tail cannot be JSON-decoded even though the entry we want is + sitting in it intact. Scanning the text finds it anyway. + """ + start = text.find(operation_id) + if start == -1: + return None + match = _MEDIA_SLOT.search(text, start, start + 800) + return match.group(1) if match else None + + +def read_media_urls(payload: Any, media_id: str) -> MediaUrls: + video = image = None + for text in _walk_strings(payload): + if not text.startswith("https://"): + continue + if MEDIA_HOST + "/video/" in text and video is None: + video = text + elif MEDIA_HOST + "/image/" in text and image is None: + image = text + return MediaUrls(media_id=media_id, video=video, image=image) diff --git a/agent/services/flow_client.py b/agent/services/flow_client.py index 8b7a78bee..4f12019df 100644 --- a/agent/services/flow_client.py +++ b/agent/services/flow_client.py @@ -1,8 +1,18 @@ """ -Flow Client — communicates with Google Flow API via Chrome extension WebSocket bridge. +Flow Client — communicates with Google Flow via the Chrome extension bridge. -Agent runs a WS server. Extension connects as client. Agent sends API requests, +Agent runs a WS server. Extension connects as client. Agent sends requests, extension executes them in browser context (residential IP, cookies, reCAPTCHA). + +Two transports live here. The current one is Flow's ``batchexecute`` endpoint on +flow.google.com, whose calls only a signed-in page can sign — the agent builds +the envelope, the extension runs it in the tab (see :mod:`agent.services.flow_batch`). +The old REST path against ``aisandbox-pa.googleapis.com`` is kept behind +``USE_BATCH_RPC=0``; it needs a ``Bearer ya29.…`` that Flow stopped minting in +the September 2026 migration, so it is a post-mortem tool, not a fallback. + +Both shape their answers the same way, so everything downstream — the worker's +parsers, the operation poller, the scene/character updaters — is transport-blind. """ import asyncio import json @@ -14,7 +24,11 @@ from agent.config import ( GOOGLE_FLOW_API, GOOGLE_API_KEY, ENDPOINTS, VIDEO_MODELS, UPSCALE_MODELS, IMAGE_MODELS, VIDEO_POLL_TIMEOUT, + USE_BATCH_RPC, FLOW_PROJECT_ID, FLOW_ALLOW_DEGRADED, + DEFAULT_PAYGATE_TIER, ) +from agent import config as _config +from agent.services import flow_batch as fb from agent.services.headers import random_headers logger = logging.getLogger(__name__) @@ -27,6 +41,13 @@ def __init__(self): self._extension_ws = None # Set by WS server when extension connects self._pending: dict[str, asyncio.Future] = {} self._flow_key: Optional[str] = None + # Per-operation poll state. `_operation_projects` says which project + # listing to look a finished media up in; `_operation_media` caches the + # id once the listing has it, so later rounds skip the listing entirely; + # `_operation_polls` counts rounds, to keep the listing off most of them. + self._operation_projects: dict[str, str] = {} + self._operation_media: dict[str, str] = {} + self._operation_polls: dict[str, int] = {} # WS stats self._ws_connect_count = 0 self._ws_disconnect_count = 0 @@ -130,7 +151,10 @@ async def _sync_tier(self): self._sync_in_progress = False _UUID_RE = __import__("re").compile(r'^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$') - _SAFE_URL_RE = __import__("re").compile(r'^https://(storage\.googleapis\.com|lh3\.googleusercontent\.com)/') + # flow-content.google is where the rewritten frontend serves media from; + # the other two are the pre-migration hosts, still seen on older media. + _SAFE_URL_RE = __import__("re").compile( + r'^https://(storage\.googleapis\.com|lh3\.googleusercontent\.com|flow-content\.google)/') async def _refresh_media_urls(self, urls: list[dict]): """Update scene/character URLs in DB from fresh TRPC-captured signed URLs. @@ -192,18 +216,65 @@ async def _refresh_media_urls(self, urls: list[dict]): await event_bus.emit("urls_refreshed", {"count": updated}) async def refresh_project_urls(self, project_id: str) -> dict: - """Refresh media URLs for a project. + """Re-sign every stored media url for a project. - Note: Google Flow's get_media API returns encoded content (base64), - not fresh signed URLs. URL refresh requires TRPC intercept from - the extension when the user opens the project in Chrome. - The video reviewer falls back to get_media content directly. + The batch path can do this properly: the media rpc answers a media id + with a freshly signed url, so we walk the project's scenes and entities + and refresh each id we hold. The legacy path could not — its media + endpoint returned base64 content rather than a url — so it still asks + the user to open the project in Chrome and let the intercept catch them. """ - logger.info("URL refresh requested for project %s — TRPC endpoint no longer available, " - "use extension passive intercept (open project in Chrome)", project_id[:12]) - return {"refreshed": 0, "found": 0, "note": "TRPC endpoint unavailable. " - "Video reviewer uses get_media fallback automatically. " - "For URL refresh, open the project in Google Flow in Chrome."} + if not USE_BATCH_RPC: + logger.info("URL refresh requested for project %s — legacy path has no " + "url-serving media endpoint", project_id[:12]) + return {"refreshed": 0, "found": 0, "note": "Legacy REST path: no URL refresh. " + "Open the project in Google Flow in Chrome and let the extension " + "intercept fresh URLs, or set USE_BATCH_RPC=1."} + + from agent.db import crud + + # (media_id, kind) -> the scene/character fields it should land in + targets: dict[tuple[str, str], list[tuple[str, str, str]]] = {} + + def want(media_id, kind, table, row_id, field): + if media_id and self._UUID_RE.match(media_id): + targets.setdefault((media_id, kind), []).append((table, row_id, field)) + + scenes = [] + for video in await crud.list_videos(project_id): + scenes.extend(await crud.list_scenes(video["id"])) + + for scene in scenes: + for prefix in ("vertical", "horizontal"): + want(scene.get(f"{prefix}_image_media_id"), "image", + "scene", scene["id"], f"{prefix}_image_url") + want(scene.get(f"{prefix}_video_media_id"), "video", + "scene", scene["id"], f"{prefix}_video_url") + want(scene.get(f"{prefix}_upscale_media_id"), "video", + "scene", scene["id"], f"{prefix}_upscale_url") + for char in await crud.get_project_characters(project_id): + want(char.get("media_id"), "image", "character", char["id"], "reference_image_url") + + refreshed = 0 + for (media_id, kind), fields in targets.items(): + try: + urls = await self._batch_media_urls(media_id) + except Exception as e: + logger.warning("Refresh failed for media %s: %s", media_id[:12], e) + continue + url = urls.video if kind == "video" else urls.image + if not url: + continue + for table, row_id, field in fields: + if table == "scene": + await crud.update_scene(row_id, **{field: url}) + else: + await crud.update_character(row_id, **{field: url}) + refreshed += 1 + + logger.info("Refreshed %d/%d media urls for project %s", + refreshed, len(targets), project_id[:12]) + return {"refreshed": refreshed, "found": len(targets)} async def _send(self, method: str, params: dict, timeout: float = 300) -> dict: """Send request to extension and wait for response. @@ -253,9 +324,427 @@ def _client_context(self, project_id: str, user_paygate_tier: str = "PAYGATE_TIE "userPaygateTier": user_paygate_tier, } + # ─── batchexecute transport ────────────────────────────── + # + # Flow's rewritten frontend signs every call with the session cookie plus a + # per-page `at` token, and a generate also carries a single-use reCAPTCHA. + # None of that can be replayed from here, so the agent builds the envelope + # and the extension runs it inside a signed-in flow.google.com tab. + + async def batch_rpc(self, rpcid: str, freq: str, + captcha_action: str | None = None, + match: str | None = None, + timeout: float = 300) -> dict: + """Run one batchexecute RPC in the Flow page. Returns the raw body. + + ``match`` asks the extension to cut the response down to an 800-byte + window around that string before handing it back. The project listing + is tens of megabytes for the one entry we want, and the cheapest place + to throw the rest away is inside the tab. + """ + params: dict = {"rpcid": rpcid, "freq": freq} + if captcha_action: + params["captchaAction"] = captcha_action + if match: + params["match"] = match + return await self._send("batch_rpc", params, timeout=timeout) + + async def _batch_payload(self, rpcid: str, freq: str, + captcha_action: str | None = None, + timeout: float = 300): + """One RPC, unwrapped to its inner payload. Raises on anything else.""" + result = await self.batch_rpc(rpcid, freq, captcha_action, timeout=timeout) + if result.get("error"): + raise fb.FlowBatchError(f"{rpcid}: {result['error']}") + return fb.first_payload(result.get("data") or "", rpcid) + + def _batch_project_id(self, project_id: str) -> str: + """The Flow project an RPC is scoped to. + + Flow Kit stores the Flow project uuid as the local project id, but a + few call sites pass "0" or "" for project-less work; those fall back to + the pinned FLOW_PROJECT_ID. + """ + if project_id and self._UUID_RE.match(str(project_id)): + return str(project_id) + if FLOW_PROJECT_ID: + return FLOW_PROJECT_ID + raise fb.FlowBatchError( + "NO_FLOW_PROJECT: every batchexecute call is scoped to a Flow project. " + "Create one in the Flow UI and pin its uuid as FLOW_PROJECT_ID." + ) + + def _batch_image_model(self, override: str | None = None) -> str: + # Read through the module: PATCH /api/models hot-reloads both of these. + nickname = _config.DEFAULT_IMAGE_MODEL + return fb.resolve_image_model( + override or _config.IMAGE_MODELS.get(nickname) or nickname + ) + + def _batch_video_model(self, tier: str, gen_type: str, aspect_ratio: str) -> str: + legacy = VIDEO_MODELS.get(tier, {}).get(gen_type, {}).get(aspect_ratio) + return fb.resolve_video_model(legacy) + + def _remember_operation(self, operation_id: str, project_id: str): + """Which project an operation belongs to — the listing lookup needs it. + + A poll record usually carries the project id, but old operations decay + to a bare id, so keep our own note. Bounded: this is a cache, and the + pinned project is always a workable fallback. + """ + if not operation_id: + return + if len(self._operation_projects) > 512: + self._operation_projects.clear() + self._operation_media.clear() + self._operation_polls.clear() + self._operation_projects[operation_id] = project_id + # ─── High-level API Methods ────────────────────────────── + def flow_project_id(self, requested: str | None = None) -> str | None: + """The Flow project to attach a new Flow Kit project to, if any. + + Project creation went with the labs.google tRPC endpoint the migration + unauthenticated, so on the batch path a project is made once in the + Flow UI and its uuid supplied here or pinned as FLOW_PROJECT_ID. + """ + if requested and self._UUID_RE.match(requested): + return requested + return FLOW_PROJECT_ID or None + async def create_project(self, project_title: str, tool_name: str = "PINHOLE") -> dict: + if not USE_BATCH_RPC: + return await self._legacy_create_project(project_title, tool_name) + pid = self.flow_project_id() + if not pid: + return {"error": _UNSUPPORTED_CREATE_PROJECT} + logger.info("Reusing pinned Flow project %s for '%s'", pid[:12], project_title) + return {"status": 200, "data": {"projectId": pid}} + + async def generate_images(self, prompt: str, project_id: str, + aspect_ratio: str = "IMAGE_ASPECT_RATIO_PORTRAIT", + user_paygate_tier: str = "PAYGATE_TIER_TWO", + character_media_ids: list[str] = None, + image_model: str = None) -> dict: + """Generate image(s). + + ``character_media_ids`` are attached as reference images, which is what + keeps an entity the same across scenes. Response is shaped like the + old REST one so the parsers downstream do not have to care which + transport produced it. + """ + if not USE_BATCH_RPC: + return await self._legacy_generate_images( + prompt, project_id, aspect_ratio, user_paygate_tier, character_media_ids) + + try: + pid = self._batch_project_id(project_id) + freq = fb.image_request( + prompt, pid, count=1, aspect=aspect_ratio, + model=self._batch_image_model(image_model), + ref_media_ids=list(character_media_ids or []) or None, + ) + payload = await self._batch_payload(fb.RPC_GEN_IMAGE, freq, fb.CAPTCHA_IMAGE) + except Exception as e: + return _batch_error(e) + + images = fb.read_images(payload) + if not images: + return {"status": 502, "error": "Image generation returned no media url"} + return {"status": 200, "data": {"media": [_as_media_record(i) for i in images]}} + + async def edit_image(self, prompt: str, source_media_id: str, + project_id: str, + aspect_ratio: str = "IMAGE_ASPECT_RATIO_PORTRAIT", + user_paygate_tier: str = "PAYGATE_TIER_ONE", + character_media_ids: list[str] = None) -> dict: + """Regenerate from an existing image plus any entity references. + + The REST path had a dedicated base-image input type; the new payload's + reference slot was captured but a base-image variant of it was not, so + here the source rides in as the first reference. In practice that + conditions the result on the source rather than editing it in place — + good enough for continuation scenes, not identical to the old edit. + Capturing the real slot is the fix; see docs/CAPTURE.md. + """ + if not USE_BATCH_RPC: + return await self._legacy_edit_image( + prompt, source_media_id, project_id, aspect_ratio, + user_paygate_tier, character_media_ids) + + refs = [source_media_id] + [ + mid for mid in (character_media_ids or []) if mid != source_media_id + ] + return await self.generate_images( + prompt=prompt, project_id=project_id, aspect_ratio=aspect_ratio, + user_paygate_tier=user_paygate_tier, character_media_ids=refs, + ) + + async def generate_video(self, start_image_media_id: str, prompt: str, + project_id: str, scene_id: str, + aspect_ratio: str = "VIDEO_ASPECT_RATIO_PORTRAIT", + end_image_media_id: str = None, + user_paygate_tier: str = "PAYGATE_TIER_TWO") -> dict: + """Submit an i2v generation. Returns operations for the poller.""" + if not USE_BATCH_RPC: + return await self._legacy_generate_video( + start_image_media_id, prompt, project_id, scene_id, + aspect_ratio, end_image_media_id, user_paygate_tier) + + if end_image_media_id: + if not FLOW_ALLOW_DEGRADED: + return {"error": _unsupported( + "start+end frame chaining", + "the new payload's end-image slot was never captured", + )} + logger.warning( + "Scene %s: dropping end frame %s — chaining is not on the batch path, " + "running plain i2v because FLOW_ALLOW_DEGRADED=1", + str(scene_id)[:12], end_image_media_id[:12]) + + gen_type = "start_end_frame_2_video" if end_image_media_id else "frame_2_video" + try: + pid = self._batch_project_id(project_id) + freq = fb.video_request( + prompt, pid, start_image_media_id, aspect=aspect_ratio, + model=self._batch_video_model(user_paygate_tier, gen_type, aspect_ratio), + ) + payload = await self._batch_payload( + fb.RPC_GEN_VIDEO, freq, fb.CAPTCHA_VIDEO, timeout=120) + operation = fb.read_operation(payload) + except Exception as e: + return _batch_error(e) + + self._remember_operation(operation.operation_id, pid) + return {"status": 200, "data": {"operations": [_as_pending_operation(operation.operation_id)]}} + + async def generate_video_from_references(self, reference_media_ids: list[str], + prompt: str, project_id: str, scene_id: str, + aspect_ratio: str = "VIDEO_ASPECT_RATIO_PORTRAIT", + user_paygate_tier: str = "PAYGATE_TIER_TWO") -> dict: + """Generate video from multiple reference images (r2v).""" + if not USE_BATCH_RPC: + return await self._legacy_generate_video_from_references( + reference_media_ids, prompt, project_id, scene_id, + aspect_ratio, user_paygate_tier) + + if not FLOW_ALLOW_DEGRADED: + return {"error": _unsupported( + "reference-to-video (r2v)", + "its payload was never captured off the new UI", + )} + if not reference_media_ids: + return {"error": "No reference media_ids for r2v"} + logger.warning( + "Scene %s: r2v is not on the batch path — running i2v off the first " + "reference %s because FLOW_ALLOW_DEGRADED=1", + str(scene_id)[:12], reference_media_ids[0][:12]) + return await self.generate_video( + start_image_media_id=reference_media_ids[0], prompt=prompt, + project_id=project_id, scene_id=scene_id, aspect_ratio=aspect_ratio, + user_paygate_tier=user_paygate_tier, + ) + + async def upscale_video(self, media_id: str, scene_id: str, + aspect_ratio: str = "VIDEO_ASPECT_RATIO_PORTRAIT", + resolution: str = "VIDEO_RESOLUTION_4K") -> dict: + """Upscale a video.""" + if not USE_BATCH_RPC: + return await self._legacy_upscale_video(media_id, scene_id, aspect_ratio, resolution) + return {"error": _unsupported( + "video upscale", + "no upsampler rpc appears in the new frontend's captures", + )} + + async def check_video_status(self, operations: list[dict]) -> dict: + """One poll round for each submitted operation. + + Three signals have to agree before a clip can be downloaded, and they + arrive out of order: + + * the operation poll says how the job is going — but it can sit at no + status at all on a job that finished, and a "Media not found." + complaint on it is survivable rather than fatal; + * the project listing is what actually gains a media id; + * the media record serves the poster image first and grows the + ``/video/`` url in later. + + So an operation only reports SUCCESSFUL once there is a video url. + Everything short of that is PENDING, and the caller's own poll loop + owns the timeout. + """ + if not USE_BATCH_RPC: + return await self._legacy_check_video_status(operations) + + out = [] + for entry in operations or []: + op_id = (entry.get("operation") or {}).get("name") or entry.get("name") or "" + if not op_id: + out.append({"operation": {}, "status": "MEDIA_GENERATION_STATUS_FAILED", + "error": "operation carried no name"}) + continue + try: + out.append(await self._poll_batch_operation(op_id)) + except Exception as e: + # A hiccup on one poll round costs a round, not the job. + logger.warning("Operation %s poll failed: %s", op_id[:20], e) + out.append(_as_pending_operation(op_id, error=str(e))) + return {"status": 200, "data": {"operations": out}} + + async def _poll_batch_operation(self, operation_id: str) -> dict: + media_id = self._operation_media.get(operation_id) + complaint = None + + if not media_id: + media_id, complaint = await self._find_operation_media(operation_id) + if not media_id: + return _as_pending_operation(operation_id, error=complaint) + self._operation_media[operation_id] = media_id + + urls = await self._batch_media_urls(media_id) + if not urls.video: + # The id landed but the clip is still being written; downloading + # now would save the poster still instead of the video. + return _as_pending_operation(operation_id, error=complaint, media_id=media_id) + + # The media id stays cached rather than being cleared here: a batch + # with several operations re-polls the finished ones alongside the + # pending ones, and a cleared entry would report them PENDING again. + # Growth is bounded by _remember_operation. + return { + "operation": { + "name": operation_id, + "metadata": {"video": {"mediaId": media_id, "fifeUrl": urls.video}}, + }, + "status": "MEDIA_GENERATION_STATUS_SUCCESSFUL", + } + + async def _find_operation_media(self, operation_id: str) -> tuple[str | None, str | None]: + """Ask the operation how it is going, then the listing where its media is. + + The listing is the authority — the poll has been seen to never report a + finished job the listing already knows about — but it is also the + expensive call, so it is only consulted when the poll says something + happened, when the poll is unreadable, or every third round regardless. + """ + rounds = self._operation_polls.get(operation_id, 0) + 1 + self._operation_polls[operation_id] = rounds + + project_id = self._operation_projects.get(operation_id) or FLOW_PROJECT_ID + complaint = None + worth_looking = rounds % 3 == 0 + try: + operation = fb.read_operation( + await self._batch_payload( + fb.RPC_OPERATION, fb.operation_request(operation_id), timeout=60) + ) + complaint = operation.error + project_id = operation.project_id or project_id + if project_id: + self._remember_operation(operation_id, project_id) + worth_looking = worth_looking or operation.done or operation.complained + except Exception as e: + # An operation that has decayed to a bare id still shows up in the + # listing, so a failed poll is a reason to look there, not to stop. + logger.debug("Operation %s poll unreadable (%s), trying the listing", + operation_id[:20], e) + worth_looking = True + + if not worth_looking: + return None, complaint + if not project_id: + return None, "no project id for the listing lookup" + return await self._media_id_for(operation_id, project_id), complaint + + async def _media_id_for(self, operation_id: str, project_id: str) -> str | None: + """Find an operation's media id in the project listing. + + Asks the extension for an 800-byte window around the operation id + rather than the whole listing — that payload is past 17 MB and grows + with every generation, so anything that ships it whole gets truncated + and loses roughly half of all lookups. + """ + result = await self.batch_rpc( + fb.RPC_PROJECT_MEDIA, fb.project_media_request(project_id), + match=operation_id, timeout=120, + ) + if result.get("error"): + raise fb.FlowBatchError(f"{fb.RPC_PROJECT_MEDIA}: {result['error']}") + raw = result.get("data") or "" + media_id = fb.find_media_id_in_text(raw, operation_id) + if not media_id and raw.lstrip().startswith(")]}"): + # an extension that cannot filter hands back the whole envelope + try: + media_id = fb.find_media_id( + fb.first_payload(raw, fb.RPC_PROJECT_MEDIA), operation_id) + except (fb.FlowBatchError, fb.RpcError, json.JSONDecodeError): + media_id = None + return media_id + + async def _batch_media_urls(self, media_id: str) -> "fb.MediaUrls": + payload = await self._batch_payload( + fb.RPC_MEDIA, fb.media_request(media_id), timeout=60) + return fb.read_media_urls(payload, media_id) + + async def get_credits(self) -> dict: + """Get user credits and tier. + + The new frontend has no captured credits rpc, and the tier no longer + selects a model — aspect is its own slot and the model names are + fixed — so on the batch path this answers with the configured default + rather than pretending to know. + """ + if not USE_BATCH_RPC: + return await self._legacy_get_credits() + return {"status": 200, "data": { + "userPaygateTier": DEFAULT_PAYGATE_TIER, + "note": "batchexecute path: tier is configured (DEFAULT_PAYGATE_TIER), not fetched", + }} + + async def validate_media_id(self, media_id: str) -> bool: + """Check if a mediaId is still valid.""" + result = await self.get_media(media_id) + status = result.get("status", 500) + return isinstance(status, int) and status == 200 + + async def get_media(self, media_id: str) -> dict: + """Fetch a media record, which is where a fresh signed url lives.""" + if not USE_BATCH_RPC: + return await self._legacy_get_media(media_id) + try: + urls = await self._batch_media_urls(media_id) + except Exception as e: + return _batch_error(e) + if not urls.video and not urls.image: + return {"status": 404, "error": f"No urls for media {media_id}"} + data: dict = {} + if urls.video: + data["video"] = {"fifeUrl": urls.video} + if urls.image: + data["image"] = {"fifeUrl": urls.image} + return {"status": 200, "data": data} + + async def upload_image(self, image_base64: str, mime_type: str = "image/jpeg", + project_id: str = "", file_name: str = "image.jpg") -> dict: + """Upload an image into the project so it can be used as a reference.""" + if not USE_BATCH_RPC: + return await self._legacy_upload_image(image_base64, mime_type, project_id, file_name) + try: + pid = self._batch_project_id(project_id) + payload = await self._batch_payload( + fb.RPC_UPLOAD_IMAGE, + fb.upload_request(image_base64, pid, mime_type, file_name), + fb.CAPTCHA_IMAGE, timeout=120, + ) + media_id = fb.read_uploaded_media_id(payload) + except Exception as e: + return _batch_error(e) + return {"status": 200, "data": {"media": {"name": media_id}}, "_mediaId": media_id} + + # ─── Legacy REST methods (aisandbox-pa, pre-migration) ─── + + async def _legacy_create_project(self, project_title: str, tool_name: str = "PINHOLE") -> dict: """Create a project on Google Flow via tRPC endpoint. Returns the full response including projectId. @@ -273,7 +762,7 @@ async def create_project(self, project_title: str, tool_name: str = "PINHOLE") - "body": body, }, timeout=30) - async def generate_images(self, prompt: str, project_id: str, + async def _legacy_generate_images(self, prompt: str, project_id: str, aspect_ratio: str = "IMAGE_ASPECT_RATIO_PORTRAIT", user_paygate_tier: str = "PAYGATE_TIER_TWO", character_media_ids: list[str] = None) -> dict: @@ -322,7 +811,7 @@ async def generate_images(self, prompt: str, project_id: str, "captchaAction": "IMAGE_GENERATION", }) - async def edit_image(self, prompt: str, source_media_id: str, + async def _legacy_edit_image(self, prompt: str, source_media_id: str, project_id: str, aspect_ratio: str = "IMAGE_ASPECT_RATIO_PORTRAIT", user_paygate_tier: str = "PAYGATE_TIER_ONE", @@ -368,7 +857,7 @@ async def edit_image(self, prompt: str, source_media_id: str, "captchaAction": "IMAGE_GENERATION", }) - async def generate_video(self, start_image_media_id: str, prompt: str, + async def _legacy_generate_video(self, start_image_media_id: str, prompt: str, project_id: str, scene_id: str, aspect_ratio: str = "VIDEO_ASPECT_RATIO_PORTRAIT", end_image_media_id: str = None, @@ -414,7 +903,7 @@ async def generate_video(self, start_image_media_id: str, prompt: str, "captchaAction": "VIDEO_GENERATION", }, timeout=60) # Submit only — polling is separate - async def generate_video_from_references(self, reference_media_ids: list[str], + async def _legacy_generate_video_from_references(self, reference_media_ids: list[str], prompt: str, project_id: str, scene_id: str, aspect_ratio: str = "VIDEO_ASPECT_RATIO_PORTRAIT", user_paygate_tier: str = "PAYGATE_TIER_TWO") -> dict: @@ -460,7 +949,7 @@ async def generate_video_from_references(self, reference_media_ids: list[str], "captchaAction": "VIDEO_GENERATION", }, timeout=60) - async def upscale_video(self, media_id: str, scene_id: str, + async def _legacy_upscale_video(self, media_id: str, scene_id: str, aspect_ratio: str = "VIDEO_ASPECT_RATIO_PORTRAIT", resolution: str = "VIDEO_RESOLUTION_4K") -> dict: """Upscale a video.""" @@ -493,7 +982,7 @@ async def upscale_video(self, media_id: str, scene_id: str, "captchaAction": "VIDEO_GENERATION", }, timeout=60) - async def check_video_status(self, operations: list[dict]) -> dict: + async def _legacy_check_video_status(self, operations: list[dict]) -> dict: """Check status of video generation operations.""" body = {"operations": operations} url = self._build_url("check_video_status") @@ -504,7 +993,7 @@ async def check_video_status(self, operations: list[dict]) -> dict: "body": body, }, timeout=30) # No captcha needed - async def get_credits(self) -> dict: + async def _legacy_get_credits(self) -> dict: """Get user credits and tier.""" url = self._build_url("get_credits") return await self._send("api_request", { @@ -513,17 +1002,7 @@ async def get_credits(self) -> dict: "headers": random_headers(), }, timeout=15) - async def validate_media_id(self, media_id: str) -> bool: - """Check if a mediaId is still valid. - - Production calls: GET /v1/media/{mediaId}?key=...&clientContext.tool=PINHOLE - Returns True on 200, False otherwise. - """ - result = await self.get_media(media_id) - status = result.get("status", 500) - return isinstance(status, int) and status == 200 - - async def get_media(self, media_id: str) -> dict: + async def _legacy_get_media(self, media_id: str) -> dict: """Fetch media metadata from Google Flow. Returns the raw API response which contains a fresh signed URL @@ -536,7 +1015,7 @@ async def get_media(self, media_id: str) -> dict: "headers": random_headers(), }, timeout=15) - async def upload_image(self, image_base64: str, mime_type: str = "image/jpeg", + async def _legacy_upload_image(self, image_base64: str, mime_type: str = "image/jpeg", project_id: str = "", file_name: str = "image.jpg") -> dict: """Upload an image for use as start/end frame. @@ -574,6 +1053,56 @@ async def upload_image(self, image_base64: str, mime_type: str = "image/jpeg", return result +# ─── Response shaping ──────────────────────────────────────── +# +# The batch path answers in Flow's positional arrays; everything downstream +# reads the old REST shapes. These put one back on the other so the parsers, +# the poller and the DB writers never learn which transport ran. + +_CAPTURE_HINT = "see docs/CAPTURE.md to record its payload off the new UI" + +_UNSUPPORTED_CREATE_PROJECT = ( + "NO_FLOW_PROJECT: Flow's project.createProject endpoint went with the September 2026 " + "migration, so Flow Kit cannot create one. Make a project in the Flow UI, then either " + "pass its uuid as flow_project_id or pin it as FLOW_PROJECT_ID." +) + + +def _unsupported(feature: str, why: str) -> str: + return f"UNSUPPORTED_ON_BATCH_API: {feature} — {why}; {_CAPTURE_HINT}." + + +def _batch_error(exc: Exception) -> dict: + """An exception from the batch path, in the error shape callers expect.""" + return {"status": 502, "error": f"{type(exc).__name__}: {exc}"} + + +def _as_media_record(image: "fb.GeneratedImage") -> dict: + """One generated image, in the REST response's `media[]` shape.""" + return { + "name": image.media_id, + "image": {"generatedImage": {"mediaId": image.media_id, "fifeUrl": image.url}}, + } + + +def _as_pending_operation(operation_id: str, error: str | None = None, + media_id: str | None = None) -> dict: + """An operation that has not produced a fetchable clip yet. + + ``error`` is carried, not acted on: a poll complaint is a diagnostic that + finished jobs also report, so it exists to make a timeout message useful. + """ + entry: dict = { + "operation": {"name": operation_id}, + "status": "MEDIA_GENERATION_STATUS_PENDING", + } + if media_id: + entry["operation"]["metadata"] = {"video": {"mediaId": media_id}} + if error: + entry["complaint"] = error + return entry + + def _is_ws_error(result: dict) -> bool: return bool(result.get("error")) or (isinstance(result.get("status"), int) and result["status"] >= 400) diff --git a/tests/unit/test_flow_batch.py b/tests/unit/test_flow_batch.py new file mode 100644 index 000000000..63b8761dc --- /dev/null +++ b/tests/unit/test_flow_batch.py @@ -0,0 +1,233 @@ +"""The batchexecute codec, builders and readers. + +These lock down the slots that cost hours to find — see docs/CAPTURE.md and the +comments in flow_batch.py. A payload Flow accepts and then ignores looks exactly +like a payload that worked, so the assertions here are about position, not shape. +""" +import json + +import pytest + +from agent.services import flow_batch as fb + + +def envelope(rpcid: str, payload) -> str: + """A response body as batchexecute serves it: sentinel, then chunks.""" + chunk = json.dumps([["wrb.fr", rpcid, json.dumps(payload)]]) + return f")]}}'\n{len(chunk)}\n{chunk}" + + +def inner(freq: str): + """The inner payload back out of an f.req envelope.""" + return json.loads(json.loads(freq)[0][0][1]) + + +class TestEnvelopeCodec: + def test_build_wraps_inner_as_a_json_string(self): + freq = fb.build_envelope("rpc1", [1, "two"]) + assert json.loads(freq) == [[["rpc1", '[1,"two"]', None, "generic"]]] + + def test_parse_reads_a_payload_back(self): + results = fb.parse_envelope(envelope("rpc1", {"a": 1})) + assert [(r.rpcid, r.data) for r in results] == [("rpc1", {"a": 1})] + + def test_parse_survives_a_truncated_tail(self): + """A response cut mid-chunk must not cost us the envelopes before it.""" + body = envelope("rpc1", {"a": 1}) + '\n50\n[["wrb.fr","rpc2","[1,2' + results = fb.parse_envelope(body) + assert [r.rpcid for r in results] == ["rpc1"] + + def test_parse_tolerates_a_missing_sentinel(self): + chunk = json.dumps([["wrb.fr", "rpc1", '{"a":1}']]) + assert fb.parse_envelope(chunk)[0].data == {"a": 1} + + def test_empty_body_is_no_results_not_a_crash(self): + assert fb.parse_envelope("") == [] + + def test_error_slot_becomes_an_error_result(self): + chunk = json.dumps([["wrb.fr", "rpc1", None, None, None, [8]]]) + result = fb.parse_envelope(f")]}}'\n{len(chunk)}\n{chunk}")[0] + assert not result.ok and result.error == [8] + + def test_first_payload_raises_on_the_error_slot(self): + chunk = json.dumps([["wrb.fr", "rpc1", None, None, None, [8]]]) + with pytest.raises(fb.RpcError): + fb.first_payload(f")]}}'\n{len(chunk)}\n{chunk}", "rpc1") + + def test_first_payload_raises_when_the_rpc_is_absent(self): + with pytest.raises(fb.FlowBatchError): + fb.first_payload(envelope("other", [1]), "rpc1") + + +class TestImageRequest: + PID = "11111111-2222-3333-4444-555555555555" + + def test_aspect_lands_in_slot_4_not_a_variant_count(self): + """Slot 4 is the aspect ratio. `count=1` only looked right because + 1 means square.""" + item = inner(fb.image_request("a cat", self.PID, + aspect="IMAGE_ASPECT_RATIO_LANDSCAPE"))[1][0] + assert item[4] == fb.ASPECT_LANDSCAPE + + def test_count_repeats_the_item_under_fresh_seeds(self): + items = inner(fb.image_request("a cat", self.PID, count=3, seed=100))[1] + assert len(items) == 3 + assert [i[3] for i in items] == [100, 100 + 9973, 100 + 2 * 9973] + + def test_prompts_give_each_variant_its_own_text(self): + items = inner(fb.image_request("fallback", self.PID, count=3, + prompts=["one", "two"]))[1] + assert [i[8][0][0][0] for i in items] == ["one", "two", "fallback"] + + def test_reference_puts_the_media_id_first_and_the_type_flag_fourth(self): + """The arrangement probing never found: wrong ones are accepted and + then quietly ignored.""" + item = inner(fb.image_request("a cat", self.PID, ref_media_ids=["mid-1"]))[1][0] + assert item[2] == [["mid-1", None, None, None, fb.REF_TYPE_IMAGE]] + + def test_no_references_leaves_the_slot_null_rather_than_empty(self): + assert inner(fb.image_request("a cat", self.PID))[1][0][2] is None + + def test_the_captcha_placeholder_is_present_for_the_extension_to_replace(self): + assert fb.CAPTCHA_SLOT in fb.image_request("a cat", self.PID) + + def test_the_project_id_rides_in_the_context(self): + assert inner(fb.image_request("a cat", self.PID))[1][0][7][5] == self.PID + + def test_the_model_is_named_in_slot_5(self): + item = inner(fb.image_request("a cat", self.PID, model="NARWHAL"))[1][0] + assert item[5] == "NARWHAL" + + +class TestVideoRequest: + PID = "11111111-2222-3333-4444-555555555555" + + def test_video_aspect_does_not_share_the_image_encoding(self): + """1 is portrait here; for an image 1 is square.""" + payload = inner(fb.video_request("go", self.PID, "mid", + aspect="VIDEO_ASPECT_RATIO_PORTRAIT")) + assert payload[0][0][2] == fb.VIDEO_ASPECT_PORTRAIT + + def test_an_image_aspect_is_refused_rather_than_rendered_wrong(self): + with pytest.raises(ValueError): + fb.video_request("go", self.PID, "mid", aspect=fb.ASPECT_LANDSCAPE) + + def test_the_source_media_id_and_a_full_frame_crop_travel_together(self): + block = inner(fb.video_request("go", self.PID, "mid-9"))[0][0][4] + assert block[1] == "mid-9" + assert block[5] == fb.FULL_FRAME_CROP + + def test_a_hand_reframed_crop_overrides_the_default(self): + crop = [None, 0.1, 1, 0.9] + assert inner(fb.video_request("go", self.PID, "mid", crop=crop))[0][0][4][5] == crop + + +class TestReaders: + OP = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee" + MID = "12345678-1234-1234-1234-1234567890ab" + + def test_images_are_read_out_of_the_url_path(self): + url = f"https://{fb.MEDIA_HOST}/image/{self.MID}?sig=x" + images = fb.read_images(["noise", [url, "more"]]) + assert images == [fb.GeneratedImage(media_id=self.MID, url=url)] + + def test_a_repeated_url_is_not_a_second_variant(self): + url = f"https://{fb.MEDIA_HOST}/image/{self.MID}?sig=x" + assert len(fb.read_images([url, url])) == 1 + + def test_operation_reads_the_id_and_status(self): + op = fb.read_operation([None, 50, [[self.OP, "proj", "scene", "CAE"]]]) + assert (op.operation_id, op.status, op.done) == (self.OP, "CAE", True) + + def test_the_third_uuid_is_the_scene_and_is_never_taken_as_media(self): + """Feeding it to the media rpc answers NOT_FOUND forever.""" + record = [self.OP, "proj", "scene-uuid", "CAE"] + op = fb.read_operation([None, 50, [record]]) + assert "scene-uuid" not in (op.operation_id, op.project_id, op.status) + + def test_a_complaint_is_carried_but_is_not_a_terminal_status(self): + detail = [None] * 8 + [[fb.OUTCOME_COMPLAINT, [None, "Media not found."]]] + op = fb.read_operation([None, 50, [[self.OP, "proj", "scene", None, None, detail]]]) + assert op.complained and op.error == "Media not found." + assert not op.done + + def test_a_healthy_outcome_carries_no_complaint(self): + detail = [None] * 8 + [[fb.OUTCOME_OK]] + op = fb.read_operation([None, 50, [[self.OP, "p", "s", "CAE", None, detail]]]) + assert op.error is None + + def test_an_empty_operation_payload_raises(self): + with pytest.raises(fb.FlowBatchError): + fb.read_operation([None, 50, []]) + + def test_the_media_id_is_found_in_an_unparsable_listing(self): + """The listing outgrows any response cap; a truncated tail still holds + the entry we came for.""" + text = f'["{self.OP}",null,null,["title",1,2,null,null,"{self.MID}"],"proj' + assert fb.find_media_id_in_text(text, self.OP) == self.MID + + def test_an_absent_operation_reads_as_not_there_yet(self): + assert fb.find_media_id_in_text("nothing here", self.OP) is None + + def test_the_media_id_is_found_in_a_decoded_listing_too(self): + payload = [[self.OP, None, None, ["t", 1, None, None, self.MID], "proj"]] + assert fb.find_media_id(payload, self.OP) == self.MID + + def test_urls_are_split_by_kind(self): + video = f"https://{fb.MEDIA_HOST}/video/{self.MID}?s=1" + image = f"https://{fb.MEDIA_HOST}/image/{self.MID}?s=1" + urls = fb.read_media_urls([image, video], self.MID) + assert (urls.video, urls.image) == (video, image) + + def test_a_poster_only_record_has_no_video_yet(self): + """A media id arrives before the clip is fetchable; downloading on the + id alone saves a still picture.""" + image = f"https://{fb.MEDIA_HOST}/image/{self.MID}?s=1" + assert fb.read_media_urls([image], self.MID).video is None + + def test_uploaded_media_id_is_the_first_slot(self): + assert fb.read_uploaded_media_id([[self.MID, "proj", "op", "CAE"]]) == self.MID + + def test_an_upload_with_no_id_raises(self): + with pytest.raises(fb.FlowBatchError): + fb.read_uploaded_media_id([[]]) + + +class TestResolvers: + def test_rest_era_aspect_names_still_work(self): + assert fb.resolve_aspect("IMAGE_ASPECT_RATIO_PORTRAIT") == fb.ASPECT_PORTRAIT + + def test_an_unknown_aspect_name_raises_rather_than_defaulting(self): + with pytest.raises(ValueError): + fb.resolve_aspect("IMAGE_ASPECT_RATIO_CINEMA") + + def test_nicknames_resolve_to_wire_names(self): + assert fb.resolve_image_model("NANO_BANANA_PRO") == "GEM_PIX_2" + assert fb.resolve_image_model("NANO_BANANA_2") == "NARWHAL" + + def test_a_wire_name_passes_through(self): + assert fb.resolve_image_model("NARWHAL") == "NARWHAL" + + def test_an_unknown_image_model_coerces_to_the_default(self): + assert fb.resolve_image_model("SOMETHING_ELSE") == fb.IMAGE_MODEL + + @pytest.mark.parametrize("legacy,expected", [ + ("veo_3_1_i2v_s_fast_ultra_relaxed", "veo_3_1_i2v_s_fast_ultra"), + ("veo_3_1_i2v_s_fast_portrait", fb.VIDEO_MODEL), + ("veo_3_1_i2v_s_fast_fl", fb.VIDEO_MODEL), + ("veo_3_1_r2v_fast_landscape_ultra_relaxed", "veo_3_1_i2v_s_fast_ultra"), + ("veo_3_1_i2v_lite", "veo_3_1_i2v_lite"), + (None, fb.VIDEO_MODEL), + ]) + def test_rest_era_video_keys_fold_onto_accepted_names(self, legacy, expected): + """Aspect and chaining are their own slots now; the suffixed names are + rejected outright, so only the tier/quality intent survives.""" + assert fb.resolve_video_model(legacy) == expected + + def test_every_resolved_video_model_is_one_flow_accepts(self): + for tier in ("PAYGATE_TIER_ONE", "PAYGATE_TIER_TWO"): + for gen in ("frame_2_video", "start_end_frame_2_video", "reference_frame_2_video"): + for aspect in ("VIDEO_ASPECT_RATIO_PORTRAIT", "VIDEO_ASPECT_RATIO_LANDSCAPE"): + from agent.config import VIDEO_MODELS + key = VIDEO_MODELS.get(tier, {}).get(gen, {}).get(aspect) + assert fb.resolve_video_model(key) in fb.VIDEO_MODELS diff --git a/tests/unit/test_flow_client_batch.py b/tests/unit/test_flow_client_batch.py new file mode 100644 index 000000000..733f52440 --- /dev/null +++ b/tests/unit/test_flow_client_batch.py @@ -0,0 +1,395 @@ +"""The batch path's answers, in the shapes the rest of the pipeline reads. + +Everything downstream of FlowClient — the worker's parsers, the operation +poller, the scene writers — was written against the old REST responses. These +tests hold the adapter to that contract, so a transport swap stays invisible. +""" +import json + +import pytest + +from agent.services import flow_batch as fb +from agent.services.flow_client import FlowClient +from agent.worker._parsing import _extract_media_id, _extract_output_url, _is_error + +PROJECT = "11111111-2222-3333-4444-555555555555" +MEDIA = "12345678-1234-1234-1234-1234567890ab" +OPERATION = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee" +IMAGE_URL = f"https://{fb.MEDIA_HOST}/image/{MEDIA}?sig=x" +VIDEO_URL = f"https://{fb.MEDIA_HOST}/video/{MEDIA}?sig=x" + + +def envelope(rpcid: str, payload) -> str: + chunk = json.dumps([["wrb.fr", rpcid, json.dumps(payload)]]) + return f")]}}'\n{len(chunk)}\n{chunk}" + + +@pytest.fixture +def client(monkeypatch): + """A FlowClient whose transport replays canned RPC responses. + + `calls` records what each rpc was asked, so a test can assert on the + envelope as well as on what came back. + """ + import agent.services.flow_client as module + monkeypatch.setattr(module, "USE_BATCH_RPC", True) + monkeypatch.setattr(module, "FLOW_PROJECT_ID", PROJECT) + monkeypatch.setattr(module, "FLOW_ALLOW_DEGRADED", False) + + c = FlowClient() + c.responses = {} + c.calls = [] + + async def fake_batch_rpc(rpcid, freq, captcha_action=None, match=None, timeout=300): + c.calls.append({"rpcid": rpcid, "freq": freq, + "captcha": captcha_action, "match": match}) + canned = c.responses.get(rpcid, {"data": ""}) + return canned(match) if callable(canned) else canned + + c.batch_rpc = fake_batch_rpc + return c + + +class TestGenerateImages: + async def test_answers_in_the_shape_the_media_parser_reads(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[IMAGE_URL]])} + result = await client.generate_images("a cat", PROJECT) + + assert not _is_error(result) + assert _extract_media_id(result, "GENERATE_IMAGE") == MEDIA + assert _extract_output_url(result, "GENERATE_IMAGE") == IMAGE_URL + + async def test_asks_for_a_captcha(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[IMAGE_URL]])} + await client.generate_images("a cat", PROJECT) + assert client.calls[0]["captcha"] == fb.CAPTCHA_IMAGE + + async def test_character_refs_ride_in_the_reference_slot(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[IMAGE_URL]])} + await client.generate_images("a cat", PROJECT, character_media_ids=["ref-a", "ref-b"]) + + item = json.loads(json.loads(client.calls[0]["freq"])[0][0][1])[1][0] + assert item[2] == [["ref-a", None, None, None, fb.REF_TYPE_IMAGE], + ["ref-b", None, None, None, fb.REF_TYPE_IMAGE]] + + async def test_a_project_less_call_falls_back_to_the_pinned_project(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[IMAGE_URL]])} + await client.generate_images("a cat", "0") + + item = json.loads(json.loads(client.calls[0]["freq"])[0][0][1])[1][0] + assert item[7][5] == PROJECT + + async def test_no_url_back_is_an_error_not_a_silent_success(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[]])} + assert _is_error(await client.generate_images("a cat", PROJECT)) + + async def test_a_transport_error_becomes_an_error_result(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"error": "CAPTCHA_FAILED: NO_FLOW_TAB"} + result = await client.generate_images("a cat", PROJECT) + assert _is_error(result) and "NO_FLOW_TAB" in result["error"] + + async def test_no_project_anywhere_is_a_named_failure(self, client, monkeypatch): + import agent.services.flow_client as module + monkeypatch.setattr(module, "FLOW_PROJECT_ID", "") + result = await client.generate_images("a cat", "0") + assert "NO_FLOW_PROJECT" in result["error"] + + +class TestEditImage: + async def test_the_source_leads_the_reference_list(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[IMAGE_URL]])} + await client.edit_image("redraw", "src-1", PROJECT, character_media_ids=["ref-a"]) + + item = json.loads(json.loads(client.calls[0]["freq"])[0][0][1])[1][0] + assert [ref[0] for ref in item[2]] == ["src-1", "ref-a"] + + async def test_the_source_is_not_repeated_when_it_is_also_a_character(self, client): + client.responses[fb.RPC_GEN_IMAGE] = {"data": envelope(fb.RPC_GEN_IMAGE, [[IMAGE_URL]])} + await client.edit_image("redraw", "src-1", PROJECT, character_media_ids=["src-1", "ref-a"]) + + item = json.loads(json.loads(client.calls[0]["freq"])[0][0][1])[1][0] + assert [ref[0] for ref in item[2]] == ["src-1", "ref-a"] + + +class TestGenerateVideo: + def _submitted(self, client): + return {"data": envelope(fb.RPC_GEN_VIDEO, [None, 50, [[OPERATION, PROJECT, "scene", None]]])} + + async def test_returns_an_operation_the_poller_can_carry(self, client): + client.responses[fb.RPC_GEN_VIDEO] = self._submitted(client) + result = await client.generate_video("mid", "go", PROJECT, "scene-1") + + ops = result["data"]["operations"] + assert ops[0]["operation"]["name"] == OPERATION + assert ops[0]["status"] == "MEDIA_GENERATION_STATUS_PENDING" + + async def test_remembers_which_project_to_look_the_media_up_in(self, client): + client.responses[fb.RPC_GEN_VIDEO] = self._submitted(client) + await client.generate_video("mid", "go", PROJECT, "scene-1") + assert client._operation_projects[OPERATION] == PROJECT + + async def test_chaining_fails_loudly_rather_than_dropping_the_end_frame(self, client): + result = await client.generate_video("mid", "go", PROJECT, "scene-1", + end_image_media_id="end-mid") + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + assert not client.calls, "nothing should have been sent" + + async def test_degraded_mode_runs_i2v_off_the_start_frame(self, client, monkeypatch): + import agent.services.flow_client as module + monkeypatch.setattr(module, "FLOW_ALLOW_DEGRADED", True) + client.responses[fb.RPC_GEN_VIDEO] = self._submitted(client) + + result = await client.generate_video("start-mid", "go", PROJECT, "scene-1", + end_image_media_id="end-mid") + assert not _is_error(result) + payload = json.loads(json.loads(client.calls[0]["freq"])[0][0][1]) + assert payload[0][0][4][1] == "start-mid" + + async def test_r2v_fails_loudly_by_default(self, client): + result = await client.generate_video_from_references(["a", "b"], "go", PROJECT, "s") + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + + async def test_degraded_r2v_uses_the_first_reference_as_the_start_frame(self, client, monkeypatch): + import agent.services.flow_client as module + monkeypatch.setattr(module, "FLOW_ALLOW_DEGRADED", True) + client.responses[fb.RPC_GEN_VIDEO] = self._submitted(client) + + await client.generate_video_from_references(["ref-a", "ref-b"], "go", PROJECT, "s") + payload = json.loads(json.loads(client.calls[0]["freq"])[0][0][1]) + assert payload[0][0][4][1] == "ref-a" + + async def test_upscale_is_unported_and_has_no_fallback(self, client, monkeypatch): + import agent.services.flow_client as module + monkeypatch.setattr(module, "FLOW_ALLOW_DEGRADED", True) + result = await client.upscale_video(MEDIA, "scene-1") + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + + +class TestCheckVideoStatus: + def _poll(self, status=None, complaint=None): + detail = None + if complaint: + detail = [None] * 8 + [[fb.OUTCOME_COMPLAINT, [None, complaint]]] + record = [OPERATION, PROJECT, "scene", status, None, detail] + return {"data": envelope(fb.RPC_OPERATION, [None, 50, [record]])} + + def _listing(self, found=True): + text = (f'["{OPERATION}",null,null,["t",1,2,null,null,"{MEDIA}"]' if found else "") + return lambda match: {"data": text} + + async def _status(self, client): + result = await client.check_video_status([{"operation": {"name": OPERATION}}]) + return result["data"]["operations"][0] + + async def test_successful_once_a_video_url_exists(self, client): + client.responses[fb.RPC_OPERATION] = self._poll(status="CAE") + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL])} + + op = await self._status(client) + assert op["status"] == "MEDIA_GENERATION_STATUS_SUCCESSFUL" + assert _extract_media_id({"data": {"operations": [op]}}, "GENERATE_VIDEO") == MEDIA + assert _extract_output_url({"data": {"operations": [op]}}, "GENERATE_VIDEO") == VIDEO_URL + + async def test_a_media_id_with_only_a_poster_is_still_pending(self, client): + """Downloading on the id alone would save a still picture.""" + client.responses[fb.RPC_OPERATION] = self._poll(status="CAE") + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [IMAGE_URL])} + + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_PENDING" + + async def test_a_complaint_is_carried_but_does_not_fail_the_job(self, client): + """Jobs report "Media not found." and still deliver a finished clip.""" + client.responses[fb.RPC_OPERATION] = self._poll(complaint="Media not found.") + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing(found=False) + + op = await self._status(client) + assert op["status"] == "MEDIA_GENERATION_STATUS_PENDING" + assert op["complaint"] == "Media not found." + + async def test_the_listing_decides_even_when_the_poll_never_says_done(self, client): + """The poll can sit at no status at all on a job that finished, so the + listing is consulted on a schedule rather than only on the poll's say-so.""" + client.responses[fb.RPC_OPERATION] = self._poll(status=None) + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL])} + + await self._status(client) + await self._status(client) + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_SUCCESSFUL" + + async def test_the_listing_is_asked_for_a_window_not_the_whole_thing(self, client): + client.responses[fb.RPC_OPERATION] = self._poll(status="CAE") + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing(found=False) + + await self._status(client) + listing = next(c for c in client.calls if c["rpcid"] == fb.RPC_PROJECT_MEDIA) + assert listing["match"] == OPERATION + + async def test_an_unreadable_poll_still_consults_the_listing(self, client): + """Old operations decay to a bare id but stay in the listing.""" + client.responses[fb.RPC_OPERATION] = {"error": "boom"} + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL])} + + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_SUCCESSFUL" + + async def test_a_quiet_poll_does_not_pay_for_the_listing_every_round(self, client): + """The listing is a 17 MB call; a poll with nothing to report skips it.""" + client.responses[fb.RPC_OPERATION] = self._poll(status=None) + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL])} + + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_PENDING" + assert not [c for c in client.calls if c["rpcid"] == fb.RPC_PROJECT_MEDIA] + + await self._status(client) + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_SUCCESSFUL" + + async def test_a_known_media_id_is_not_looked_up_again(self, client): + """Once the listing has answered, later rounds go straight to the media.""" + client.responses[fb.RPC_OPERATION] = self._poll(status="CAE") + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [IMAGE_URL])} + + await self._status(client) # poster only — still pending + client.calls.clear() + await self._status(client) + assert not [c for c in client.calls if c["rpcid"] == fb.RPC_PROJECT_MEDIA] + assert not [c for c in client.calls if c["rpcid"] == fb.RPC_OPERATION] + + async def test_a_finished_operation_stays_finished_when_re_polled(self, client): + """A batch re-polls its finished operations alongside its pending ones.""" + client.responses[fb.RPC_OPERATION] = self._poll(status="CAE") + client.responses[fb.RPC_PROJECT_MEDIA] = self._listing() + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL])} + + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_SUCCESSFUL" + assert (await self._status(client))["status"] == "MEDIA_GENERATION_STATUS_SUCCESSFUL" + + async def test_a_nameless_operation_fails_instead_of_polling_forever(self, client): + result = await client.check_video_status([{"operation": {}}]) + assert result["data"]["operations"][0]["status"] == "MEDIA_GENERATION_STATUS_FAILED" + + +class TestMediaAndUpload: + async def test_get_media_reports_the_signed_urls(self, client): + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL, IMAGE_URL])} + result = await client.get_media(MEDIA) + assert result["status"] == 200 + assert result["data"]["video"]["fifeUrl"] == VIDEO_URL + + async def test_a_media_id_with_no_urls_reads_as_404(self, client): + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [])} + assert (await client.get_media(MEDIA))["status"] == 404 + + async def test_validate_media_id_follows_the_status(self, client): + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL])} + assert await client.validate_media_id(MEDIA) is True + client.responses[fb.RPC_MEDIA] = {"data": envelope(fb.RPC_MEDIA, [])} + assert await client.validate_media_id(MEDIA) is False + + async def test_upload_returns_the_media_id_the_callers_look_for(self, client): + client.responses[fb.RPC_UPLOAD_IMAGE] = { + "data": envelope(fb.RPC_UPLOAD_IMAGE, [[MEDIA, PROJECT, OPERATION, "CAE"]]) + } + result = await client.upload_image("Ym9keQ==", project_id=PROJECT) + assert result["_mediaId"] == MEDIA + assert result["data"]["media"]["name"] == MEDIA + + async def test_upload_carries_a_captcha_like_a_generate(self, client): + client.responses[fb.RPC_UPLOAD_IMAGE] = { + "data": envelope(fb.RPC_UPLOAD_IMAGE, [[MEDIA, PROJECT, OPERATION, "CAE"]]) + } + await client.upload_image("Ym9keQ==", project_id=PROJECT) + assert client.calls[0]["captcha"] == fb.CAPTCHA_IMAGE + + +class TestProjectAndCredits: + async def test_create_project_hands_back_the_pinned_one(self, client): + result = await client.create_project("My Film") + assert result["data"]["projectId"] == PROJECT + + async def test_create_project_without_a_pin_explains_itself(self, client, monkeypatch): + import agent.services.flow_client as module + monkeypatch.setattr(module, "FLOW_PROJECT_ID", "") + result = await client.create_project("My Film") + assert "NO_FLOW_PROJECT" in result["error"] + assert "FLOW_PROJECT_ID" in result["error"] + + async def test_credits_answers_the_configured_tier_rather_than_guessing(self, client): + result = await client.get_credits() + assert result["data"]["userPaygateTier"] + assert not client.calls, "there is no credits rpc to call" + + +class TestRefreshProjectUrls: + """Re-signing every stored media id — what `/fk-refresh-urls` runs.""" + + @pytest.fixture + def db(self, monkeypatch): + """A stand-in for the crud layer, recording what got written.""" + from agent.db import crud + + state = { + "videos": [{"id": "vid-1"}], + "scenes": [{ + "id": "scene-1", + "vertical_image_media_id": MEDIA, + "vertical_video_media_id": "22222222-2222-2222-2222-222222222222", + "horizontal_image_media_id": "CAMSnot-a-uuid", + }], + "characters": [{"id": "char-1", "media_id": "33333333-3333-3333-3333-333333333333"}], + "writes": [], + } + + async def list_videos(pid): return state["videos"] + async def list_scenes(vid): return state["scenes"] + async def get_project_characters(pid): return state["characters"] + async def update_scene(sid, **kw): state["writes"].append(("scene", sid, kw)) + async def update_character(cid, **kw): state["writes"].append(("character", cid, kw)) + + for name, fn in [("list_videos", list_videos), ("list_scenes", list_scenes), + ("get_project_characters", get_project_characters), + ("update_scene", update_scene), ("update_character", update_character)]: + monkeypatch.setattr(crud, name, fn) + return state + + async def test_writes_a_fresh_url_into_each_field_that_holds_the_id(self, client, db): + def media(match): + return {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL, IMAGE_URL])} + client.responses[fb.RPC_MEDIA] = media + + result = await client.refresh_project_urls(PROJECT) + + assert result["found"] == 3, "the CAMS id is not a media id and is skipped" + assert result["refreshed"] == 3 + written = {(table, tuple(kw)[0]) for table, _, kw in db["writes"]} + assert written == { + ("scene", "vertical_image_url"), + ("scene", "vertical_video_url"), + ("character", "reference_image_url"), + } + + async def test_an_image_field_takes_the_image_url_not_the_video_one(self, client, db): + client.responses[fb.RPC_MEDIA] = lambda m: { + "data": envelope(fb.RPC_MEDIA, [VIDEO_URL, IMAGE_URL])} + + await client.refresh_project_urls(PROJECT) + by_field = {tuple(kw)[0]: tuple(kw.values())[0] for _, _, kw in db["writes"]} + assert by_field["vertical_image_url"] == IMAGE_URL + assert by_field["vertical_video_url"] == VIDEO_URL + + async def test_one_dead_media_id_does_not_sink_the_rest(self, client, db): + seen = [] + + def media(match): + seen.append(1) + if len(seen) == 1: + return {"error": "NOT_FOUND"} + return {"data": envelope(fb.RPC_MEDIA, [VIDEO_URL, IMAGE_URL])} + client.responses[fb.RPC_MEDIA] = media + + result = await client.refresh_project_urls(PROJECT) + assert result["found"] == 3 and result["refreshed"] == 2 From aa3c917d49a9b4a10e5a3c062c3a2690e5b606d9 Mon Sep 17 00:00:00 2001 From: Hoang Tuan Nguyen Date: Sun, 6 Sep 2026 22:50:03 +0700 Subject: [PATCH 2/6] feat(extension): run Flow RPCs inside the signed-in page MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Only the page can sign a batchexecute call, so the new `batch_rpc` handler mints a reCAPTCHA through the existing grecaptcha bridge and runs the POST in the tab's MAIN world, where `at` / `f.sid` / `bl` live. The project listing is past 17 MB for the one entry we want, so the runner cuts it to an 800-byte window around the operation id before handing it back rather than pushing the rest through the bridge. Every tab lookup now goes through one `flowUrls` list that includes flow.google.com. Matching only the old URL is what produced `CAPTCHA_FAILED: NO_FLOW_TAB` against a Flow tab sitting right there. Captcha solving tries each Flow tab in turn and reloads discarded ones instead of returning on the first miss — one stale tab used to veto every generation. The bearer capture, the aisandbox-pa proxy and the telemetry are kept for USE_BATCH_RPC=0 and an old pinned tab, and are labelled as such. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015JkgMuy89iJ9qc1NkRdPk5 --- extension/background.js | 259 ++++++++++++++++++++++++++++++++++------ extension/injected.js | 11 +- extension/manifest.json | 30 ++++- 3 files changed, 257 insertions(+), 43 deletions(-) diff --git a/extension/background.js b/extension/background.js index fb1b1b415..f7fd1fcef 100644 --- a/extension/background.js +++ b/extension/background.js @@ -2,13 +2,30 @@ * Flow Kit — Chrome Extension Background Service Worker * * Connects to local Python agent via WebSocket (agent runs WS server). - * Captures bearer token, solves reCAPTCHA, proxies API calls through browser. + * Mints reCAPTCHA and runs Flow's batchexecute RPCs inside the Flow tab. + * + * Flow moved to flow.google.com in September 2026 and stopped minting the + * `Bearer ya29.…` the old REST host needed. The current path is `batch_rpc`: + * the agent builds an `f.req` envelope, this worker mints a captcha for it and + * runs the POST in the page's MAIN world, where the `at` CSRF token lives. + * The bearer capture and `api_request` proxy below are the legacy path, kept + * for USE_BATCH_RPC=0 and for an old pinned labs.google tab. */ const AGENT_WS_URL = 'ws://127.0.0.1:9222'; // NOTE: This is a browser-restricted public API key — safe to ship in extension bundles. const API_KEY = 'AIzaSyBtrm0o5ab1c-Ec8ZuLcGt3oJAA5VWt3pY'; +// labs.google/fx/tools/flow still resolves but redirects here, so in practice a +// signed-in tab is only ever flow.google.com/*. The legacy patterns stay for an +// old pinned tab. Every tab lookup in this file goes through this list. +const flowUrls = [ + 'https://flow.google.com/*', + 'https://labs.google/fx/tools/flow*', + 'https://labs.google/fx/*/tools/flow*', +]; +const FLOW_TAB_URL = 'https://flow.google.com/'; + let ws = null; let flowKey = null; let callbackSecret = null; // Auth secret for HTTP callback, received from server on WS connect @@ -113,9 +130,7 @@ chrome.webRequest.onBeforeSendHeaders.addListener( let _openingFlowTab = false; async function captureTokenFromFlowTab() { - const tabs = await chrome.tabs.query({ - url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'], - }); + const tabs = await chrome.tabs.query({ url: flowUrls }); if (!tabs.length) { if (_openingFlowTab) { console.log('[FlowAgent] Flow tab already opening, skipping'); @@ -124,11 +139,9 @@ async function captureTokenFromFlowTab() { _openingFlowTab = true; try { console.log('[FlowAgent] No Flow tab found — opening one in background'); - await chrome.tabs.create({ url: 'https://labs.google/fx/tools/flow', active: false }); + await chrome.tabs.create({ url: FLOW_TAB_URL, active: false }); await sleep(3000); - const retryTabs = await chrome.tabs.query({ - url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'], - }); + const retryTabs = await chrome.tabs.query({ url: flowUrls }); if (!retryTabs.length) { console.log('[FlowAgent] Flow tab not ready yet after open'); return; @@ -194,7 +207,9 @@ function connectToAgent() { try { const msg = JSON.parse(data); - if (msg.method === 'api_request') { + if (msg.method === 'batch_rpc') { + await handleBatchRpc(msg); + } else if (msg.method === 'api_request') { await handleApiRequest(msg); } else if (msg.method === 'trpc_request') { await handleTrpcRequest(msg); @@ -297,39 +312,83 @@ async function requestCaptchaFromTab(tabId, requestId, pageAction) { } } +/** Try to wake a discarded Flow tab so `sendMessage` can reach it. + * Chrome auto-discards backgrounded tabs to save memory; the tab still shows + * up in `chrome.tabs.query` but cross-context calls fail with "No current + * window" / "No tab with id". A reload re-hydrates it. */ +async function reviveTabIfNeeded(tab) { + if (!tab?.discarded) return tab; + try { + await chrome.tabs.reload(tab.id); + await sleep(2500); + return await chrome.tabs.get(tab.id); + } catch { + return null; + } +} + +function captchaFromTab(tabId, requestId, captchaAction) { + return Promise.race([ + requestCaptchaFromTab(tabId, requestId, captchaAction), + new Promise((_, rej) => setTimeout(() => rej(new Error('CAPTCHA_TIMEOUT')), 30000)), + ]); +} + async function solveCaptcha(requestId, captchaAction) { - const tabs = await chrome.tabs.query({ - url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'], - }); + let tabs = await chrome.tabs.query({ url: flowUrls }); + // No Flow tab at all — spawn one and let it settle. if (!tabs.length) { - // Auto-open Flow tab and wait briefly before returning error try { - await chrome.tabs.create({ url: 'https://labs.google/fx/tools/flow', active: false }); + await chrome.tabs.create({ url: FLOW_TAB_URL, active: false }); await sleep(3000); - // Retry tab query after opening - const retryTabs = await chrome.tabs.query({ - url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'], - }); - if (!retryTabs.length) return { error: 'NO_FLOW_TAB' }; - const resp = await Promise.race([ - requestCaptchaFromTab(retryTabs[0].id, requestId, captchaAction), - new Promise((_, rej) => setTimeout(() => rej(new Error('CAPTCHA_TIMEOUT')), 30000)), - ]); - return resp; + tabs = await chrome.tabs.query({ url: flowUrls }); } catch (e) { return { error: e.message || 'NO_FLOW_TAB' }; } + if (!tabs.length) return { error: 'NO_FLOW_TAB' }; + } + + // Try each Flow tab in turn. A tab that answers "no grecaptcha" is a tab + // sitting on a page that never loaded it — another Flow tab may well be + // fine. Returning on the first one let one stale tab veto every generation. + const errors = []; + for (const candidate of tabs) { + const tab = await reviveTabIfNeeded(candidate); + if (!tab) continue; + try { + const resp = await captchaFromTab(tab.id, requestId, captchaAction); + if (!resp?.token) { + errors.push(resp?.error || 'NO_TOKEN'); + continue; + } + return resp; + } catch (e) { + const msg = e?.message || ''; + errors.push(msg); + // Tab evaporated mid-call (window closed, discarded again, navigated + // away). Move on to the next candidate rather than failing the job. + if ( + msg.includes('No current window') || + msg.includes('No tab with id') || + msg.includes('Receiving end does not exist') + ) { + continue; + } + return { error: msg }; + } } + // Every candidate failed — last-ditch, spawn a fresh tab and try it once. try { - const resp = await Promise.race([ - requestCaptchaFromTab(tabs[0].id, requestId, captchaAction), - new Promise((_, rej) => setTimeout(() => rej(new Error('CAPTCHA_TIMEOUT')), 30000)), - ]); - return resp; + await chrome.tabs.create({ url: FLOW_TAB_URL, active: false }); + await sleep(3000); + const fresh = await chrome.tabs.query({ url: flowUrls }); + const target = fresh.find((t) => !t.discarded) || fresh[0]; + if (!target) return { error: 'NO_FLOW_TAB' }; + return await captchaFromTab(target.id, requestId, captchaAction); } catch (e) { - return { error: e.message }; + return { error: e?.message || errors[0] || 'NO_FLOW_TAB' }; } } @@ -350,6 +409,136 @@ async function handleSolveCaptcha(msg) { sendToAgent({ id, result }); } +// ─── Page-context RPC runner (the current path) ───────────── +// +// Flow's frontend signs its calls with cookies and a per-page `at` token, and +// every generate carries a single-use reCAPTCHA. None of that can be replayed +// from the service worker, so the request has to be issued by the Flow page +// itself: mint a fresh captcha through the grecaptcha bridge, then run the +// batchexecute POST in the page's MAIN world, where at / f.sid / bl live. + +const CAPTCHA_SLOT = '__CAPTCHA__'; +const MAX_RPC_TEXT = 32000000; // the project listing alone is past 17 MB + +async function runBatchRpc(cmd) { + const tabs = await chrome.tabs.query({ url: flowUrls }); + let candidate = tabs.find((t) => !t.discarded) || tabs[0]; + if (!candidate) { + // No Flow tab — open one and give the app a moment to boot, otherwise + // WIZ_global_data is not on the page yet and `at` comes back empty. + try { + await chrome.tabs.create({ url: FLOW_TAB_URL, active: false }); + await sleep(5000); + const fresh = await chrome.tabs.query({ url: flowUrls }); + candidate = fresh.find((t) => !t.discarded) || fresh[0]; + } catch (e) { + return { error: e?.message || 'NO_FLOW_TAB' }; + } + if (!candidate) return { error: 'NO_FLOW_TAB' }; + } + // Chrome discards backgrounded tabs; executeScript throws on a dead one. + const tab = await reviveTabIfNeeded(candidate); + if (!tab) return { error: 'FLOW_TAB_DISCARDED' }; + + let freq = cmd.freq; + if (cmd.captchaAction) { + const solved = await solveCaptcha(cmd.id, cmd.captchaAction); + if (!solved?.token) return { error: `CAPTCHA_FAILED: ${solved?.error || 'no token'}` }; + freq = freq.split(CAPTCHA_SLOT).join(solved.token); + } + + const [injected] = await chrome.scripting.executeScript({ + target: { tabId: tab.id }, + world: 'MAIN', + args: [cmd.rpcid, freq, MAX_RPC_TEXT, cmd.match || null], + func: async (rpcid, freqStr, maxText, match) => { + const wiz = globalThis.WIZ_global_data || {}; + const at = wiz.SNlM0e; + const sid = wiz.FdrFJe; + const bl = wiz.cfb2h; + if (!at) return { error: 'NO_AT_TOKEN' }; + const reqid = Math.floor(Math.random() * 900000) + 100000; + const url = + `/_/AiSandboxAngularFrontend/data/batchexecute?rpcids=${encodeURIComponent(rpcid)}` + + `&f.sid=${encodeURIComponent(sid || '')}&bl=${encodeURIComponent(bl || '')}` + + `&hl=en-AU&_reqid=${reqid}&rt=c`; + const resp = await fetch(url, { + method: 'POST', + credentials: 'include', + headers: { + 'content-type': 'application/x-www-form-urlencoded;charset=UTF-8', + 'x-same-domain': '1', + }, + body: new URLSearchParams({ 'f.req': freqStr, at }), + }); + const text = await resp.text(); + // The project listing is tens of megabytes and all we ever want from it + // is one entry. Cutting it down here keeps that payload inside the tab + // instead of pushing it through the bridge on every poll. + if (match) { + const found = text.indexOf(match); // not `at` — that is the CSRF token above + return { + status: resp.status, + matched: found !== -1, + text: found === -1 ? '' : text.slice(found, found + 800), + }; + } + return { status: resp.status, text: text.slice(0, maxText) }; + }, + }); + + return injected?.result || { error: 'NO_INJECTION_RESULT' }; +} + +async function handleBatchRpc(msg) { + const { id, params } = msg; + const { rpcid, freq, captchaAction, match } = params || {}; + if (!rpcid || !freq) { + sendToAgent({ id, status: 400, error: 'INVALID_BATCH_RPC' }); + return; + } + + setState('running'); + const hasCaptcha = !!captchaAction; + if (hasCaptcha) metrics.requestCount++; + // Polls and listing lookups run constantly; only the generates are worth + // a row in the log the popup shows. + const visible = hasCaptcha; + if (visible) { + addRequestLog({ + id, type: `RPC:${rpcid}`, time: new Date().toISOString(), + status: 'processing', error: null, outputUrl: null, url: rpcid, + payloadSummary: freq.slice(0, 200), + }); + } + + try { + const out = await runBatchRpc({ id, rpcid, freq, captchaAction, match }); + if (out.error) { + if (hasCaptcha) { metrics.failedCount++; metrics.lastError = out.error; } + if (visible) updateRequestLog(id, { status: 'failed', error: out.error }); + sendToAgent({ id, status: 502, error: out.error }); + } else { + if (hasCaptcha) { metrics.successCount++; metrics.lastError = null; } + if (visible) { + updateRequestLog(id, { + status: 'success', httpStatus: out.status, + responseSummary: (out.text || '').slice(0, 300), + }); + } + sendToAgent({ id, status: out.status, data: out.text }); + } + } catch (e) { + const err = e?.message || 'BATCH_RPC_FAILED'; + if (hasCaptcha) { metrics.failedCount++; metrics.lastError = err; } + if (visible) updateRequestLog(id, { status: 'failed', error: err }); + sendToAgent({ id, status: 500, error: err }); + } + + chrome.storage.local.set({ metrics }); + setState('idle'); +} + // ─── API Request Proxy ────────────────────────────────────── async function handleTrpcRequest(msg) { @@ -394,6 +583,8 @@ async function handleTrpcRequest(msg) { } } +// Legacy REST proxy against aisandbox-pa. Reachable only with USE_BATCH_RPC=0 +// on a profile that still holds a `Bearer ya29.…`; Flow stopped minting those. async function handleApiRequest(msg) { const { id, params } = msg; const { url, method, headers, body, captchaAction } = params; @@ -565,14 +756,12 @@ chrome.runtime.onMessage.addListener((msg, _, reply) => { } if (msg.type === 'OPEN_FLOW_TAB') { - chrome.tabs.query({ - url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'], - }).then((tabs) => { + chrome.tabs.query({ url: flowUrls }).then((tabs) => { if (tabs.length) { chrome.tabs.update(tabs[0].id, { active: true }); reply({ ok: true, tabId: tabs[0].id }); } else { - chrome.tabs.create({ url: 'https://labs.google/fx/tools/flow' }) + chrome.tabs.create({ url: FLOW_TAB_URL }) .then((tab) => reply({ ok: true, tabId: tab.id })) .catch((e) => reply({ error: e.message })); } @@ -709,6 +898,8 @@ function _buildFrontendEventsPayload() { } async function sendTelemetry() { + // Legacy-path camouflage: these endpoints want the bearer Flow no longer + // mints, so on the batch path there is no flowKey and this is a no-op. if (!flowKey || state === 'off') return; const headers = { diff --git a/extension/injected.js b/extension/injected.js index 2e15672b5..e26b690a8 100644 --- a/extension/injected.js +++ b/extension/injected.js @@ -1,6 +1,11 @@ /** - * Injected into MAIN world on labs.google — has access to window.grecaptcha - * Also intercepts TRPC fetch responses to capture fresh signed media URLs. + * Injected into the page's MAIN world on flow.google.com (and an old pinned + * labs.google tab) — has access to window.grecaptcha. + * + * The reCAPTCHA site key survived the September 2026 migration unchanged. The + * TRPC fetch intercept below did not: it belongs to the labs.google frontend + * and is inert on flow.google.com, where media urls come back inline on the + * generate call and from the media rpc. */ const SITE_KEY = '6LdsFiUsAAAAAIjVDZcuLhaHiDn5nnHVXVRQGeMV'; @@ -46,7 +51,7 @@ window.addEventListener('GET_CAPTCHA', async ({ detail }) => { } }); -function waitForGrecaptcha(timeout = 10000) { +function waitForGrecaptcha(timeout = 22000) { // it loads lazily; 10s was optimistic return new Promise((resolve, reject) => { const start = Date.now(); const check = () => { diff --git a/extension/manifest.json b/extension/manifest.json index dff8b6ef4..1ae492a63 100644 --- a/extension/manifest.json +++ b/extension/manifest.json @@ -1,10 +1,20 @@ { "manifest_version": 3, "name": "Flow Kit", - "version": "0.2.0", - "description": "Local agent bridge for Google Flow API — captures tokens, solves reCAPTCHA, proxies API calls", - "permissions": ["storage", "alarms", "tabs", "webRequest", "scripting", "declarativeNetRequest", "sidePanel"], + "version": "0.3.0", + "description": "Local agent bridge for Google Flow \u2014 runs batchexecute RPCs inside a signed-in flow.google.com tab, mints reCAPTCHA, proxies API calls", + "permissions": [ + "storage", + "alarms", + "tabs", + "webRequest", + "scripting", + "declarativeNetRequest", + "sidePanel" + ], "host_permissions": [ + "https://flow.google.com/*", + "https://flow-content.google/*", "https://labs.google/*", "https://aisandbox-pa.googleapis.com/*", "http://127.0.0.1:8100/*" @@ -15,17 +25,25 @@ "content_scripts": [ { "matches": [ + "https://flow.google.com/*", "https://labs.google/fx/tools/flow*", "https://labs.google/fx/*/tools/flow*" ], - "js": ["content.js"], + "js": [ + "content.js" + ], "run_at": "document_start" } ], "web_accessible_resources": [ { - "resources": ["injected.js"], - "matches": ["https://labs.google/*"] + "resources": [ + "injected.js" + ], + "matches": [ + "https://flow.google.com/*", + "https://labs.google/*" + ] } ], "declarative_net_request": { From 62a31a8bfa4425027f516e3dfd1c6e1691264bbf Mon Sep 17 00:00:00 2001 From: Hoang Tuan Nguyen Date: Sun, 6 Sep 2026 22:50:03 +0700 Subject: [PATCH 3/6] feat(projects): pin a Flow project instead of creating one MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit project.createProject went with the migration — the labs.google tRPC endpoint it lived on is unauthenticated now, and calling it just fails. Every batchexecute call is scoped to a project, so one is made once in the Flow UI and its uuid supplied per project as `flow_project_id` or pinned as FLOW_PROJECT_ID. Reading the id back is split out so it works either way: the batch path answers `{"projectId": …}`, the legacy tRPC path buried it under result/data/json/result. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015JkgMuy89iJ9qc1NkRdPk5 --- agent/api/projects.py | 45 ++++++++++++++++++++++++++++------------- agent/models/project.py | 4 ++++ 2 files changed, 35 insertions(+), 14 deletions(-) diff --git a/agent/api/projects.py b/agent/api/projects.py index c87a76b17..600ebe8a9 100644 --- a/agent/api/projects.py +++ b/agent/api/projects.py @@ -7,7 +7,7 @@ from fastapi import APIRouter, HTTPException from pydantic import BaseModel -from agent.config import BASE_DIR +from agent.config import BASE_DIR, USE_BATCH_RPC from agent.models.project import Project, ProjectCreate, ProjectUpdate from agent.models.character import Character from agent.sdk.persistence.sqlite_repository import SQLiteRepository @@ -125,6 +125,24 @@ async def _detect_user_tier(client) -> str: return "PAYGATE_TIER_ONE" +def _read_flow_project_id(flow_result: dict) -> str: + """Pull the project uuid out of whichever transport answered. + + The batch path answers `{"projectId": …}`; the legacy tRPC path buries it + under result/data/json/result. + """ + data = flow_result.get("data") or {} + if isinstance(data, dict): + if isinstance(data.get("projectId"), str): + return data["projectId"] + try: + return data["result"]["data"]["json"]["result"]["projectId"] + except (KeyError, TypeError): + pass + logger.error("Unexpected Flow response: %s", flow_result) + raise HTTPException(502, "Could not read a Flow project id from the response") + + def _get_repo() -> SQLiteRepository: return SQLiteRepository() @@ -154,25 +172,24 @@ async def create(body: ProjectCreate): detected_tier = await _detect_user_tier(client) - flow_result = await client.create_project(body.name, body.tool_name) - if flow_result.get("error"): - raise HTTPException(502, f"Flow API error: {flow_result['error']}") - - try: - data = flow_result.get("data", {}) - result = data["result"]["data"]["json"]["result"] - flow_project_id = result["projectId"] - except (KeyError, TypeError) as e: - logger.error("Unexpected Flow response: %s", flow_result) - raise HTTPException(502, f"Failed to parse Flow response: {e}") - - logger.info("Flow project created: %s", flow_project_id) + # On the batch path Flow no longer creates projects for us — the uuid comes + # from the request or from FLOW_PROJECT_ID. The legacy path still mints one. + flow_project_id = client.flow_project_id(body.flow_project_id) if USE_BATCH_RPC else None + if flow_project_id: + logger.info("Flow project reused: %s", flow_project_id) + else: + flow_result = await client.create_project(body.name, body.tool_name) + if flow_result.get("error"): + raise HTTPException(502, f"Flow API error: {flow_result['error']}") + flow_project_id = _read_flow_project_id(flow_result) + logger.info("Flow project created: %s", flow_project_id) repo = _get_repo() # Step 2: Create local project with the Flow-assigned ID and detected tier create_data = body.model_dump(exclude_none=True) create_data.pop("tool_name", None) + create_data.pop("flow_project_id", None) create_data.pop("style", None) characters_input = create_data.pop("characters", None) diff --git a/agent/models/project.py b/agent/models/project.py index 1e266b98f..bf3a677fb 100644 --- a/agent/models/project.py +++ b/agent/models/project.py @@ -18,6 +18,10 @@ class ProjectCreate(BaseModel): language: str = "en" user_paygate_tier: PaygateTier = "PAYGATE_TIER_ONE" tool_name: str = "PINHOLE" + # Flow stopped letting clients create projects in the September 2026 + # migration, so a project is made once in the Flow UI and its uuid supplied + # here. Falls back to FLOW_PROJECT_ID when omitted. + flow_project_id: Optional[str] = None material: str = Field("realistic", pattern=r"^[a-z0-9][a-z0-9_]{1,63}$") # material ID from GET /api/materials style: Optional[str] = None # deprecated: use material instead; "3D"→"3d_pixar", "photorealistic"→"realistic" allow_music: bool = False # when True, skip "no background music" suffix in video prompts From 99e3e451acef57b2fc9e6b741e8ccfc8e340b8dc Mon Sep 17 00:00:00 2001 From: Hoang Tuan Nguyen Date: Sun, 6 Sep 2026 22:50:19 +0700 Subject: [PATCH 4/6] fix(worker): capability gaps are terminal, poll complaints are not MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An unported capability and a missing Flow project are configuration answers, not something five retries can reach — UNSUPPORTED_ON_BATCH_API and NO_FLOW_PROJECT now fail once instead of burning the retry budget on every scene of a pipeline run. The other direction for the poll: an operation can report "Media not found." and still deliver a finished clip, so the complaint is carried alongside a still-pending round and only quoted if we time out, where it is the one piece of context worth having. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015JkgMuy89iJ9qc1NkRdPk5 --- agent/sdk/services/operations.py | 9 ++++++++- agent/worker/processor.py | 8 ++++++++ 2 files changed, 16 insertions(+), 1 deletion(-) diff --git a/agent/sdk/services/operations.py b/agent/sdk/services/operations.py index 99c8f8326..23c228df5 100644 --- a/agent/sdk/services/operations.py +++ b/agent/sdk/services/operations.py @@ -115,6 +115,10 @@ async def _poll_operations( poll_interval = VIDEO_POLL_INTERVAL elapsed = 0 current_ops = operations + # The batch path attaches the operation's own grumble ("Media not found.") + # to a still-pending round. It is a diagnostic, not a verdict — finished + # jobs report it too — so it is only worth quoting if we time out. + last_complaint = None while elapsed < timeout: await asyncio.sleep(poll_interval) @@ -136,6 +140,8 @@ async def _poll_operations( error_msg = "" for op in ops: + if op.get("complaint"): + last_complaint = op["complaint"] status = op.get("status", "") if status == "MEDIA_GENERATION_STATUS_SUCCESSFUL": continue @@ -159,7 +165,8 @@ async def _poll_operations( done_count = sum(1 for o in ops if o.get("status") == "MEDIA_GENERATION_STATUS_SUCCESSFUL") logger.debug("Poll %ds/%ds: %d/%d done", elapsed, timeout, done_count, len(ops)) - return {"error": f"Polling timeout after {timeout}s"} + detail = f": {last_complaint}" if last_complaint else "" + return {"error": f"Polling timeout after {timeout}s{detail}"} class OperationService: diff --git a/agent/worker/processor.py b/agent/worker/processor.py index 50297089f..c23039afe 100644 --- a/agent/worker/processor.py +++ b/agent/worker/processor.py @@ -444,6 +444,14 @@ async def _handle_failure(rid: str, req: dict, result: dict, retry_after: dict = error_lower = str(error_msg).lower() + # A capability the batch path does not have, or a missing Flow project, is + # a configuration answer — not something a retry can reach. Fail it once. + if "unsupported_on_batch_api" in error_lower or "no_flow_project" in error_lower: + await crud.update_request(rid, status="FAILED", error_message=str(error_msg)) + await _mark_scene_failed(req) + logger.error("Request %s FAILED (not retryable): %s", rid[:8], error_msg) + return + # WS transient errors (extension disconnect/reconnect): retry without incrementing count if "extension reconnected" in error_lower or "extension disconnected" in error_lower or "extension not connected" in error_lower: await crud.update_request(rid, status="PENDING", error_message=str(error_msg)) From ad9d6854531945defdc3b16d013a0de42380524f Mon Sep 17 00:00:00 2001 From: Hoang Tuan Nguyen Date: Sun, 6 Sep 2026 22:50:19 +0700 Subject: [PATCH 5/6] fix(review): download from the media record's fresh url The two transports answer a media lookup differently: the batch path hands back a freshly signed url, the legacy REST path inlined the bytes as base64. Try the url first and keep the encoded content as the fallback, so an expired scene url still resolves on either path. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015JkgMuy89iJ9qc1NkRdPk5 --- agent/services/video_reviewer.py | 30 +++++++++++++++++++----------- 1 file changed, 19 insertions(+), 11 deletions(-) diff --git a/agent/services/video_reviewer.py b/agent/services/video_reviewer.py index cb52a4c7d..5bf06666c 100644 --- a/agent/services/video_reviewer.py +++ b/agent/services/video_reviewer.py @@ -125,7 +125,12 @@ async def _download_video(url: str, dest: Path) -> None: async def _download_via_get_media(media_id: str, dest: Path) -> None: - """Download video by fetching encoded content from get_media API.""" + """Re-fetch a clip through the media record when its stored url has expired. + + The two transports answer differently: the batch path hands back a freshly + signed url to download, the legacy REST path inlined the bytes as base64. + Try the url first, fall back to the encoded content. + """ from agent.services.flow_client import get_flow_client client = get_flow_client() @@ -134,18 +139,21 @@ async def _download_via_get_media(media_id: str, dest: Path) -> None: raise ValueError(f"get_media failed for {media_id}: {result['error']}") data = result.get("data", result) - # Video content is in video.encodedVideo or image.encodedImage (base64) - encoded = None - if isinstance(data, dict): - if "video" in data and isinstance(data["video"], dict): - encoded = data["video"].get("encodedVideo") - elif "image" in data and isinstance(data["image"], dict): - encoded = data["image"].get("encodedImage") - elif "encodedVideo" in data: - encoded = data["encodedVideo"] + if not isinstance(data, dict): + raise ValueError(f"Unreadable get_media response for {media_id}") + + video = data.get("video") if isinstance(data.get("video"), dict) else {} + image = data.get("image") if isinstance(data.get("image"), dict) else {} + + url = video.get("fifeUrl") or data.get("fifeUrl") + if url: + await _download_video(url, dest) + logger.info("Downloaded %s via a freshly signed url", media_id[:12]) + return + encoded = video.get("encodedVideo") or image.get("encodedImage") or data.get("encodedVideo") if not encoded: - raise ValueError(f"No encoded content in get_media response for {media_id}") + raise ValueError(f"No video url or encoded content in get_media response for {media_id}") video_bytes = base64.standard_b64decode(encoded) with open(dest, "wb") as f: From 323ed8532e75abe3fd029a183ef9147a0366a392 Mon Sep 17 00:00:00 2001 From: Hoang Tuan Nguyen Date: Sun, 6 Sep 2026 22:50:19 +0700 Subject: [PATCH 6/6] docs: follow Flow to its new house MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit What the migration changes for whoever runs this next: a Flow project has to be pinned, one signed-in flow.google.com tab has to stay open, and `flow_key_present: false` is now normal rather than a fault. fk-doctor gains the batch path's error taxonomy — NO_AT_TOKEN, NO_FLOW_PROJECT, UNSUPPORTED_ON_BATCH_API — and loses the advice that sent people looking for a bearer token that is not coming back. The skills whose behaviour actually changed say so where you would hit it: chaining and upscale name what they cannot do and offer the degraded fallback, change-model explains that only the quality tier survives the fold onto the three accepted names, refresh-urls describes what it does now that it works. docs/CAPTURE.md is the way back for the three unported payloads: record the action, diff the slots, throw the capture away. Guessing at Google's positional payloads does not work — a reference image in the wrong slot is accepted and then silently ignored. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_015JkgMuy89iJ9qc1NkRdPk5 --- CLAUDE.md | 25 ++++++++++- README.md | 79 ++++++++++++++++++++++++++++------- docs/CAPTURE.md | 77 ++++++++++++++++++++++++++++++++++ skills/fk-change-model.md | 22 ++++++++++ skills/fk-create-project.md | 21 ++++++++++ skills/fk-doctor.md | 60 ++++++++++++++++++++++---- skills/fk-gen-chain-videos.md | 27 ++++++++++++ skills/fk-pipeline.md | 2 +- skills/fk-refresh-urls.md | 34 +++++++++------ skills/fk-upload-image.md | 8 +++- 10 files changed, 313 insertions(+), 42 deletions(-) create mode 100644 docs/CAPTURE.md diff --git a/CLAUDE.md b/CLAUDE.md index 31938df77..483aa64d5 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -7,14 +7,35 @@ Base URL: `http://127.0.0.1:8100` ```bash curl -s http://127.0.0.1:8100/health # Must return: {"extension_connected": true} + +curl -s http://127.0.0.1:8100/api/flow/status +# Must return: {"transport": "batch", "flow_project_id": "", ...} ``` +Also needed: **one signed-in `https://flow.google.com/` tab left open**. Only the +page can sign a Flow request, so nothing works headless. + ## How to work - Always use `/fk:*` skills — all rules and workflows live inside each skill - Never write scripts to loop API calls — use `POST /api/requests/batch` - `media_id` is always UUID format (`xxxxxxxx-xxxx-xxxx-xxxx-xxxxxxxxxxxx`), never `CAMS...` strings -- **On any pipeline error** (request `FAILED`, stuck `PROCESSING`, `extension_connected: false`, HTTP 4xx/5xx from `:8100`, YouTube `HttpError`, error strings like `UNSAFE_GENERATION` / `not found` / `CAPTCHA` / `NO_FLOW_KEY`): invoke `/fk-doctor` before guessing a fix +- **On any pipeline error** (request `FAILED`, stuck `PROCESSING`, `extension_connected: false`, HTTP 4xx/5xx from `:8100`, YouTube `HttpError`, error strings like `UNSAFE_GENERATION` / `not found` / `CAPTCHA` / `NO_AT_TOKEN` / `NO_FLOW_PROJECT` / `UNSUPPORTED_ON_BATCH_API`): invoke `/fk-doctor` before guessing a fix +- `flow_key_present: false` is **normal** — the current transport has no bearer token + +## Since Flow moved (September 2026) + +Flow lives at `flow.google.com` and signs every call in the page. Consequences +that change how you work: + +- **Projects are not created by Flow Kit any more.** Make one in the Flow UI and + pin its uuid as `FLOW_PROJECT_ID`, or pass `flow_project_id` to `POST /api/projects`. +- **Three capabilities are unported** because their payloads were never captured: + 4K upscale, r2v, and start+end-frame chaining. They fail with + `UNSUPPORTED_ON_BATCH_API` rather than silently producing the wrong thing. + `FLOW_ALLOW_DEGRADED=1` drops chaining and r2v to plain i2v; upscale has no + fallback. To restore one properly, see `docs/CAPTURE.md`. +- **A poll saying "Media not found." is not a failure.** Finished jobs report it. ## Skills @@ -40,7 +61,7 @@ curl -s http://127.0.0.1:8100/health | `/fk-status` | Project status dashboard | | `/fk-switch-project` | Switch active project | | `/fk-fix-uuids` | Fix non-UUID media_ids | -| `/fk-refresh-urls` | Refresh expired GCS URLs | +| `/fk-refresh-urls` | Refresh expired signed media URLs | | `/fk-doctor` | Diagnose errors + prescribe fixes (Flow/extension/worker/YT) | | `/fk-add-material` | Set image material style | | `/fk-change-model` | Change video/image model | diff --git a/README.md b/README.md index ef59a5c9f..249cff91f 100644 --- a/README.md +++ b/README.md @@ -50,7 +50,7 @@ # FLOW KIT -Standalone system to generate AI videos via Google Flow API. Uses a Chrome extension as browser bridge for authentication, reCAPTCHA solving, and API proxying. +Standalone system to generate AI videos via Google Flow. Uses a Chrome extension as a browser bridge: it mints reCAPTCHA and runs Flow's batchexecute RPCs inside a signed-in `flow.google.com` tab, which is the only place they can be signed. ## Showcase @@ -149,17 +149,29 @@ Each project goes through: **story → entities → reference images → scene i ## Architecture ``` -┌──────────────────┐ WebSocket ┌──────────────────────┐ -│ Python Agent │◄──────────────────►│ Chrome Extension │ -│ (FastAPI+SQLite)│ localhost:9222 │ (MV3 Service Worker) │ -│ │ │ │ -│ - REST API :8100│ ── commands ──► │ - Token capture │ -│ - Queue worker │ ◄── results ── │ - reCAPTCHA solve │ -│ - Post-process │ │ - API proxy │ -│ - SQLite DB │ │ (on labs.google) │ -└──────────────────┘ └──────────────────────┘ +┌──────────────────┐ WebSocket ┌──────────────────────┐ ┌──────────────────┐ +│ Python Agent │◄──────────────────►│ Chrome Extension │────►│ flow.google.com │ +│ (FastAPI+SQLite)│ localhost:9222 │ (MV3 Service Worker) │ │ (signed-in tab) │ +│ │ │ │ │ │ +│ - REST API :8100│ ── envelopes ──► │ - reCAPTCHA mint │ │ batchexecute │ +│ - Queue worker │ ◄── responses ── │ - runs the RPC in │ │ cookie + `at` │ +│ - Post-process │ │ the page's world │ │ │ +│ - SQLite DB │ │ │ │ │ +└──────────────────┘ └──────────────────────┘ └──────────────────┘ ``` +Flow signs every call with the session cookie plus a per-page `at` token, and a +generate also carries a single-use reCAPTCHA. None of that can be replayed from +outside the browser, so the agent builds the request and the **page** issues it. +One signed-in Flow tab has to stay open; nothing here works headless. + +> **September 2026 — Flow moved.** It now lives at `flow.google.com` and the old +> `aisandbox-pa.googleapis.com` REST API has no caller: the `Bearer ya29.…` it +> needed stopped being minted. If you are upgrading from an older Flow Kit, +> reload the extension (v0.3.0+) and pin `FLOW_PROJECT_ID` — see +> [Configuration](#configuration). The legacy path is still there behind +> `USE_BATCH_RPC=0`, but it is a post-mortem tool, not a fallback. + ## Quick Start ### One-command setup @@ -183,16 +195,48 @@ pip install -r requirements.txt ```bash # 1. Load Chrome extension: chrome://extensions → Developer mode → Load unpacked → extension/ -# 2. Open https://labs.google/fx/tools/flow and sign in -# 3. Start agent +# 2. Open https://flow.google.com/ and sign in — leave the tab open +# 3. Create a project in the Flow UI and copy its uuid out of the URL +export FLOW_PROJECT_ID= + +# 4. Start agent source venv/bin/activate # if using setup.sh python -m agent.main -# 4. Verify +# 5. Verify curl http://127.0.0.1:8100/health # {"status":"ok","extension_connected":true} +curl http://127.0.0.1:8100/api/flow/status +# {"connected":true,"transport":"batch","flow_project_id":"…","flow_key_present":false} ``` +`flow_key_present: false` is expected — the current transport has no bearer +token. Step 3 is not optional: Flow's project-creation endpoint went with the +migration, so without a pinned project every request fails `NO_FLOW_PROJECT`. +You can also pass `flow_project_id` per project on `POST /api/projects`. + +### Configuration + +| Env var | Default | What it does | +|---------|---------|--------------| +| `FLOW_PROJECT_ID` | — | The Flow project every RPC is scoped to. Required. | +| `USE_BATCH_RPC` | `1` | `0` falls back to the pre-migration REST path (dead auth). | +| `FLOW_ALLOW_DEGRADED` | `0` | `1` lets scene chaining and r2v fall back to plain i2v instead of failing. | +| `DEFAULT_PAYGATE_TIER` | `PAYGATE_TIER_TWO` | Carried for the DB and dashboard; no longer selects a model. | + +### What does not work on the new API yet + +Three capabilities have no captured payload, so they fail with +`UNSUPPORTED_ON_BATCH_API` rather than quietly producing the wrong thing: + +| Capability | Status | Workaround | +|---|---|---| +| 4K/1080p upscale (`/fk-pipeline` last step) | unported | none — keep the 1080p render | +| Reference-to-video (r2v) | unported | `FLOW_ALLOW_DEGRADED=1` → i2v off the first reference | +| Start+end-frame chaining (`/fk-gen-chain-videos`) | unported | `FLOW_ALLOW_DEGRADED=1` → i2v off the start frame | + +Restoring one starts with a capture, not a guess: [`docs/CAPTURE.md`](docs/CAPTURE.md). + ## End-to-End Example: "Pippip the Fish Merchant" A chubby cat sells fish at a market. 3 scenes, vertical, Pixar 3D style. @@ -679,7 +723,7 @@ These arrive in the response body as `data.error.details[].reason`. The worker a | Status | Source | Meaning | Handling | |--------|--------|---------|----------| | **400** | Flow API | Invalid payload, UNSAFE_GENERATION, entity not found (sometimes) | Route by `details.reason` — some are auto-recoverable, others terminal | -| **401** | Flow API | Bearer token expired | Extension re-captures token from labs.google tab; request retries | +| **401** | Flow API (legacy path only) | Bearer expired — on a post-migration profile it was never minted | Switch to the batch path (`USE_BATCH_RPC=1`) | | **403** | Extension (`background.js:432`) | `CAPTCHA_FAILED`, `NO_FLOW_TAB`, or `MODEL_ACCESS_DENIED` | CAPTCHA → retry loop; NO_FLOW_TAB → fail (user must open Flow); tier → fail | | **404** | Flow API | `media_id` not found (expired upload) | Same as "Requested entity was not found" — auto re-upload | | **429** | Flow API | Rate limited / quota | Back off + retry; if `USER_QUOTA_REACHED` appears, fail | @@ -699,7 +743,10 @@ String patterns in `error_message` that the worker recognizes: | `Extension not connected` | Chrome extension offline or WS dropped | 503 returned; worker re-queues PENDING and waits | | `extension reconnected` / `extension disconnected` | WS bounce mid-request | Re-queue PENDING without incrementing `retry_count` | | `extension_switched` | User switched Flow tabs mid-generation | Re-queue PENDING | -| `NO_FLOW_KEY` | Extension has no captured bearer token | User must open `labs.google/fx/tools/flow` and sign in | +| `NO_FLOW_KEY` | Extension has no captured bearer token — **legacy path only**, expected on the batch path | Only meaningful with `USE_BATCH_RPC=0` | +| `NO_AT_TOKEN` | Flow tab is signed out, on an interstitial, or still booting | Open `flow.google.com`, sign in, let the app load | +| `NO_FLOW_PROJECT` | No Flow project to scope the RPC to | Pin `FLOW_PROJECT_ID` — **terminal, not retried** | +| `UNSUPPORTED_ON_BATCH_API` | Upscale / r2v / chaining — payload never captured | See `docs/CAPTURE.md` — **terminal, not retried** | | `NO_FLOW_TAB` | No Google Flow tab available for reCAPTCHA | User must open a Flow tab | | `Failed to fetch` | Network drop inside extension service worker | Retry with backoff | | `timeout` / WS 60s no response | Extension hung mid-request | Re-queue PENDING | @@ -731,7 +778,7 @@ From `youtube/upload.py` (HTTP errors from YouTube Data API v3): | Problem | Solution | |---------|----------| | Extension shows "Agent disconnected" | Start `python -m agent.main` | -| Extension shows "No token" | Open `labs.google/fx/tools/flow` and sign in | +| Extension shows "No token" | Expected on the batch path — there is no bearer token any more | | `CAPTCHA_FAILED: NO_FLOW_TAB` | Open a Google Flow tab | | 403 `MODEL_ACCESS_DENIED` | Tier mismatch — check `/api/flow/credits`, downgrade model in `models.json` | | Scene images inconsistent | Check all refs have UUID `media_id` — run `/fk-fix-uuids` | diff --git a/docs/CAPTURE.md b/docs/CAPTURE.md new file mode 100644 index 000000000..366e6169b --- /dev/null +++ b/docs/CAPTURE.md @@ -0,0 +1,77 @@ +# Capturing a Flow batchexecute payload + +`agent/services/flow_batch.py` only knows the RPC shapes that were captured off a +real UI action. Adding one — video upscale, reference-to-video, start+end-frame +chaining, a base-image edit — starts by watching the browser do it, because +guessing at Google's positional payloads does not work. Thirty generations were +spent proving that a reference image in the wrong slot is *accepted* and then +silently ignored. + +The recorder is deliberately **not** in the shipped extension: it writes live +request bodies to disk, so it goes in for one session and comes straight back out. + +## What already exists + +| step | rpcid | notes | +|---|---|---| +| generate image | `ogiZ0b` | signed CDN url comes back inline | +| generate video | `eb1hJf` | returns an operation id | +| poll operation | `jwpduf` | status `CAE` means finished | +| operation → media id | `Zzl0ze` | `projects/`; the listing is ~17 MB | +| media id → urls | `as29s` | signed `/video/` + poster `/image/` | +| upload an image | `maseQ` | base64 in the payload, captcha like a generate | + +Missing, and each blocked behind a capture: **video upscale**, **r2v**, +**start+end-frame chaining**, and the **base-image** variant of the image edit. + +## Recording one + +1. Add to `extension/background.js`, temporarily: + +```js +const NETLOG_HOSTS = ['https://flow.google.com/_/*']; +const pending = new Map(); + +chrome.webRequest.onBeforeRequest.addListener((d) => { + const body = d.requestBody?.raw?.length + ? new TextDecoder().decode(new Uint8Array(d.requestBody.raw[0].bytes)) + : null; + pending.set(d.requestId, { ts: new Date().toISOString(), url: d.url, body }); +}, { urls: NETLOG_HOSTS }, ['requestBody']); + +chrome.webRequest.onCompleted.addListener((d) => { + const rec = pending.get(d.requestId); + if (!rec) return; + pending.delete(d.requestId); + fetch('http://127.0.0.1:8100/api/ext/netlog', { + method: 'POST', headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ ...rec, statusCode: d.statusCode }), + }).catch(() => {}); +}, { urls: NETLOG_HOSTS }); +``` + + Headers are omitted on purpose — the payload shape is what is wanted, and the + credentials are what caused the mess. Add an `onBeforeSendHeaders` listener + only if a header itself is the open question, and delete the log afterwards. + +2. Add a matching throwaway route to `agent/main.py` that appends the posts to a + file under the scratch directory. +3. Reload the extension, perform **one** action in Flow, then reload again with + the listener removed. +4. Read the `f.req` out of the capture: it is + `[[[rpcid, "", null, "generic"]]]`. Diff the + inner payload against the closest builder in `flow_batch.py` to find which + slot changed. +5. Delete the capture file. It is evidence, not a fixture — put what you learned + into a builder and a test instead. + +## Two things worth knowing before you diff + +**Accepted ≠ used.** A wrong arrangement inside a slot comes back 200 and is +then ignored. Prove a reference image with a prompt that never names its +subject; prove an aspect ratio by reading the JPEG header, not by trusting the +field name. + +**Slots do not share encodings.** Image aspect is 1 square / 2 portrait / +3 landscape / 4 is 3:4 / 5 is 4:3. Video aspect is 1 portrait / 2 landscape, in +its own slot. Conflating them renders the wrong shape silently. diff --git a/skills/fk-change-model.md b/skills/fk-change-model.md index 2ab0be41d..c3b8223e9 100644 --- a/skills/fk-change-model.md +++ b/skills/fk-change-model.md @@ -11,6 +11,28 @@ Usage: --- +## What the model keys mean on the current Flow API + +Since Flow moved to `flow.google.com`, aspect ratio is its own payload slot and +the REST-era suffixed names (`…_portrait`, `…_fl`, `…_relaxed`) are **rejected +outright**. `models.json` still stores the old [tier][gen_type][aspect] map for +the legacy path, and `resolve_video_model()` folds whatever it finds onto the +three names the new path accepts: + +| Wire name | What it is | +|---|---| +| `veo_3_1_i2v_s_fast_ultra` | anything whose key mentions `ultra` | +| `veo_3_1_i2v_lite` | anything whose key mentions `lite` | +| `veo_3_1_i2v_lite_low_priority` | the default, and anything unrecognised | + +So changing a Landscape vs Portrait key has **no effect** on the batch path — +only the quality tier survives the fold. `r2v` and `start_end` keys are moot +there too: both are unported. + +Image models are unaffected: `GEM_PIX_2` (Nano Banana Pro) and `NARWHAL` +(Banana 2) are both accepted, and `default_image_model` in `models.json` picks +which nickname is used. + ## Step 1: Show Current Models ```bash diff --git a/skills/fk-create-project.md b/skills/fk-create-project.md index 803372cf7..e72c6a626 100644 --- a/skills/fk-create-project.md +++ b/skills/fk-create-project.md @@ -104,8 +104,29 @@ Camera stays behind. Viewers see the leader's power through body language, not f | Military rank + origin | The Field Marshal, The Admiral | Military figures | | Generic role | The Royal Advisor, The Strategist | Secondary characters | +## Step 0: Make sure there is a Flow project to attach to + +Since Flow moved to `flow.google.com`, Flow Kit cannot create Flow projects — +the endpoint that did it went with the migration. Every generation is scoped to +an existing one. + +```bash +curl -s http://127.0.0.1:8100/api/flow/status | python3 -c " +import sys, json +s = json.load(sys.stdin) +print('Flow project:', s.get('flow_project_id') or 'NONE — create one in the Flow UI') +" +``` + +If it prints NONE, ask the user to open `https://flow.google.com/`, create a +project, and copy the uuid out of the URL. Then either pin it +(`export FLOW_PROJECT_ID=` before starting the agent) or pass it as +`flow_project_id` in Step 1. Without it every request fails `NO_FLOW_PROJECT`. + ## Step 1: Create project with all entities +Add `"flow_project_id": ""` if you are not using the pinned one. + ```bash curl -X POST http://127.0.0.1:8100/api/projects \ -H "Content-Type: application/json" \ diff --git a/skills/fk-doctor.md b/skills/fk-doctor.md index 8891f54bd..f02ddd61b 100644 --- a/skills/fk-doctor.md +++ b/skills/fk-doctor.md @@ -6,7 +6,7 @@ Diagnose any FlowKit error and prescribe a fix. Knows the full error taxonomy ac - Any `/api/requests/*` response has `status=FAILED` or `error_message` is set - A request has been `PROCESSING` for > 10 minutes with no progress - `GET /health` returns `extension_connected: false` -- User reports any error string containing: `UNSAFE_GENERATION`, `QUOTA`, `not found`, `CAPTCHA`, `NO_FLOW_KEY`, `NO_FLOW_TAB`, `extension_switched`, `Failed to fetch`, `MODEL_ACCESS_DENIED`, `PAYGATE_TIER_TWO`, `invalidTags`, `quotaExceeded`, `invalid_grant` +- User reports any error string containing: `UNSAFE_GENERATION`, `QUOTA`, `not found`, `CAPTCHA`, `NO_AT_TOKEN`, `NO_FLOW_PROJECT`, `UNSUPPORTED_ON_BATCH_API`, `NO_FLOW_KEY`, `NO_FLOW_TAB`, `FLOW_TAB_DISCARDED`, `extension_switched`, `Failed to fetch`, `MODEL_ACCESS_DENIED`, `PAYGATE_TIER_TWO`, `invalidTags`, `quotaExceeded`, `invalid_grant` - User asks "why did X fail", "what's wrong with the pipeline", "why is this stuck", "tại sao X lỗi", "lỗi gì vậy" - An HTTP 4xx/5xx reaches the main agent from any endpoint under `127.0.0.1:8100` - A YouTube upload returns `HttpError` from `googleapiclient` @@ -67,6 +67,26 @@ Cross-reference `error_message` against the taxonomy below. Print: **Diagnosis / Match against taxonomy — even partial matches (`"not found"`, `"captcha"`, `"quota"`). +## Which transport is running + +Since Flow moved to `flow.google.com` (September 2026) there are two paths, and +the taxonomy below splits on which one is live: + +```bash +python3 -c "from agent.config import USE_BATCH_RPC, FLOW_PROJECT_ID; \ + print('batch' if USE_BATCH_RPC else 'legacy REST', '| project:', FLOW_PROJECT_ID or 'UNPINNED')" +``` + +- **batch** (default) — the agent builds an `f.req` envelope, the extension runs + it inside a signed-in `flow.google.com` tab. No bearer token exists on this + path, so `flow_key_present: false` in `/api/flow/status` is **normal**, not a + fault. Requires a Flow tab open and `FLOW_PROJECT_ID` pinned. +- **legacy REST** (`USE_BATCH_RPC=0`) — the pre-migration `aisandbox-pa` path. + It needs a `Bearer ya29.…` Flow no longer mints, so it will 401 on any fresh + profile. Treat any report of it "suddenly breaking" as the migration, not a + regression: if `token_age_s` only climbs across tab reloads, the token is not + stale, it is gone. + ## Error Taxonomy ### A. Flow-native structured errors (from `data.error.details[].reason`) @@ -85,7 +105,7 @@ Match against taxonomy — even partial matches (`"not found"`, `"captcha"`, `"q | Status | Origin | When you see it | |--------|--------|-----------------| | **400** | Flow API | Invalid payload / UNSAFE / entity-not-found — **route by `details.reason`** | -| **401** | Flow API | Bearer token expired — extension should auto-recapture from labs.google tab | +| **401** | Flow API (legacy path only) | Bearer expired — and on a post-migration profile it is not expired, it was never minted. Switch to the batch path | | **403** | Extension (`background.js:432`) | `CAPTCHA_FAILED`, `NO_FLOW_TAB`, or `MODEL_ACCESS_DENIED` — read the suffix | | **404** | Flow API | `media_id` not found — same handler as "entity not found" | | **429** | Flow API | Rate-limit or quota — backoff; if message mentions QUOTA_REACHED, terminal | @@ -106,11 +126,33 @@ Detection lives in `agent/worker/_parsing.py:_is_error`. A result is treated as | `Extension not connected` | WS dropped or extension offline | Reload extension at `chrome://extensions`; worker auto-retries | | `extension reconnected` / `extension disconnected` | WS bounce mid-request | Auto re-queue, `retry_count` NOT incremented | | `extension_switched` | User switched active Flow tab | Auto re-queue | -| `NO_FLOW_KEY` | No bearer token captured | Open `labs.google/fx/tools/flow` and sign in | -| `NO_FLOW_TAB` | No Flow tab for CAPTCHA solve | Open any Flow tab | +| `NO_FLOW_KEY` | No bearer token captured — **legacy path only**; expected and harmless on the batch path | Only meaningful with `USE_BATCH_RPC=0`; otherwise ignore | +| `NO_FLOW_TAB` | No Flow tab for CAPTCHA solve or RPC signing | Open `https://flow.google.com/` and sign in | | `Failed to fetch` | Network drop inside service worker | Auto-retry with backoff | | WS 60s timeout | Extension hung | Reload extension; worker re-queues | +### C2. Batch path (`flow.google.com`) errors + +| Error contains | Cause | Auto-handling | Fix | +|----------------|-------|---------------|-----| +| `NO_AT_TOKEN` | The Flow tab loaded but `WIZ_global_data.SNlM0e` is absent — the page is signed out, on an interstitial, or still booting | Retried with backoff | Open `https://flow.google.com/`, confirm you are signed in, let the app finish loading | +| `NO_FLOW_TAB` | No Flow tab to sign the request | Extension opens one and retries once | Leave one signed-in Flow tab open; nothing here works headless | +| `FLOW_TAB_DISCARDED` | Chrome discarded the backgrounded tab and the reload did not revive it | Retried with backoff | Pin the Flow tab, or keep its window visible | +| `NO_FLOW_PROJECT` | No Flow project to scope the RPC to | **Terminal — not retried** | Create a project in the Flow UI, pin its uuid as `FLOW_PROJECT_ID` (or pass `flow_project_id` on `POST /api/projects`) | +| `UNSUPPORTED_ON_BATCH_API` | A capability whose payload was never captured off the new UI: **video upscale**, **r2v**, **start+end-frame chaining** | **Terminal — not retried** | For chaining and r2v, `FLOW_ALLOW_DEGRADED=1` falls back to plain i2v off the start frame. Upscale has no fallback. Real fix: capture the payload — `docs/CAPTURE.md` | +| `PUBLIC_ERROR_UNUSUAL_ACTIVITY` | A reCAPTCHA token was replayed — they are single-use | Retried as a captcha error | Usually self-clears; if it persists the extension is reusing a token, reload it | +| `no ogiZ0b envelope in response` | The RPC answered but not with the payload we came for — usually a signed-out page returning an HTML redirect | Retried with backoff | Re-sign in on the Flow tab | +| `Polling timeout after Ns: Media not found.` | The job never produced media inside the budget | Terminal after `MAX_RETRIES` | The quoted complaint is a **diagnostic, not the cause** — finished jobs report it too. Check the Flow UI: if the clip is there, raise `VIDEO_POLL_TIMEOUT` | + +Two behaviours on this path routinely look like bugs and are not: + +- **A poll can say "Media not found." and the job still finishes.** The project + listing is what decides; the poll is a hint. Never treat the complaint as fatal. +- **A media id arrives before the clip is fetchable.** The media record serves + the poster image first and grows the `/video/` url in later, so a scene sits + PENDING for a while after its id exists. Downloading on the id alone saves a + still picture. + ### D. YouTube upload errors (`youtube/upload.py`) | Error | Cause | Fix | @@ -137,15 +179,15 @@ When the user describes a symptom in plain language, map it here first. | Problem | Solution | |---------|----------| | Extension shows "Agent disconnected" | Start `python -m agent.main` | -| Extension shows "No token" | Open `labs.google/fx/tools/flow` and sign in | -| `CAPTCHA_FAILED: NO_FLOW_TAB` | Open a Google Flow tab | +| Extension shows "No token" | Expected on the batch path — there is no bearer any more. Only act on it with `USE_BATCH_RPC=0` | +| `CAPTCHA_FAILED: NO_FLOW_TAB` | Open `https://flow.google.com/` — and check the extension is v0.3.0+, older builds only matched the dead labs.google URL and could not see the tab that was right there | | 403 `MODEL_ACCESS_DENIED` | Tier mismatch — `GET /api/flow/credits`, downgrade model in `models.json` via `/fk-change-model` | | Scene images inconsistent across scenes | Check all refs have UUID `media_id` — run `/fk-fix-uuids` | | `media_id` starts with `CAMS...` | Run `/fk-fix-uuids` to extract UUID from URL | -| Upscale "permission denied" | Requires `PAYGATE_TIER_TWO` account — TIER_ONE cannot upscale | +| Upscale fails on every scene | On the batch path upscale is unported (`UNSUPPORTED_ON_BATCH_API`) — no upsampler rpc has been captured. On the legacy path it needs `PAYGATE_TIER_TWO` | | Request stuck in PROCESSING > 10 min | Check `error_message` history; if extension dropped, reload it at `chrome://extensions` | | "Requested entity was not found" spam | Image URLs expired — re-upload via `POST /api/upload-image` or wait for `_recover_entity_not_found` | -| Expired GCS signed URLs | Run `/fk-refresh-urls` to regenerate | +| Expired signed URLs | Run `/fk-refresh-urls` — on the batch path this re-signs every stored media id through the media rpc | | YouTube upload `invalidTags` | Tag-char overflow — quote overhead counts (spaces → +2 per tag) | | Python `cryptography` arch mismatch | Use `python3.10`, not `python3.13` (x86/arm64 binary mismatch) | | `curl: (7) Failed to connect to 127.0.0.1:8100` | Agent not running — `python -m agent.main` | @@ -154,6 +196,7 @@ When the user describes a symptom in plain language, map it here first. Decision order — stop at first match: +0. **`UNSUPPORTED_ON_BATCH_API` / `NO_FLOW_PROJECT`** → FAILED immediately. These are configuration answers, not something a retry can reach. 1. **`"not found"` in message** → `_recover_entity_not_found()` re-uploads media, marks PENDING. 2. **`reconnected` / `disconnected` / `switched`** → PENDING, keep `retry_count`. 3. **`captcha` / `recaptcha`** → PENDING if retry_count < 10; else FAILED. @@ -168,6 +211,7 @@ Always end with a prescription block: Symptom: Root cause: Layer: Flow | Extension | FastAPI | Worker | YouTube | Env +Transport: batch (flow.google.com) | legacy REST (aisandbox-pa) Auto-handler: === FIX === diff --git a/skills/fk-gen-chain-videos.md b/skills/fk-gen-chain-videos.md index 0781ec629..e0e58b882 100644 --- a/skills/fk-gen-chain-videos.md +++ b/skills/fk-gen-chain-videos.md @@ -4,6 +4,33 @@ Usage: `/gen-chain-videos ` This creates smooth transitions between scenes in a chain by using the **NEXT scene's image as the endImage** of the current scene's video, so the last frame of scene N matches the first frame of scene N+1 → seamless concat. +## Before you start: chaining is not on the new Flow API + +Flow's `flow.google.com` payload has an aspect slot and a single source-image +slot; the **end-image slot was never captured**, so a start+end frame request +cannot be built. `GENERATE_VIDEO` with an `endImage` fails immediately with +`UNSUPPORTED_ON_BATCH_API` (terminal — it is not retried). + +```bash +curl -s http://127.0.0.1:8100/api/flow/status | python3 -c " +import sys, json +s = json.load(sys.stdin) +print('transport:', s['transport'], '| degraded fallback:', s['allow_degraded']) +" +``` + +Two honest options — tell the user which one you are taking: + +1. **`FLOW_ALLOW_DEGRADED=1`** (restart the agent with it): each scene renders as + plain i2v off its own start frame. The pipeline completes, but the cut + between scenes is **not** seamless — the chain invariant below does not hold. + Prefer `/fk-concat-fit-narrator` with a crossfade to hide the seams. +2. **Restore chaining properly** by capturing the end-image slot off the Flow + UI — `docs/CAPTURE.md`. This is the only way to get the real behaviour back. + +Everything below describes the intended behaviour, which is what option 2 +restores. + ## How chaining works For any scene that has a CHILD in the chain (i.e. some other scene with `parent_scene_id == this.id`): diff --git a/skills/fk-pipeline.md b/skills/fk-pipeline.md index b09d1373e..4b09b6bc0 100644 --- a/skills/fk-pipeline.md +++ b/skills/fk-pipeline.md @@ -5,7 +5,7 @@ Auto-detect project state and run the correct stages (continuation or full run). Usage: `/fk-pipeline [project_id] [orientation] [options]` Options: -- `--upscale` — include 4K upscale stage (TIER_TWO only) +- `--upscale` — include 4K upscale stage. **Unavailable since Flow moved** — no upsampler rpc has been captured on `flow.google.com`, so every upscale fails `UNSUPPORTED_ON_BATCH_API` (terminal, not retried). Warn the user and run without it; the 1080p render is the deliverable. See `docs/CAPTURE.md`. - `--tts` — include TTS narration stage (parallel with upscale) - `--download` — auto-download 4K files as upscales complete - `--concat` — run concat after all stages done diff --git a/skills/fk-refresh-urls.md b/skills/fk-refresh-urls.md index 7415b864f..5c364eb18 100644 --- a/skills/fk-refresh-urls.md +++ b/skills/fk-refresh-urls.md @@ -1,4 +1,4 @@ -Refresh expired GCS signed URLs for all scenes in a video (images, videos, upscale videos) and character reference images. +Re-sign expired media URLs for all scenes in a video (images, videos, upscale videos) and character reference images. Usage: `/fk-refresh-urls [--project-id ]` @@ -11,10 +11,10 @@ Usage: `/fk-refresh-urls [--project-id ]` ## Pre-flight ```bash -# Extension must be connected with flow key curl -s http://127.0.0.1:8100/api/flow/status -# Must show: {"connected": true, "flow_key_present": true} -# If flow_key_present is false: open/refresh a Google Flow tab in Chrome +# Must show: {"connected": true, "transport": "batch"} +# Ignore flow_key_present — the batch path has no bearer token. +# If connected is false: open https://flow.google.com/ and sign in. ``` ## Step 1: Get project_id from video @@ -25,9 +25,11 @@ PID=$(curl -s "http://127.0.0.1:8100/api/videos/${VID}" | python3 -c "import sys echo "Project: $PID" ``` -## Step 2: Bulk refresh via TRPC +## Step 2: Bulk re-sign every stored media id -This calls Google Flow's TRPC `flow.getFlow` endpoint, extracts ALL fresh signed URLs from the response, and updates scenes + characters in DB. +This walks the project's scenes and entities, asks Flow's media rpc to re-sign +each `*_media_id` it holds, and writes the fresh urls back to the DB. It is a +call per media id, so a large project takes a moment. ```bash curl -s -X POST "http://127.0.0.1:8100/api/flow/refresh-urls/${PID}" | python3 -c " @@ -87,16 +89,21 @@ else: " ``` -## Step 4: Per-media fallback (if TRPC fails) +## Step 4: Per-media fallback (for anything the bulk pass missed) -If the bulk TRPC refresh doesn't cover all media (e.g., TRPC response is partial), fall back to per-media refresh: +If a media id was not covered — because it is not stored on a scene or entity +row — re-sign it directly: ```bash -# Get a fresh URL for a specific media_id curl -s "http://127.0.0.1:8100/api/flow/media/" -# Returns: {fifeUrl: "https://...", servingUri: "https://...", ...} +# Returns: {"video": {"fifeUrl": "https://flow-content.google/video/…"}, +# "image": {"fifeUrl": "https://flow-content.google/image/…"}} ``` +A record with only `image` and no `video` means the clip is not finished being +written yet — the poster arrives before the video url does. Wait and retry; +do not save the poster as the video. + Then update the scene manually: ```bash @@ -109,8 +116,9 @@ curl -X PATCH "http://127.0.0.1:8100/api/scenes/" \ | Issue | Cause | Fix | |-------|-------|-----| -| `flow_key_present: false` | Extension hasn't captured auth token | Open/refresh a Google Flow tab in Chrome | -| `Extension not connected` | Chrome extension WS disconnected | Check Chrome extension is enabled, refresh Flow tab | -| `refreshed: 0` | TRPC returned no URLs | Project may not exist on Google Flow, or auth expired | +| `flow_key_present: false` | No bearer token — **expected on the batch path** | Ignore; only meaningful with `USE_BATCH_RPC=0` | +| `Extension not connected` | Chrome extension WS disconnected | Check the extension is enabled, refresh the Flow tab | +| `refreshed: 0`, `found: 0` | No media ids stored for this project | Nothing to refresh — check the project id | +| `refreshed: 0`, `found: N` | Every re-sign failed | Read the agent log; usually `NO_FLOW_TAB` or a signed-out Flow tab | | Some URLs still expired after refresh | media_id mismatch (upscale overwrote video_media_id) | Use per-media fallback with correct media_id | | `get_media` returns error for media_id | Media deleted or expired on Google's side | Re-generate the video/image | diff --git a/skills/fk-upload-image.md b/skills/fk-upload-image.md index e852f6a74..5aad8457e 100644 --- a/skills/fk-upload-image.md +++ b/skills/fk-upload-image.md @@ -9,13 +9,17 @@ Useful for: setting channel icons, covers, or any local image as an entity refer ```bash curl -s http://127.0.0.1:8100/health ``` -Must have `extension_connected: true` AND flow key present. Abort if not. +Must have `extension_connected: true`. Abort if not. ```bash curl -s http://127.0.0.1:8100/api/flow/status -# Must return: {"connected": true, "flow_key_present": true} +# Must return: {"connected": true, "transport": "batch", "flow_project_id": ""} +# flow_key_present is a legacy-path signal — false is expected here. ``` +The upload is scoped to a Flow project: pass `project_id`, or leave it out to +use the pinned `FLOW_PROJECT_ID`. Without either you get `NO_FLOW_PROJECT`. + ## Step 2: Upload image ```bash