diff --git a/CLAUDE.md b/CLAUDE.md index 075463575..3b76e2197 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`. +- **Four capabilities are unported** because their payloads were never captured: + 4K upscale, r2v, start+end-frame chaining, and Omni Flash. 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 c94e3a18d..e4dd67eab 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 @@ -169,17 +169,29 @@ A local React dashboard (`dashboard/`) for monitoring and driving the pipeline ## 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 @@ -203,16 +215,49 @@ 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 | +| Omni Flash (`model_family=omni_flash`) | unported | use `model_family=veo` | + +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. @@ -751,7 +796,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 | @@ -771,7 +816,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 | @@ -803,7 +851,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` | | 403 `PUBLIC_ERROR_UNUSUAL_ACTIVITY` / `reCAPTCHA evaluation failed` | Pause submits, clear cookies for `google.com` + `labs.google` in Chrome, sign back in, then resubmit with ≥1s gap and ≤5 concurrent. Switch network or wait 1–6 h if still blocked | diff --git a/agent/api/flow.py b/agent/api/flow.py index e215fbb31..99394096a 100644 --- a/agent/api/flow.py +++ b/agent/api/flow.py @@ -3,6 +3,7 @@ from pydantic import BaseModel from typing import Literal, Optional +from agent.config import USE_BATCH_RPC, FLOW_PROJECT_ID, FLOW_ALLOW_DEGRADED from agent.services.flow_client import get_flow_client from agent.services.omni_flash import ( check_omni_flash_status, @@ -97,10 +98,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 4b1f35885..43634f208 100644 --- a/agent/api/models.py +++ b/agent/api/models.py @@ -36,6 +36,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("") @@ -68,6 +69,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", 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/config.py b/agent/config.py index 0e779f2de..5d6edf2fe 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 c44c789b5..8c690e59f 100644 --- a/agent/models.json +++ b/agent/models.json @@ -56,5 +56,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/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 diff --git a/agent/sdk/services/operations.py b/agent/sdk/services/operations.py index 1cb9ecbe1..b0199e82e 100644 --- a/agent/sdk/services/operations.py +++ b/agent/sdk/services/operations.py @@ -40,7 +40,7 @@ def _char_matches(c: dict, name_set: set) -> bool: import aiohttp from agent.db import crud -from agent.config import VIDEO_POLL_INTERVAL, VIDEO_POLL_TIMEOUT +from agent.config import USE_BATCH_RPC, VIDEO_POLL_INTERVAL, VIDEO_POLL_TIMEOUT from agent.utils.paths import scene_4k_path from agent.utils.slugify import slugify from agent.worker._parsing import ( @@ -248,6 +248,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) @@ -269,6 +273,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 @@ -292,7 +298,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: @@ -449,8 +456,14 @@ async def generate_scene_video(self, scene: dict, orientation: str, req_row = await crud.get_request(request_id) existing_op = req_row.get("request_id") if req_row else None - # Heuristic: bare UUID = workflow name → skip shortcut. Slash/colon = old operation path. - looks_like_workflow_uuid = bool(existing_op and len(existing_op) == 36 and existing_op.count("-") == 4) + # A bare uuid means different things on the two transports. On the + # legacy path it is a Low Priority workflow name that cannot be + # re-polled, so the retry resubmits. On the batch path it is the + # operation id, and looking it up in the project listing is exactly + # what the status poll does — resubmitting there would abandon a + # running render and pay for a second one. + bare_uuid = bool(existing_op and len(existing_op) == 36 and existing_op.count("-") == 4) + looks_like_workflow_uuid = bare_uuid and not USE_BATCH_RPC if existing_op and not looks_like_workflow_uuid: logger.info("Video gen already submitted (op=%s), re-polling", existing_op[:30]) operations = [{"operation": {"name": existing_op}, "status": "MEDIA_GENERATION_STATUS_PENDING"}] 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 8ba0e542e..667e2d51d 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__) @@ -29,6 +43,13 @@ def __init__(self): self._pending: dict[str, asyncio.Future] = {} self._pending_ws: dict[str, object] = {} 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 @@ -136,6 +157,11 @@ def _should_failover(result: dict) -> bool: return any(marker in message for marker in ( "no_flow_key", "no_flow_tab", + # Batch path: this profile's Flow tab cannot sign a request — it is + # signed out, still booting, or Chrome discarded it. Another + # profile's tab may be perfectly able to. + "no_at_token", + "flow_tab_discarded", "no current window", "extension not connected", "extension disconnected", @@ -235,7 +261,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. @@ -297,18 +326,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. @@ -320,9 +396,16 @@ async def _send(self, method: str, params: dict, timeout: float = 300) -> dict: if not self.connected: return {"error": "Extension not connected"} - extension_candidates = self._extension_candidates(require_token=True) - if not extension_candidates: + # A bearer token is only worth routing on when something is going to + # send one. The batchexecute path authenticates in the page with the + # session cookie, so demanding a flow key there would reject every + # profile — no profile on that path ever captures one. + needs_token = not USE_BATCH_RPC + extension_candidates = self._extension_candidates(require_token=needs_token) + if not extension_candidates and needs_token: return {"error": "NO_FLOW_KEY"} + if not extension_candidates: + return {"error": "Extension not connected"} last_result = {"error": "Extension not connected"} for index, extension_ws in enumerate(extension_candidates): @@ -387,9 +470,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. @@ -407,7 +908,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: @@ -456,7 +957,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", @@ -502,7 +1003,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, @@ -548,7 +1049,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: @@ -594,7 +1095,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.""" @@ -627,7 +1128,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") @@ -638,7 +1139,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", { @@ -647,17 +1148,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 @@ -670,7 +1161,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. @@ -708,6 +1199,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/agent/services/omni_flash.py b/agent/services/omni_flash.py index 2324b812f..0db4c8cbe 100644 --- a/agent/services/omni_flash.py +++ b/agent/services/omni_flash.py @@ -26,11 +26,30 @@ from pathlib import Path from urllib.parse import quote +from agent.config import USE_BATCH_RPC from agent.services.flow_client import get_flow_client from agent.services.headers import random_headers _MODELS_FILE = Path(__file__).parent.parent / "models.json" +#: Every Omni surface here rides the pre-migration transports — the REST +#: endpoints on aisandbox-pa and the labs.google tRPC snapshot it polls +#: through. Flow moved to flow.google.com in September 2026 and stopped +#: minting the bearer both of those need, and no Omni payload has been +#: captured off the new frontend, so on the batch path these fail with a +#: name rather than dying on a 401 five retries deep. +_UNSUPPORTED_ON_BATCH = ( + "UNSUPPORTED_ON_BATCH_API: Omni Flash — it speaks the pre-migration REST " + "and tRPC endpoints, and no batchexecute payload for it has been captured; " + "see docs/CAPTURE.md. Use the Veo path (model_family=veo), or set " + "USE_BATCH_RPC=0 on a profile that still holds a bearer token." +) + + +def _batch_path_blocks_omni() -> dict | None: + """The error to return instead of reaching for auth that is gone.""" + return {"error": _UNSUPPORTED_ON_BATCH} if USE_BATCH_RPC else None + OMNI_FLASH_VALID_DURATIONS = (4, 6, 8, 10) OMNI_FLASH_VALID_ASPECTS = { "VIDEO_ASPECT_RATIO_PORTRAIT", @@ -211,6 +230,9 @@ async def _submit_omni_frame_video( seed: int | None = None, ) -> dict: """Submit Omni first-frame or First+Last generation.""" + blocked = _batch_path_blocks_omni() + if blocked: + return blocked _validate_frame_inputs( start_image_media_id, end_image_media_id, @@ -330,6 +352,9 @@ async def generate_omni_flash_video( the workflow names and primary media IDs required by the Omni polling path. Do not feed Omni operation handles to ``check_video_status``. """ + blocked = _batch_path_blocks_omni() + if blocked: + return blocked refs = _validate_reference_inputs(reference_media_ids, duration_s, aspect_ratio) model_key = _load_model_key(duration_s, mode="reference_to_video") client = get_flow_client() @@ -384,6 +409,9 @@ async def check_omni_flash_status( ``flow.projectInitialData`` tRPC response. The old ``/v1/media`` transport currently returns ``INVALID_ARGUMENT`` for these workflow media IDs. """ + blocked = _batch_path_blocks_omni() + if blocked: + return blocked normalized = [] for workflow in workflows or []: item = _normalize_workflow(workflow) diff --git a/agent/services/video_reviewer.py b/agent/services/video_reviewer.py index 5a1e1b872..013d4e15b 100644 --- a/agent/services/video_reviewer.py +++ b/agent/services/video_reviewer.py @@ -135,7 +135,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() @@ -144,18 +149,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: 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)) 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/extension/background.js b/extension/background.js index b105f086d..98228da85 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 @@ -135,9 +152,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'); @@ -146,11 +161,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; @@ -216,7 +229,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); @@ -319,39 +334,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' }; } } @@ -372,6 +431,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) { @@ -428,6 +617,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; @@ -599,14 +790,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 })); } @@ -743,6 +932,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 01105d548..1ae492a63 100644 --- a/extension/manifest.json +++ b/extension/manifest.json @@ -1,13 +1,22 @@ { "manifest_version": 3, "name": "Flow Kit", - "version": "0.2.1", - "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/*", - "https://flow-content.google/*", "http://127.0.0.1:8100/*" ], "background": { @@ -16,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": { diff --git a/skills/fk-change-model.md b/skills/fk-change-model.md index 0013fd6bc..5cba5f6ba 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 130b562cb..1eacc340c 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 656b1d381..c134a557b 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`, `UNUSUAL_ACTIVITY`, `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`, `UNUSUAL_ACTIVITY`, `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`) @@ -86,7 +106,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 | @@ -107,11 +127,44 @@ 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` | +| `UNSUPPORTED_ON_BATCH_API: Omni Flash` | Omni speaks the pre-migration REST + tRPC endpoints; no batchexecute payload captured | **Terminal — not retried** | Use `model_family=veo`, or `USE_BATCH_RPC=0` on a profile that still holds a bearer | +| `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` | + +`NO_AT_TOKEN`, `NO_FLOW_TAB` and `FLOW_TAB_DISCARDED` are profile-local, so +with several extension profiles connected the agent fails the request over to +another one before reporting it. A single such error in the log with the job +still succeeding is that failover working, not a fault. + +Three 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. +- **A retry does not resubmit.** A batch operation id is a bare uuid, which the + Low Priority workflow path treats as unrecoverable and resubmits. On the batch + path it is recoverable — the status poll finds it in the project listing — so + a retried video request re-polls the render already running instead of paying + for a second one. + ### D. YouTube upload errors (`youtube/upload.py`) | Error | Cause | Fix | @@ -138,16 +191,16 @@ 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` | | 403 `PUBLIC_ERROR_UNUSUAL_ACTIVITY` / `reCAPTCHA evaluation failed` | Google flagged the session as bot-like (rapid bursts, VPN/shared IP, stale cookies). **Pause submits**, then in Chrome: `chrome://settings/cookies` → remove cookies for `google.com` and `labs.google` → reload `labs.google/fx/tools/flow` → sign in & solve any captcha → resubmit with ≥1s gap and ≤5 concurrent. Switch network or wait 1–6 h if still blocked | | 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` | @@ -156,6 +209,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. @@ -170,6 +224,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 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 diff --git a/tests/unit/test_omni_flash.py b/tests/unit/test_omni_flash.py index 38d5ed9fc..2333e6ee4 100644 --- a/tests/unit/test_omni_flash.py +++ b/tests/unit/test_omni_flash.py @@ -1,9 +1,16 @@ -"""Unit tests for Gemini Omni Flash Flow submissions and workflow polling.""" +"""Unit tests for Gemini Omni Flash Flow submissions and workflow polling. + +Omni speaks the pre-migration transports — the REST endpoints on aisandbox-pa +and the labs.google tRPC snapshot it polls through — so the wire contracts +asserted here are legacy-path contracts and the module is pinned to that path +for the file. What happens on the batch path is one test at the bottom. +""" from unittest.mock import AsyncMock, MagicMock, patch import pytest +import agent.services.omni_flash as omni_flash from agent.services.omni_flash import ( OMNI_FLASH_MAX_REFERENCE_IMAGES, _fetch_media_url, @@ -17,6 +24,12 @@ ) +@pytest.fixture(autouse=True) +def legacy_transport(monkeypatch): + """Omni is only reachable on the pre-migration path; assert it there.""" + monkeypatch.setattr(omni_flash, "USE_BATCH_RPC", False) + + @pytest.mark.parametrize( ("duration", "expected"), [ @@ -462,3 +475,52 @@ async def test_submit_rejects_empty_reference_set(): project_id="p", duration_s=8, ) + + +class TestBatchPathIsRefusedRatherThanAttempted: + """Flow stopped minting the bearer these endpoints need, and no Omni + payload has been captured off the new frontend. Saying so beats a 401 + five retries deep.""" + + @pytest.fixture(autouse=True) + def batch_transport(self, monkeypatch): + monkeypatch.setattr(omni_flash, "USE_BATCH_RPC", True) + + @pytest.fixture + def client(self): + with patch("agent.services.omni_flash.get_flow_client") as factory: + stub = MagicMock() + stub._send = AsyncMock() + factory.return_value = stub + yield stub + + async def test_first_frame_names_the_gap_and_sends_nothing(self, client): + result = await generate_omni_flash_first_frame_video( + start_image_media_id="mid", prompt="go", project_id="pid") + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + client._send.assert_not_called() + + async def test_first_last_names_the_gap_and_sends_nothing(self, client): + result = await generate_omni_flash_first_last_video( + start_image_media_id="a", end_image_media_id="b", + prompt="go", project_id="pid") + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + client._send.assert_not_called() + + async def test_reference_to_video_names_the_gap_and_sends_nothing(self, client): + result = await generate_omni_flash_video( + reference_media_ids=["a"], prompt="go", project_id="pid") + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + client._send.assert_not_called() + + async def test_polling_names_the_gap_and_sends_nothing(self, client): + result = await check_omni_flash_status( + [{"name": "wf", "primary_media_id": "mid", "project_id": "pid"}]) + assert "UNSUPPORTED_ON_BATCH_API" in result["error"] + client._send.assert_not_called() + + async def test_the_message_points_at_both_ways_out(self, client): + result = await generate_omni_flash_video( + reference_media_ids=["a"], prompt="go", project_id="pid") + assert "docs/CAPTURE.md" in result["error"] + assert "USE_BATCH_RPC=0" in result["error"]