`
+- 响应:
+```json
+{
+ "ok": true,
+ "commands": [
+ {
+ "id": "req-uuid",
+ "method": "api_request",
+ "params": { "...": "..." }
+ }
+ ],
+ "server_time": 1784488000000
+}
+```
+
+`POST /api/ext/callback`(已有,增强)
+- 继续接收命令结果、`token_captured`、`extension_ready`、`ping`、`media_urls_refresh`
+- 必须支持 Bearer secret
+
+#### B. 后端 → 扩展
+仍使用现有 method,只是 `send_message()` 在 HTTP 模式下写入命令队列。
+
+#### C. 连接健康定义
+```python
+extension_connected = session_last_seen_within(15s)
+has_flow_key = bool(flow_key)
+healthy = extension_connected and has_flow_key
+```
+
+---
+
+## 文件结构
+
+### 将创建
+- `flow-agent/omniflash/http_bridge.py`
+- `flow-agent/tests/test_http_bridge.py`
+- `flow-agent/tests/test_ext_http_api.py`
+- `docs/superpowers/plans/2026-07-20-extension-http-bridge.md`(本文件)
+
+### 将修改
+- `flow-chrome-extension/manifest.json`
+- `flow-chrome-extension/background.js`
+- `flow-agent/omniflash/bridge.py`
+- `flow-agent/cli/api.py`
+- `flow-agent/omniflash/config.py`(如需)
+- `flow-agent/config.env`
+- `README.md`
+
+### 不改
+- Google Flow 生成业务逻辑
+- OpenAI 兼容外部 API 形状
+- MCP `/sse`(那是 MCP 客户端通道)
+
+---
+
+### 任务 0:工作区确认与基线
+
+**文件:**
+- 工作目录:`F:\Code\Flow-Agent-New`
+
+- [x] **步骤 1:确认目录**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+dir
+Test-Path .\flow-agent\omniflash\bridge.py
+Test-Path .\flow-chrome-extension\background.js
+```
+
+预期:两者均为 `True`
+
+- [x] **步骤 2:安装/确认 Python 依赖(在本仓库)**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+# 优先使用 uv 或现有 venv;示例如下
+python -m pip install -e . -q
+python -c "import omniflash, cli.api; print('ok')"
+```
+
+- [x] **步骤 3:基线状态**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+git status
+git rev-parse --show-toplevel
+```
+
+预期:toplevel 为 `F:/Code/Flow-Agent-New` 或该仓库根
+
+---
+
+### 任务 1:后端 HTTP session + 命令队列模型
+
+**文件:**
+- 创建:`flow-agent/omniflash/http_bridge.py`
+- 修改:`flow-agent/omniflash/bridge.py`(仅预留接入点,可放任务 2)
+- 测试:`flow-agent/tests/test_http_bridge.py`
+
+- [x] **步骤 1:编写失败测试**
+
+```python
+# flow-agent/tests/test_http_bridge.py
+import time
+from omniflash.http_bridge import ExtensionHttpRegistry
+
+def test_hello_registers_session_and_accepts_token():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ out = reg.hello(session_id="s1", flow_key="tok", secret="test-secret")
+ assert out["ok"] is True
+ assert reg.is_connected("s1") is True
+ assert reg.get_flow_key("s1") == "tok"
+
+def test_enqueue_and_poll_returns_command_once():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ reg.hello(session_id="s1", flow_key="tok", secret="sec")
+ reg.enqueue("s1", {"id": "r1", "method": "get_status", "params": {}})
+ first = reg.poll("s1")
+ second = reg.poll("s1")
+ assert len(first["commands"]) == 1
+ assert first["commands"][0]["id"] == "r1"
+ assert second["commands"] == []
+
+def test_session_expires():
+ reg = ExtensionHttpRegistry(session_ttl_sec=1)
+ reg.hello(session_id="s1", flow_key="tok", secret="sec")
+ time.sleep(1.2)
+ assert reg.is_connected("s1") is False
+```
+
+- [x] **步骤 2:运行测试确认失败**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+python -m pytest tests/test_http_bridge.py -v
+```
+
+预期:FAIL,找不到 `ExtensionHttpRegistry`
+
+- [x] **步骤 3:实现最小 registry**
+
+`flow-agent/omniflash/http_bridge.py` 需包含:
+- `hello(session_id, flow_key, secret, meta=None)`
+- `touch(session_id)`
+- `is_connected(session_id=None)`
+- `has_online_session()`
+- `enqueue(session_id|None, command)`
+- `poll(session_id, max_commands=10)`
+- `get_flow_key(session_id=None)`
+- 线程安全(`threading.Lock`)
+
+```python
+@dataclass
+class Session:
+ session_id: str
+ secret: str
+ flow_key: str | None
+ last_seen: float
+ queue: deque
+```
+
+- [x] **步骤 4:运行测试确认通过**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+python -m pytest tests/test_http_bridge.py -v
+```
+
+- [x] **步骤 5:Commit**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+git add flow-agent/omniflash/http_bridge.py flow-agent/tests/test_http_bridge.py
+git commit -m "feat(bridge): add HTTP extension session registry and command queue"
+```
+
+---
+
+### 任务 2:FastAPI 增加 hello/poll,并增强 callback/health
+
+**文件:**
+- 修改:`flow-agent/cli/api.py`
+- 修改:`flow-agent/omniflash/bridge.py`
+- 测试:`flow-agent/tests/test_ext_http_api.py`
+
+- [x] **步骤 1:写 API 级失败测试**
+
+```python
+# flow-agent/tests/test_ext_http_api.py
+from fastapi.testclient import TestClient
+
+def test_hello_and_poll_roundtrip():
+ from cli.api import app, bridge
+ client = TestClient(app)
+ r = client.post("/api/ext/hello", json={
+ "session_id": "s-test",
+ "flowKeyPresent": True,
+ "flowKey": "tok-1",
+ "extension_version": "1.0.0",
+ })
+ assert r.status_code == 200
+ body = r.json()
+ assert body["ok"] is True
+ secret = body["secret"]
+ bridge.enqueue_http_command({"id": "c1", "method": "get_status", "params": {}})
+ p = client.get(
+ "/api/ext/poll",
+ params={"session_id": "s-test"},
+ headers={"Authorization": f"Bearer {secret}"},
+ )
+ assert p.status_code == 200
+ assert p.json()["commands"][0]["id"] == "c1"
+```
+
+- [x] **步骤 2:运行测试确认失败**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+python -m pytest tests/test_ext_http_api.py -v
+```
+
+- [x] **步骤 3:实现端点与 bridge 接入**
+
+1. `POST /api/ext/hello`
+2. `GET /api/ext/poll`
+3. 增强 `POST /api/ext/callback`
+4. `/health` 增加:
+```python
+"extension_connected": bridge.is_extension_connected(),
+"has_flow_key": bridge._flow_key is not None,
+"transport": bridge.active_transport(), # http | ws | none
+```
+
+`bridge.send_message`:
+```python
+async def send_message(self, msg):
+ if self.http_registry and self.http_registry.has_online_session():
+ self.http_registry.enqueue(None, msg)
+ return
+ # fallback existing WS send
+```
+
+`api_request` 前置:
+```python
+if not self.is_extension_connected():
+ return {"error": "Extension not connected"}
+```
+
+- [x] **步骤 4:测试通过**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+python -m pytest tests/test_http_bridge.py tests/test_ext_http_api.py -v
+```
+
+- [x] **步骤 5:Commit**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+git add flow-agent/cli/api.py flow-agent/omniflash/bridge.py flow-agent/tests/test_ext_http_api.py
+git commit -m "feat(api): add extension hello/poll HTTP endpoints and health transport"
+```
+
+---
+
+### 任务 3:扩展 manifest + HTTP transport
+
+**文件:**
+- 修改:`flow-chrome-extension/manifest.json`
+- 修改:`flow-chrome-extension/background.js`
+
+- [x] **步骤 1:更新 host_permissions**
+
+```json
+"host_permissions": [
+ "https://labs.google/*",
+ "https://aisandbox-pa.googleapis.com/*",
+ "https://aisandbox-pa.sandbox.googleapis.com/*",
+ "https://storage.googleapis.com/*",
+ "http://127.0.0.1:8001/*",
+ "http://localhost:8001/*",
+ "http://127.0.0.1:8100/*"
+]
+```
+
+- [x] **步骤 2:HTTP 优先常量**
+
+```javascript
+const AGENT_BASE = 'http://127.0.0.1:8001';
+const AGENT_WS_URL = 'ws://127.0.0.1:8001/ws';
+const AGENT_HELLO_URL = `${AGENT_BASE}/api/ext/hello`;
+const AGENT_POLL_URL = `${AGENT_BASE}/api/ext/poll`;
+const AGENT_CALLBACK_URL = `${AGENT_BASE}/api/ext/callback`;
+const TRANSPORT_MODE = 'auto'; // auto | http | ws
+```
+
+- [x] **步骤 3:实现 `connectViaHttp()` / poll loop / `handleAgentMessage()`**
+
+- [x] **步骤 4:`sendToAgent` 优先 HTTP callback + Bearer**
+
+- [x] **步骤 5:`connectToAgent()` auto 模式:HTTP 成功则不强制 WS**
+
+- [x] **步骤 6:popup `connected` 判定包含 http session**
+
+- [x] **步骤 7:sniff 转发改用 `callbackUrl || AGENT_CALLBACK_URL`**
+
+- [x] **步骤 8:手动验证**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+python -m flow_cli serve
+# 另开终端
+curl http://127.0.0.1:8001/health
+```
+
+官方 Chrome 加载:`F:\Code\Flow-Agent-New\flow-chrome-extension`
+
+期望 health:
+```json
+{"extension_connected": true, "has_flow_key": true, "transport": "http"}
+```
+
+- [x] **步骤 9:Commit**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+git add flow-chrome-extension/manifest.json flow-chrome-extension/background.js
+git commit -m "feat(extension): prefer HTTP hello/poll bridge over WebSocket"
+```
+
+---
+
+### 任务 4:配置与文档
+
+**文件:**
+- 修改:`flow-agent/omniflash/config.py`
+- 修改:`flow-agent/config.env`
+- 修改:`README.md`
+
+- [x] **步骤 1:配置项**
+
+```env
+EXT_TRANSPORT=auto
+EXT_SESSION_TTL_SEC=20
+EXT_POLL_INTERVAL_MS=1000
+ENABLE_EXTENSION_WS=1
+```
+
+- [x] **步骤 2:README 增加“指纹浏览器 / HTTP 桥”说明**
+
+- [x] **步骤 3:Commit**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+git add flow-agent/omniflash/config.py flow-agent/config.env README.md
+git commit -m "docs: document HTTP extension transport for fingerprint browsers"
+```
+
+---
+
+### 任务 5:端到端验证矩阵
+
+- [x] **步骤 1:自动化**
+
+```powershell
+cd F:\Code\Flow-Agent-New\flow-agent
+python -m pytest tests/test_http_bridge.py tests/test_ext_http_api.py -v
+```
+
+- [ ] **步骤 2:官方 Chrome 回归**
+
+- [ ] **步骤 3:Hubstudio 验证(需本地 HTTP 可达)**
+
+- [ ] **步骤 4:AdsPower 验证**
+
+- [ ] **步骤 5:必要时修复后最终 commit**
+
+```powershell
+cd F:\Code\Flow-Agent-New
+git add -A
+git commit -m "fix: stabilize HTTP extension bridge for fingerprint browsers"
+```
+
+---
+
+## 风险与缓解
+
+| 风险 | 影响 | 缓解 |
+|------|------|------|
+| MV3 SW 休眠 | 命令延迟 | alarms 保活 + 后端队列 |
+| 指纹浏览器连 HTTP 也拦 | 仍不可用 | 记录环境限制;后续 LAN IP / 隧道附加计划 |
+| HTTP+WS 双投递 | 重复生成 | 命令 id 幂等 + 单一 active transport |
+| 旧扩展只支持 WS | 兼容 | 保留 `/ws` |
+
+## 不在本计划内
+1. 公网 WSS / Cloudflare Tunnel
+2. 默认绑定 `0.0.0.0`(可另开附加任务)
+3. 账号池 / OpenAI API 重写
+4. 自动改 Hubstudio/AdsPower 代理
+
+## 成功标准
+1. 官方 Chrome:HTTP 模式连通并可完成 credits/生成
+2. 至少一个指纹浏览器:在本地 HTTP 可达时不依赖 WS 也能连通
+3. `/health` 显示 `transport`
+4. `EXT_TRANSPORT=ws` 仍可用
+
+## 建议工期
+半天到 1 天:任务 1–2(后端)→ 任务 3(扩展)→ 任务 4–5(文档与验证)
+
+---
+
+## 自检
+1. 规格覆盖:HTTP 轮询替代 WS 下行 + callback 回传 + 健康状态重定义
+2. 工作目录已锁定到 `F:\Code\Flow-Agent-New`
+3. 关键路径均相对本仓库,不再引用旧 Cloak/Hubstudio 路径
diff --git a/flow-agent/cli/api.py b/flow-agent/cli/api.py
index 85745f5..bbc1c6a 100644
--- a/flow-agent/cli/api.py
+++ b/flow-agent/cli/api.py
@@ -30,6 +30,7 @@
from omniflash.config import CREDITS_PER_VIDEO
from omniflash.generators.t2i import generate_image, download_image, _parse_image_results
+
# Setup logging (format configured centrally in omniflash/__init__.py, imported above)
log = logging.getLogger("omniflash.openai_api")
@@ -92,25 +93,91 @@ async def recover_orphan_response(data: dict, meta: dict):
except Exception:
log.exception("Failed to recover orphan response")
-# ExtensionBridge lifecycle
+# ExtensionBridge / durable control-center lifecycle
bridge: Optional[ExtensionBridge] = None
+
+def _control_center_enabled() -> bool:
+ return os.environ.get("FLOW_CONTROL_CENTER_ENABLED", "0").strip().lower() in {
+ "1", "true", "yes", "on"
+ }
+
+
@asynccontextmanager
async def lifespan(app: FastAPI):
global bridge
- log.info("Starting Flow Agent Extension Bridge (OpenAI Interface)...")
- bridge = ExtensionBridge()
- bridge.set_orphan_handler(recover_orphan_response)
- await bridge.start()
-
- # Run extension connection in background
- asyncio.create_task(bridge.wait_for_extension(timeout=30))
-
- yield
+ control_services: ControlServices | None = None
+ bridge = None
+ if _control_center_enabled():
+ # Optional durable control-center path. Imported lazily so the default
+ # HTTP/WS extension bridge works without that package present.
+ try:
+ from omniflash.control.db import create_database
+ from omniflash.control.events import EventBus
+ from omniflash.control.extension_registry import ExtensionRegistry
+ from omniflash.control.flow_adapter import FlowAdapter, FlowAdapterSettings
+ from omniflash.control.material_store import MaterialStore
+ from omniflash.control.profile_manager import ProfileManager
+ from omniflash.control.recovery import RecoveryService
+ from omniflash.control.repositories import Repositories
+ from omniflash.control.scheduler import Scheduler
+ from omniflash.control.services import ControlServices
+ from omniflash.control.settings import ControlSettings
+ except ImportError as error:
+ raise RuntimeError(
+ "FLOW_CONTROL_CENTER_ENABLED is set but omniflash.control is unavailable"
+ ) from error
+
+ settings = ControlSettings()
+ settings.ensure_directories()
+ database = await create_database(settings.db_path)
+ repositories = Repositories(database.session_factory)
+ events = EventBus()
+ profiles = ProfileManager(settings.profiles_dir)
+ extensions = ExtensionRegistry(
+ database.session_factory, request_timeout=settings.submit_timeout_seconds
+ )
+ materials = MaterialStore(
+ settings.materials_dir, database.session_factory,
+ max_material_bytes=settings.max_material_bytes,
+ )
+ flow = FlowAdapter(
+ extensions, materials, repositories,
+ settings=FlowAdapterSettings(
+ project_id=DEFAULT_PROJECT,
+ submit_timeout_seconds=settings.submit_timeout_seconds,
+ outputs_dir=settings.outputs_dir,
+ ),
+ )
+ scheduler = Scheduler(
+ repositories, flow=flow, settings=settings, events=events
+ )
+ recovery = RecoveryService(repositories, events=events)
+ control_services = ControlServices(
+ settings=settings, database=database, repositories=repositories,
+ events=events, profiles=profiles, extensions=extensions,
+ materials=materials, flow=flow, scheduler=scheduler,
+ recovery=recovery,
+ )
+ await control_services.start()
+ app.state.services = control_services
+ log.info("Durable control-center services started")
+ else:
+ log.info("Starting Flow Agent Extension Bridge (OpenAI Interface)...")
+ bridge = ExtensionBridge()
+ bridge.set_orphan_handler(recover_orphan_response)
+ await bridge.start()
+ asyncio.create_task(bridge.wait_for_extension(timeout=30))
- log.info("Closing Flow Agent Extension Bridge...")
- if bridge:
- await bridge.close()
+ try:
+ yield
+ finally:
+ if control_services is not None:
+ await control_services.close()
+ app.state.services = None
+ elif bridge:
+ log.info("Closing Flow Agent Extension Bridge...")
+ await bridge.close()
app = FastAPI(
title="Flow Agent OpenAI API Wrapper",
@@ -209,6 +276,28 @@ class VideoGenerationRequest(BaseModel):
# Extension WebSocket and Callback Endpoints
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
+ enable_ws = os.environ.get("ENABLE_EXTENSION_WS", "1").strip().lower() in {
+ "1", "true", "yes", "on"
+ }
+ if not enable_ws:
+ await websocket.close(code=1000)
+ return
+
+ control_center_enabled = os.environ.get(
+ "FLOW_CONTROL_CENTER_ENABLED", "0"
+ ).strip().lower() in {"1", "true", "yes", "on"}
+ if control_center_enabled:
+ services = getattr(app.state, "services", None)
+ extensions = getattr(services, "extensions", None)
+ if extensions is not None:
+ await extensions.handle_websocket(websocket)
+ return
+
+ log.error("Control center enabled but extension registry is unavailable")
+ await websocket.accept()
+ await websocket.close(code=1013)
+ return
+
await websocket.accept()
log.info("Extension connecting via FastAPI WebSocket...")
global bridge
@@ -218,13 +307,102 @@ async def websocket_endpoint(websocket: WebSocket):
log.error("Bridge not initialized")
await websocket.close()
+
+def _public_base_url() -> str:
+ space_id = os.environ.get("SPACE_ID")
+ if space_id:
+ author, name = space_id.split("/")
+ subdomain = f"{author.lower()}-{name.lower()}".replace("_", "-")
+ return f"https://{subdomain}.hf.space"
+ port = os.environ.get("OPENAI_API_PORT", "8001")
+ return f"http://127.0.0.1:{port}"
+
+
+@app.post("/api/ext/hello")
+async def extension_hello(body: dict):
+ """Register/refresh an HTTP extension session and issue callback secret."""
+ global bridge
+ if bridge is None:
+ return JSONResponse(status_code=503, content={"ok": False})
+
+ session_id = body.get("session_id") or body.get("sessionId")
+ if not session_id:
+ return JSONResponse(
+ status_code=400,
+ content={"ok": False, "error": "session_id required"},
+ )
+
+ flow_key = body.get("flowKey") or body.get("flow_key")
+ if not flow_key and body.get("flowKeyPresent") and bridge._flow_key:
+ flow_key = bridge._flow_key
+
+ # Prefer existing per-process callback secret so WS and HTTP share auth.
+ secret = bridge._callback_secret
+ meta = {
+ "extension_version": body.get("extension_version") or body.get("extensionVersion"),
+ "capabilities": body.get("capabilities"),
+ }
+ registered = bridge.register_http_session(
+ session_id=str(session_id),
+ flow_key=flow_key,
+ secret=secret,
+ meta=meta,
+ )
+ base = _public_base_url()
+ poll_interval_ms = int(os.environ.get("EXT_POLL_INTERVAL_MS", "1000"))
+ return {
+ "ok": True,
+ "session_id": registered["session_id"],
+ "secret": secret,
+ "callback_url": f"{base}/api/ext/callback",
+ "poll_url": f"{base}/api/ext/poll",
+ "poll_interval_ms": poll_interval_ms,
+ "events_url": f"{base}/api/ext/events",
+ }
+
+
+@app.get("/api/ext/poll")
+async def extension_poll(request: Request, session_id: str = Query(...)):
+ """Pull pending backend commands for an authenticated HTTP session."""
+ global bridge
+ if bridge is None:
+ return JSONResponse(status_code=503, content={"ok": False, "commands": []})
+
+ authorization = request.headers.get("Authorization")
+ if authorization is None:
+ return JSONResponse(status_code=401, content={"ok": False, "commands": []})
+ if not bridge.verify_http_session_authorization(session_id, authorization):
+ return JSONResponse(status_code=403, content={"ok": False, "commands": []})
+
+ result = bridge.poll_http_commands(session_id)
+ if not result.get("ok"):
+ return JSONResponse(status_code=404, content=result)
+ return result
+
+
@app.post("/api/ext/callback")
-async def http_callback(body: dict):
+async def http_callback(request: Request, body: dict):
global bridge
- if bridge:
- success = bridge.handle_http_callback(body)
- return {"ok": success}
- return {"ok": False}
+ if bridge is None:
+ return JSONResponse(status_code=503, content={"ok": False})
+
+ authorization = request.headers.get("Authorization")
+ if authorization is None:
+ return JSONResponse(status_code=401, content={"ok": False})
+ if not bridge.verify_callback_authorization(authorization):
+ # Also accept a valid per-session HTTP secret (shared secret is default).
+ session_id = body.get("session_id") or body.get("sessionId")
+ if not (
+ session_id
+ and bridge.verify_http_session_authorization(str(session_id), authorization)
+ ):
+ return JSONResponse(status_code=403, content={"ok": False})
+
+ accepted = bridge.handle_http_callback(body)
+ return JSONResponse(
+ status_code=200 if accepted else 404,
+ content={"ok": bool(accepted)},
+ )
# OpenAI Endpoints
@@ -877,21 +1055,46 @@ async def get_flow_credits():
raise HTTPException(status_code=500, detail=str(e))
+
+def _runtime_identity() -> dict:
+ """Expose which on-disk code this process is running (for test safety)."""
+ import omniflash.bridge as _bridge_mod
+
+ return {
+ "code_root": os.path.dirname(os.path.dirname(os.path.abspath(__file__))),
+ "cli_api_file": os.path.abspath(__file__),
+ "bridge_file": os.path.abspath(_bridge_mod.__file__),
+ }
+
+
@app.get("/")
async def root():
- return {"status": "running", "service": "Flow Agent API"}
+ identity = _runtime_identity()
+ return {
+ "status": "running",
+ "service": "Flow Agent API",
+ **identity,
+ }
# Health Check
@app.get("/health")
async def health():
global bridge
+ identity = _runtime_identity()
if not bridge:
- return {"status": "starting", "connected": False}
+ return {
+ "status": "starting",
+ "connected": False,
+ "transport": "none",
+ **identity,
+ }
return {
"status": "healthy" if await bridge.health_check() else "unauthorized_or_disconnected",
- "extension_connected": bridge._ws is not None,
- "has_flow_key": bridge._flow_key is not None
+ "extension_connected": bridge.is_extension_connected(),
+ "has_flow_key": bridge.has_flow_key(),
+ "transport": bridge.active_transport(),
+ **identity,
}
@@ -1030,7 +1233,7 @@ def get_mcp_tools_list():
},
{
"name": "generate_flow_video",
- "description": "Generate a 10-second cinematic video clip using Google Flow. Optionally supports a starting frame reference image (Image-to-Video).",
+ "description": "Generate a cinematic video clip using Google Flow (4/6/8/10 seconds). Optionally supports a starting frame reference image (Image-to-Video).",
"inputSchema": {
"type": "object",
"properties": {
@@ -1044,6 +1247,12 @@ def get_mcp_tools_list():
"enum": ["landscape", "portrait"],
"default": "landscape"
},
+ "duration": {
+ "type": "integer",
+ "description": "Video duration in seconds (default: 10)",
+ "enum": [4, 6, 8, 10],
+ "default": 10
+ },
"start_image_path": {
"type": "string",
"description": "Optional local file path to a starting reference image on the host for Image-to-Video"
@@ -1157,6 +1366,20 @@ async def execute_mcp_tool(request_id, tool_name, arguments):
aspect = arguments.get("aspect", "landscape")
start_image_path = arguments.get("start_image_path")
start_media_id = arguments.get("start_media_id")
+ allowed_durations = {4, 6, 8, 10}
+ try:
+ duration = int(arguments.get("duration", 10))
+ except (TypeError, ValueError):
+ duration = 10
+ if duration not in allowed_durations:
+ return {
+ "jsonrpc": "2.0",
+ "id": request_id,
+ "error": {
+ "code": -32602,
+ "message": f"Error: duration must be one of {sorted(allowed_durations)}, got {duration}",
+ },
+ }
start_image_base64 = None
if start_image_path:
@@ -1186,7 +1409,7 @@ async def execute_mcp_tool(request_id, tool_name, arguments):
prompt=prompt,
aspect=aspect,
n=1,
- duration=10,
+ duration=duration,
image_base64=start_image_base64,
start_media_id=start_media_id
)
diff --git a/flow-agent/config.env b/flow-agent/config.env
index 05304eb..4d38b4a 100644
--- a/flow-agent/config.env
+++ b/flow-agent/config.env
@@ -40,3 +40,14 @@ API_REQUEST_TIMEOUT=180
# ─── Google Flow backend (do not change unless Google moves it) ───
API_BASE=https://aisandbox-pa.googleapis.com
+
+# ─── Extension bridge (fingerprint-browser friendly) ───
+# auto = HTTP hello/poll first, WebSocket fallback; http | ws force one mode
+EXT_TRANSPORT=auto
+# HTTP session considered online if hello/poll seen within this many seconds
+EXT_SESSION_TTL_SEC=20
+# Suggested poll interval returned to the extension (ms)
+EXT_POLL_INTERVAL_MS=1000
+# Keep legacy /ws endpoint available
+ENABLE_EXTENSION_WS=1
+
diff --git a/flow-agent/flow_cli/__main__.py b/flow-agent/flow_cli/__main__.py
new file mode 100644
index 0000000..0adb9bc
--- /dev/null
+++ b/flow-agent/flow_cli/__main__.py
@@ -0,0 +1,5 @@
+"""Allow `python -m flow_cli serve` from the workspace flow-agent/ directory."""
+from .main import main
+
+if __name__ == "__main__":
+ main()
diff --git a/flow-agent/flow_cli/main.py b/flow-agent/flow_cli/main.py
index dc6d23d..57fc7c3 100644
--- a/flow-agent/flow_cli/main.py
+++ b/flow-agent/flow_cli/main.py
@@ -56,10 +56,25 @@ def cmd_serve(argv):
args = parser.parse_args(argv)
import uvicorn
+
+ # Always serve from this repo's flow-agent/ tree (not a globally installed flow.exe).
+ os.chdir(ROOT_DIR)
+ if ROOT_DIR not in sys.path:
+ sys.path.insert(0, ROOT_DIR)
+
print(f"Flow Agent starting on http://{args.host}:{args.port}")
+ print(f" code_root: {ROOT_DIR}")
+ print(f" python: {sys.executable}")
print("Waiting for the Chrome extension (open Google Flow in Chrome to connect).")
- uvicorn.run("cli.api:app", host=args.host, port=args.port,
- reload=args.reload, access_log=False)
+ print("Tip: prefer .\\scripts\\dev-serve.ps1 so old global flow.exe cannot steal :8001")
+ uvicorn.run(
+ "cli.api:app",
+ host=args.host,
+ port=args.port,
+ reload=args.reload,
+ access_log=False,
+ app_dir=ROOT_DIR,
+ )
def cmd_credits(argv):
@@ -86,6 +101,8 @@ def cmd_status(argv):
print(f" status: {data.get('status')}")
print(f" extension_connected: {data.get('extension_connected')}")
print(f" has_flow_key: {data.get('has_flow_key')}")
+ if "transport" in data:
+ print(f" transport: {data.get('transport')}")
except urllib.error.URLError:
print(f"Backend: down. Start it with `flow serve`.")
sys.exit(1)
diff --git a/flow-agent/flow_mcp_server.py b/flow-agent/flow_mcp_server.py
index 780af92..f1fc553 100755
--- a/flow-agent/flow_mcp_server.py
+++ b/flow-agent/flow_mcp_server.py
@@ -125,7 +125,7 @@ def handle_tools_list(request_id):
},
{
"name": "generate_flow_video",
- "description": "Generate a 10-second cinematic video clip using Google Flow. Optionally supports a starting frame reference image (Image-to-Video).",
+ "description": "Generate a cinematic video clip using Google Flow (4/6/8/10 seconds). Optionally supports a starting frame reference image (Image-to-Video).",
"inputSchema": {
"type": "object",
"properties": {
@@ -139,6 +139,12 @@ def handle_tools_list(request_id):
"enum": ["landscape", "portrait"],
"default": "landscape"
},
+ "duration": {
+ "type": "integer",
+ "description": "Video duration in seconds (default: 10)",
+ "enum": [4, 6, 8, 10],
+ "default": 10
+ },
"start_image_path": {
"type": "string",
"description": "Optional local file path to a starting reference image on the host for Image-to-Video"
@@ -251,15 +257,22 @@ def call_generate_flow_image(prompt, size="1280x720", ref_image_path=None, model
except Exception as e:
return f"Failed to communicate with Flow Agent server: {str(e)}", None
-def call_generate_flow_video(prompt, aspect="landscape", start_image_path=None):
+def call_generate_flow_video(prompt, aspect="landscape", start_image_path=None, duration=10):
if not prompt or not str(prompt).strip():
return "Error: 'prompt' is required and cannot be empty."
prompt = str(prompt).strip()
+ allowed_durations = {4, 6, 8, 10}
+ try:
+ duration = int(duration)
+ except (TypeError, ValueError):
+ duration = 10
+ if duration not in allowed_durations:
+ return f"Error: duration must be one of {sorted(allowed_durations)}, got {duration}"
payload = {
"prompt": prompt,
"aspect": aspect,
"n": 1,
- "duration": 10
+ "duration": duration,
}
if start_image_path:
@@ -321,11 +334,12 @@ def handle_tool_call(request_id, tool_name, arguments):
"data": image_data_b64,
"mimeType": "image/png"
})
- elif tool_name == "generate_flow_video":
- prompt = arguments.get("prompt")
- aspect = arguments.get("aspect", "landscape")
- start_image_path = arguments.get("start_image_path")
- text = call_generate_flow_video(prompt, aspect, start_image_path)
+ elif tool_name == "generate_flow_video":
+ prompt = arguments.get("prompt")
+ aspect = arguments.get("aspect", "landscape")
+ start_image_path = arguments.get("start_image_path")
+ duration = arguments.get("duration", 10)
+ text = call_generate_flow_video(prompt, aspect, start_image_path, duration)
content = [{"type": "text", "text": text}]
else:
return {
diff --git a/flow-agent/omniflash/bridge.py b/flow-agent/omniflash/bridge.py
index 3692e26..084691b 100644
--- a/flow-agent/omniflash/bridge.py
+++ b/flow-agent/omniflash/bridge.py
@@ -5,9 +5,11 @@
"""
import asyncio
+import hmac
import json
import logging
import random
+import secrets
import threading
import uuid
from http.server import HTTPServer, BaseHTTPRequestHandler
@@ -18,7 +20,9 @@
WS_PORT, HTTP_PORT, API_BASE, API_KEY,
CLIENT_CTX, USER_AGENTS, API_REQUEST_TIMEOUT,
MAX_CONCURRENT_REQUESTS, REQUEST_MIN_INTERVAL,
+ EXT_SESSION_TTL_SEC,
)
+from .http_bridge import ExtensionHttpRegistry
log = logging.getLogger("omniflash.bridge")
@@ -26,7 +30,7 @@
class ExtensionBridge:
"""WebSocket server that Chrome extension connects to."""
- def __init__(self):
+ def __init__(self, session_ttl_sec: float | None = None):
self._ws = None
self._pending: dict[str, asyncio.Future] = {}
self._flow_key = None
@@ -47,6 +51,12 @@ def __init__(self):
self._rate_sem: asyncio.Semaphore | None = None
self._rate_lock: asyncio.Lock | None = None
self._last_request_at: float = 0.0
+ # Per-process callback credential shared only with the connected extension.
+ # Keep this out of configuration, logs, and error responses.
+ self._callback_secret = secrets.token_urlsafe(32)
+ # HTTP hello/poll transport (preferred for fingerprint browsers).
+ ttl = EXT_SESSION_TTL_SEC if session_ttl_sec is None else float(session_ttl_sec)
+ self.http_registry = ExtensionHttpRegistry(session_ttl_sec=ttl)
def _get_rate_limit(self):
"""Lazily build the concurrency semaphore + spacing lock on the active
@@ -75,7 +85,66 @@ def _remember_request(self, req_id, meta, max_keep=64):
oldest = next(iter(self._req_meta))
self._req_meta.pop(oldest, None)
+ def is_extension_connected(self) -> bool:
+ """True when an HTTP session is online or a WebSocket is open."""
+ if self.http_registry.has_online_session():
+ return True
+ return self._ws is not None
+
+ def has_flow_key(self) -> bool:
+ if self._flow_key:
+ return True
+ return bool(self.http_registry.get_flow_key())
+
+ def active_transport(self) -> str:
+ """Return preferred active transport: http | ws | none."""
+ if self.http_registry.has_online_session():
+ return "http"
+ if self._ws is not None:
+ return "ws"
+ return "none"
+
+ def enqueue_http_command(self, command: dict) -> bool:
+ """Enqueue a command for the latest online HTTP session."""
+ return self.http_registry.enqueue(None, command)
+
+ def register_http_session(
+ self,
+ session_id: str,
+ flow_key: str | None = None,
+ *,
+ secret: str | None = None,
+ meta: dict | None = None,
+ ) -> dict:
+ """Register/refresh an HTTP extension session and mirror flow key."""
+ resolved_secret = secret or self._callback_secret
+ out = self.http_registry.hello(
+ session_id=session_id,
+ flow_key=flow_key,
+ secret=resolved_secret,
+ meta=meta,
+ )
+ if flow_key:
+ self._flow_key = flow_key
+ if self._loop is not None:
+ self._loop.call_soon_threadsafe(self._connected.set)
+ else:
+ self._connected.set()
+ return out
+
+ def poll_http_commands(self, session_id: str, max_commands: int = 10) -> dict:
+ return self.http_registry.poll(session_id, max_commands=max_commands)
+
+ def verify_http_session_authorization(
+ self, session_id: str, authorization: str | None
+ ) -> bool:
+ return self.http_registry.verify_authorization(session_id, authorization)
+
async def send_message(self, msg):
+ # Prefer HTTP command queue when an extension is polling.
+ if self.http_registry.has_online_session():
+ self.http_registry.enqueue(None, msg)
+ return
if not self._ws:
return
try:
@@ -86,6 +155,13 @@ async def send_message(self, msg):
except Exception as e:
log.warning("Failed to send message: %s", e)
+ def _callback_config(self, callback_url):
+ return {
+ "type": "callback_config",
+ "secret": self._callback_secret,
+ "callback_url": callback_url,
+ }
+
async def handle_fastapi_ws(self, ws):
self._ws = ws
log.info("Extension connected via FastAPI WebSocket!")
@@ -101,11 +177,7 @@ async def handle_fastapi_ws(self, ws):
else:
callback_url = f"http://127.0.0.1:{os.environ.get('OPENAI_API_PORT', '8001')}/api/ext/callback"
- await self.send_message({
- "type": "callback_config",
- "secret": "flow_secret",
- "callback_url": callback_url
- })
+ await self.send_message(self._callback_config(callback_url))
# Send current state + resend token if we have one
await self.send_message({
@@ -130,15 +202,24 @@ async def handle_fastapi_ws(self, ws):
self._connected.clear()
async def start(self):
- """Start WS server and HTTP callback server."""
- self._loop = asyncio.get_event_loop()
- self._start_http_server()
+ """Start optional standalone WS/HTTP servers for non-FastAPI callers."""
+ self._loop = asyncio.get_running_loop()
+ self._ws_server = None
+ try:
+ self._start_http_server()
+ log.info("HTTP callback on http://127.0.0.1:%d", HTTP_PORT)
+ except OSError as error:
+ log.warning("Standalone HTTP callback server unavailable: %s", error)
+
+ try:
+ self._ws_server = await websockets.serve(
+ self._on_connect, "127.0.0.1", WS_PORT
+ )
+ log.info("WebSocket server on ws://127.0.0.1:%d", WS_PORT)
+ except OSError as error:
+ # FastAPI mode uses /ws on the API port; standalone WS is optional.
+ log.warning("Standalone WebSocket server unavailable: %s", error)
- self._ws_server = await websockets.serve(
- self._on_connect, "127.0.0.1", WS_PORT
- )
- log.info("WebSocket server on ws://127.0.0.1:%d", WS_PORT)
- log.info("HTTP callback on http://127.0.0.1:%d", HTTP_PORT)
log.info("Waiting for Chrome extension to connect...")
async def wait_for_extension(self, timeout=90, max_retries=3):
@@ -179,8 +260,8 @@ async def wait_for_extension(self, timeout=90, max_retries=3):
return False
async def _wait_for_ws(self):
- """Wait until a WebSocket connection is established."""
- while not self._ws:
+ """Wait until WebSocket or HTTP session is established."""
+ while not self.is_extension_connected():
await asyncio.sleep(0.5)
async def _wait_for_token(self, timeout):
@@ -194,7 +275,7 @@ async def _wait_for_token(self, timeout):
async def _request_flow_tab(self):
"""Ask extension to open or refresh a Flow tab."""
- if not self._ws:
+ if not self.is_extension_connected():
return
try:
log.info("Requesting extension to open/refresh Flow tab...")
@@ -208,8 +289,11 @@ async def _request_flow_tab(self):
async def health_check(self):
"""Quick check if extension is ready with valid token."""
- if not self._ws or not self._flow_key:
+ if not self.is_extension_connected() or not self.has_flow_key():
return False
+ # HTTP mode: presence of recent session + flow key is enough health.
+ if self.active_transport() == "http":
+ return True
try:
req_id = str(uuid.uuid4())
future = self._loop.create_future()
@@ -228,6 +312,9 @@ async def health_check(self):
async def _on_connect(self, ws):
self._ws = ws
log.info("Extension connected!")
+ await self.send_message(self._callback_config(
+ f"http://127.0.0.1:{HTTP_PORT}/api/ext/callback"
+ ))
try:
async for raw in ws:
data = json.loads(raw)
@@ -290,6 +377,15 @@ def _route_response(self, req_id, data):
else:
log.warning("Dropped orphan response for %s (no handler)", req_id)
+ def verify_callback_authorization(self, authorization):
+ """Constant-time verification for the callback Bearer credential."""
+ candidate = ""
+ if isinstance(authorization, str):
+ scheme, separator, value = authorization.partition(" ")
+ if separator and scheme.lower() == "bearer":
+ candidate = value.strip()
+ return hmac.compare_digest(candidate, self._callback_secret)
+
def handle_http_callback(self, data):
"""Called from HTTP thread when extension sends callback."""
req_id = data.get("id")
@@ -299,13 +395,32 @@ def handle_http_callback(self, data):
# loop thread; _route_response dedups and recovers as needed.
if (req_id in self._pending or req_id in self._req_meta
or req_id in self._seen_ids):
- self._loop.call_soon_threadsafe(
- self._resolve_pending, req_id, data
- )
+ if self._loop is not None:
+ self._loop.call_soon_threadsafe(
+ self._resolve_pending, req_id, data
+ )
+ else:
+ self._resolve_pending(req_id, data)
return True
- if data.get("type") == "token_captured":
+ msg_type = data.get("type")
+ if msg_type == "token_captured":
self._flow_key = data.get("flowKey")
- self._loop.call_soon_threadsafe(self._connected.set)
+ session_id = data.get("session_id") or data.get("sessionId")
+ if session_id and self._flow_key:
+ self.http_registry.hello(
+ session_id=str(session_id),
+ flow_key=self._flow_key,
+ secret=self._callback_secret,
+ )
+ if self._loop is not None:
+ self._loop.call_soon_threadsafe(self._connected.set)
+ else:
+ self._connected.set()
+ return True
+ if msg_type in ("extension_ready", "ping", "media_urls_refresh"):
+ session_id = data.get("session_id") or data.get("sessionId")
+ if session_id:
+ self.http_registry.touch(str(session_id))
return True
return False
@@ -320,7 +435,7 @@ async def api_request(self, url_path, body, captcha_action="VIDEO_GENERATION", m
REQUEST_MIN_INTERVAL seconds apart. Non-generation calls (polling,
credits — captcha_action="") bypass the limiter so they stay responsive.
"""
- if not self._ws:
+ if not self.is_extension_connected():
return {"error": "Extension not connected"}
# Only throttle credit/captcha-consuming generation calls.
@@ -389,52 +504,70 @@ async def _do_api_request(self, url_path, body, captcha_action, method, timeout,
finally:
self._pending.pop(req_id, None)
- def _start_http_server(self):
- """Start HTTP server for extension callbacks (runs in thread)."""
+ def _make_http_handler(self):
+ """Build the authenticated standalone callback request handler."""
bridge = self
class Handler(BaseHTTPRequestHandler):
+ def _send_json(self, status, payload):
+ encoded = json.dumps(payload, separators=(",", ":")).encode()
+ self.send_response(status)
+ self.send_header("Content-Type", "application/json")
+ self.send_header("Content-Length", str(len(encoded)))
+ self.end_headers()
+ self.wfile.write(encoded)
+
def do_POST(self):
- if self.path == "/api/ext/callback":
+ if self.path != "/api/ext/callback":
+ self._send_json(404, {"ok": False})
+ return
+
+ authorization = self.headers.get("Authorization")
+ if authorization is None:
+ self._send_json(401, {"ok": False})
+ return
+ if not bridge.verify_callback_authorization(authorization):
+ self._send_json(403, {"ok": False})
+ return
+
+ try:
length = int(self.headers.get("Content-Length", 0))
body = json.loads(self.rfile.read(length)) if length else {}
- bridge.handle_http_callback(body)
- self.send_response(200)
- self.send_header("Content-Type", "application/json")
- self.send_header("Access-Control-Allow-Origin", "*")
- self.end_headers()
- self.wfile.write(b'{"ok":true}')
- else:
- self.send_response(404)
- self.end_headers()
+ except (TypeError, ValueError, json.JSONDecodeError):
+ self._send_json(400, {"ok": False})
+ return
+
+ accepted = bridge.handle_http_callback(body)
+ self._send_json(200 if accepted else 404, {"ok": bool(accepted)})
def do_GET(self):
if self.path == "/health":
- self.send_response(200)
- self.send_header("Content-Type", "application/json")
- self.end_headers()
- self.wfile.write(json.dumps({
+ self._send_json(200, {
"status": "ok",
- "extension_connected": bridge._ws is not None,
- }).encode())
+ "extension_connected": bridge.is_extension_connected(),
+ "has_flow_key": bridge.has_flow_key(),
+ "transport": bridge.active_transport(),
+ })
else:
- self.send_response(404)
- self.end_headers()
+ self._send_json(404, {"ok": False})
def do_OPTIONS(self):
- self.send_response(200)
- self.send_header("Access-Control-Allow-Origin", "*")
- self.send_header("Access-Control-Allow-Methods", "POST, GET, OPTIONS")
- self.send_header("Access-Control-Allow-Headers", "Content-Type")
- self.end_headers()
+ # This loopback endpoint is not a public cross-origin API.
+ self._send_json(405, {"ok": False})
def log_message(self, *args):
pass
- server = HTTPServer(("127.0.0.1", HTTP_PORT), Handler)
+ return Handler
+
+ def _start_http_server(self):
+ """Start HTTP server for extension callbacks (runs in thread)."""
+ server = HTTPServer(("127.0.0.1", HTTP_PORT), self._make_http_handler())
thread = threading.Thread(target=server.serve_forever, daemon=True)
thread.start()
async def close(self):
- self._ws_server.close()
- await self._ws_server.wait_closed()
+ if self._ws_server is not None:
+ self._ws_server.close()
+ await self._ws_server.wait_closed()
+ self._ws_server = None
diff --git a/flow-agent/omniflash/config.py b/flow-agent/omniflash/config.py
index 515fbf9..1ed01da 100644
--- a/flow-agent/omniflash/config.py
+++ b/flow-agent/omniflash/config.py
@@ -109,6 +109,14 @@ def _load_env_file(path):
WS_PORT = int(os.environ.get("WS_PORT", "9227"))
HTTP_PORT = int(os.environ.get("HTTP_PORT", "8100"))
+# Extension transport: auto (HTTP first, WS fallback) | http | ws
+EXT_TRANSPORT = os.environ.get("EXT_TRANSPORT", "auto").strip().lower()
+EXT_SESSION_TTL_SEC = float(os.environ.get("EXT_SESSION_TTL_SEC", "20"))
+EXT_POLL_INTERVAL_MS = int(os.environ.get("EXT_POLL_INTERVAL_MS", "1000"))
+ENABLE_EXTENSION_WS = os.environ.get("ENABLE_EXTENSION_WS", "1").strip().lower() in {
+ "1", "true", "yes", "on",
+}
+
POLL_INTERVAL = int(os.environ.get("POLL_INTERVAL", "10"))
POLL_TIMEOUT = int(os.environ.get("POLL_TIMEOUT", "420"))
diff --git a/flow-agent/omniflash/http_bridge.py b/flow-agent/omniflash/http_bridge.py
new file mode 100644
index 0000000..c0bb38b
--- /dev/null
+++ b/flow-agent/omniflash/http_bridge.py
@@ -0,0 +1,194 @@
+"""In-memory HTTP extension session registry and command queue."""
+
+from __future__ import annotations
+
+import hmac
+import threading
+import time
+from collections import deque
+from dataclasses import dataclass, field
+from typing import Any
+
+
+@dataclass
+class Session:
+ session_id: str
+ secret: str
+ flow_key: str | None
+ last_seen: float
+ queue: deque = field(default_factory=deque)
+
+
+class ExtensionHttpRegistry:
+ """Thread-safe registry for extension HTTP sessions and command queues."""
+
+ def __init__(self, session_ttl_sec: float = 15.0) -> None:
+ self._session_ttl_sec = float(session_ttl_sec)
+ self._sessions: dict[str, Session] = {}
+ self._lock = threading.Lock()
+
+ def hello(
+ self,
+ session_id: str,
+ flow_key: str | None = None,
+ secret: str = "",
+ meta: dict[str, Any] | None = None,
+ ) -> dict[str, Any]:
+ """Register or refresh a session."""
+ now = time.monotonic()
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ existing = self._sessions.get(session_id)
+ if existing is not None:
+ existing.last_seen = now
+ if flow_key is not None:
+ existing.flow_key = flow_key
+ if secret:
+ existing.secret = secret
+ resolved_secret = existing.secret
+ else:
+ resolved_secret = secret or ""
+ self._sessions[session_id] = Session(
+ session_id=session_id,
+ secret=resolved_secret,
+ flow_key=flow_key,
+ last_seen=now,
+ )
+ return {
+ "ok": True,
+ "session_id": session_id,
+ "secret": resolved_secret,
+ }
+
+ def touch(self, session_id: str) -> bool:
+ """Refresh last_seen for an existing online session."""
+ now = time.monotonic()
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ session = self._get_online_unlocked(session_id, now)
+ if session is None:
+ return False
+ session.last_seen = now
+ return True
+
+ def is_connected(self, session_id: str | None = None) -> bool:
+ """Return whether the given session (or any session) is online."""
+ now = time.monotonic()
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ if session_id is None:
+ return self._has_online_unlocked(now)
+ return self._get_online_unlocked(session_id, now) is not None
+
+ def has_online_session(self) -> bool:
+ return self.is_connected(None)
+
+ def verify_authorization(self, session_id: str, authorization: str | None) -> bool:
+ """Constant-time Bearer check against the session secret."""
+ candidate = ""
+ if isinstance(authorization, str):
+ scheme, separator, value = authorization.partition(" ")
+ if separator and scheme.lower() == "bearer":
+ candidate = value.strip()
+ now = time.monotonic()
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ session = self._get_online_unlocked(session_id, now)
+ if session is None:
+ return False
+ return hmac.compare_digest(candidate, session.secret or "")
+
+ def enqueue(self, session_id: str | None, command: dict[str, Any]) -> bool:
+ """
+ Enqueue a command for a session.
+
+ If session_id is None, deliver to the most recently seen online session.
+ Duplicate command ids still present in the session queue are ignored
+ (in-flight/queue-only idempotency). After poll removes a command, the
+ same id may be enqueued again for business retries.
+ """
+ now = time.monotonic()
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ session = self._resolve_target_unlocked(session_id, now)
+ if session is None:
+ return False
+
+ cmd_id = command.get("id")
+ if cmd_id is not None:
+ cmd_id_str = str(cmd_id)
+ if any(str(item.get("id")) == cmd_id_str for item in session.queue):
+ return True
+
+ session.queue.append(dict(command))
+ return True
+
+ def poll(self, session_id: str, max_commands: int = 10) -> dict[str, Any]:
+ """Pop up to max_commands from the session queue. Each command is delivered once."""
+ now = time.monotonic()
+ max_n = max(0, int(max_commands))
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ session = self._get_online_unlocked(session_id, now)
+ if session is None:
+ return {"commands": [], "ok": False, "reason": "session_not_found"}
+ session.last_seen = now
+ commands: list[dict[str, Any]] = []
+ while session.queue and len(commands) < max_n:
+ commands.append(session.queue.popleft())
+ return {
+ "ok": True,
+ "commands": commands,
+ "session_id": session_id,
+ "server_time": int(time.time() * 1000),
+ }
+
+ def get_flow_key(self, session_id: str | None = None) -> str | None:
+ """Return flow_key for a session, or the most recently seen online session."""
+ now = time.monotonic()
+ with self._lock:
+ self._purge_expired_unlocked(now)
+ if session_id is not None:
+ session = self._get_online_unlocked(session_id, now)
+ return None if session is None else session.flow_key
+ latest = self._latest_online_unlocked(now)
+ return None if latest is None else latest.flow_key
+
+ def _is_online_unlocked(self, session: Session, now: float) -> bool:
+ return (now - session.last_seen) <= self._session_ttl_sec
+
+ def _purge_expired_unlocked(self, now: float) -> None:
+ expired = [
+ session_id
+ for session_id, session in self._sessions.items()
+ if not self._is_online_unlocked(session, now)
+ ]
+ for session_id in expired:
+ del self._sessions[session_id]
+
+ def _get_online_unlocked(self, session_id: str, now: float) -> Session | None:
+ session = self._sessions.get(session_id)
+ if session is None:
+ return None
+ if not self._is_online_unlocked(session, now):
+ return None
+ return session
+
+ def _has_online_unlocked(self, now: float) -> bool:
+ return self._latest_online_unlocked(now) is not None
+
+ def _latest_online_unlocked(self, now: float) -> Session | None:
+ latest: Session | None = None
+ for session in self._sessions.values():
+ if not self._is_online_unlocked(session, now):
+ continue
+ if latest is None or session.last_seen > latest.last_seen:
+ latest = session
+ return latest
+
+ def _resolve_target_unlocked(
+ self, session_id: str | None, now: float
+ ) -> Session | None:
+ if session_id is not None:
+ return self._get_online_unlocked(session_id, now)
+ return self._latest_online_unlocked(now)
diff --git a/flow-agent/tests/test_ext_http_api.py b/flow-agent/tests/test_ext_http_api.py
new file mode 100644
index 0000000..0d7f651
--- /dev/null
+++ b/flow-agent/tests/test_ext_http_api.py
@@ -0,0 +1,56 @@
+from fastapi.testclient import TestClient
+
+import cli.api as api
+
+
+def test_hello_and_poll_roundtrip():
+ with TestClient(api.app) as client:
+ r = client.post("/api/ext/hello", json={
+ "session_id": "s-test",
+ "flowKeyPresent": True,
+ "flowKey": "tok-1",
+ "extension_version": "1.0.0",
+ })
+ assert r.status_code == 200
+ body = r.json()
+ assert body["ok"] is True
+ secret = body["secret"]
+ assert secret
+ assert body["callback_url"].endswith("/api/ext/callback")
+ assert body["poll_url"].endswith("/api/ext/poll")
+ assert api.bridge is not None
+ api.bridge.enqueue_http_command({"id": "c1", "method": "get_status", "params": {}})
+ p = client.get(
+ "/api/ext/poll",
+ params={"session_id": "s-test"},
+ headers={"Authorization": f"Bearer {secret}"},
+ )
+ assert p.status_code == 200
+ assert p.json()["commands"][0]["id"] == "c1"
+
+
+def test_poll_requires_authorization():
+ with TestClient(api.app) as client:
+ r = client.post("/api/ext/hello", json={
+ "session_id": "s-auth",
+ "flowKey": "tok",
+ })
+ assert r.status_code == 200
+ p = client.get("/api/ext/poll", params={"session_id": "s-auth"})
+ assert p.status_code == 401
+
+
+def test_health_reports_http_transport():
+ with TestClient(api.app) as client:
+ r = client.post("/api/ext/hello", json={
+ "session_id": "s-health",
+ "flowKeyPresent": True,
+ "flowKey": "tok-health",
+ })
+ assert r.status_code == 200
+ h = client.get("/health")
+ assert h.status_code == 200
+ body = h.json()
+ assert body["extension_connected"] is True
+ assert body["has_flow_key"] is True
+ assert body["transport"] == "http"
diff --git a/flow-agent/tests/test_http_bridge.py b/flow-agent/tests/test_http_bridge.py
new file mode 100644
index 0000000..9a21de5
--- /dev/null
+++ b/flow-agent/tests/test_http_bridge.py
@@ -0,0 +1,70 @@
+# flow-agent/tests/test_http_bridge.py
+import time
+from omniflash.http_bridge import ExtensionHttpRegistry
+
+
+def test_hello_registers_session_and_accepts_token():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ out = reg.hello(session_id="s1", flow_key="tok", secret="test-secret")
+ assert out["ok"] is True
+ assert reg.is_connected("s1") is True
+ assert reg.get_flow_key("s1") == "tok"
+
+
+def test_enqueue_and_poll_returns_command_once():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ reg.hello(session_id="s1", flow_key="tok", secret="sec")
+ reg.enqueue("s1", {"id": "r1", "method": "get_status", "params": {}})
+ first = reg.poll("s1")
+ second = reg.poll("s1")
+ assert len(first["commands"]) == 1
+ assert first["commands"][0]["id"] == "r1"
+ assert second["commands"] == []
+
+
+def test_session_expires():
+ reg = ExtensionHttpRegistry(session_ttl_sec=1)
+ reg.hello(session_id="s1", flow_key="tok", secret="sec")
+ time.sleep(1.2)
+ assert reg.is_connected("s1") is False
+
+
+def test_duplicate_command_id_not_double_queued():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ reg.hello(session_id="s1", flow_key="tok", secret="sec")
+ cmd = {"id": "r1", "method": "get_status", "params": {}}
+ assert reg.enqueue("s1", cmd) is True
+ assert reg.enqueue("s1", cmd) is True
+ first = reg.poll("s1")
+ second = reg.poll("s1")
+ assert len(first["commands"]) == 1
+ assert first["commands"][0]["id"] == "r1"
+ assert second["commands"] == []
+
+
+def test_same_command_id_can_requeue_after_poll():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ reg.hello(session_id="s1", flow_key="tok", secret="sec")
+ cmd = {"id": "r1", "method": "get_status", "params": {}}
+ assert reg.enqueue("s1", cmd) is True
+ first = reg.poll("s1")
+ assert len(first["commands"]) == 1
+ assert first["commands"][0]["id"] == "r1"
+ assert reg.enqueue("s1", cmd) is True
+ second = reg.poll("s1")
+ assert len(second["commands"]) == 1
+ assert second["commands"][0]["id"] == "r1"
+
+
+def test_enqueue_none_goes_to_latest_online():
+ reg = ExtensionHttpRegistry(session_ttl_sec=15)
+ reg.hello(session_id="s1", flow_key="tok1", secret="sec1")
+ # Windows monotonic clock resolution can be ~15ms; sleep enough to order last_seen.
+ time.sleep(0.05)
+ reg.hello(session_id="s2", flow_key="tok2", secret="sec2")
+ assert reg.enqueue(None, {"id": "r1", "method": "get_status", "params": {}}) is True
+ s1 = reg.poll("s1")
+ s2 = reg.poll("s2")
+ assert s1["commands"] == []
+ assert len(s2["commands"]) == 1
+ assert s2["commands"][0]["id"] == "r1"
diff --git a/flow-chrome-extension/background.js b/flow-chrome-extension/background.js
index bffd9e6..1e4773f 100644
--- a/flow-chrome-extension/background.js
+++ b/flow-chrome-extension/background.js
@@ -5,14 +5,24 @@
* Captures bearer token, solves reCAPTCHA, proxies API calls through browser.
*/
+const AGENT_BASE = 'http://127.0.0.1:8001';
const AGENT_WS_URL = 'ws://127.0.0.1:8001/ws';
-let callbackUrl = 'http://127.0.0.1:3001/api/ext/callback';
+const AGENT_HELLO_URL = `${AGENT_BASE}/api/ext/hello`;
+const AGENT_POLL_URL = `${AGENT_BASE}/api/ext/poll`;
+const AGENT_CALLBACK_URL = `${AGENT_BASE}/api/ext/callback`;
+const TRANSPORT_MODE = 'auto'; // auto | http | ws
+let callbackUrl = AGENT_CALLBACK_URL;
// NOTE: This is a browser-restricted public API key — safe to ship in extension bundles.
const API_KEY = 'AIzaSyBtrm0o5ab1c-Ec8ZuLcGt3oJAA5VWt3pY';
let ws = null;
let flowKey = null;
-let callbackSecret = null; // Auth secret for HTTP callback, received from server on WS connect
+let callbackSecret = null; // Auth secret for HTTP callback / poll
+let httpSessionId = null;
+let httpConnected = false;
+let httpPollTimer = null;
+let httpPollIntervalMs = 1000;
+let activeTransport = 'none'; // none | http | ws
let state = 'off'; // off | idle | running
let manualDisconnect = false;
let metrics = {
@@ -69,17 +79,25 @@ chrome.alarms.onAlarm.addListener(async (alarm) => {
if (alarm.name === 'reconnect') connectToAgent();
if (alarm.name === 'keepAlive') keepAlive();
if (alarm.name === 'flushOutbox') flushOutbox();
+ if (alarm.name === 'http-poll') pollAgentCommands();
if (alarm.name === 'token-refresh') {
await captureTokenFromFlowTab();
}
});
async function init() {
- const data = await chrome.storage.local.get(['flowKey', 'metrics', 'callbackSecret', 'callbackUrl']);
+ const data = await chrome.storage.local.get([
+ 'flowKey', 'metrics', 'callbackSecret', 'callbackUrl', 'httpSessionId',
+ ]);
if (data.flowKey) flowKey = data.flowKey;
if (data.metrics) Object.assign(metrics, data.metrics);
if (data.callbackSecret) callbackSecret = data.callbackSecret;
if (data.callbackUrl) callbackUrl = data.callbackUrl;
+ if (data.httpSessionId) httpSessionId = data.httpSessionId;
+ if (!httpSessionId) {
+ httpSessionId = crypto.randomUUID();
+ chrome.storage.local.set({ httpSessionId });
+ }
await loadOutbox();
connectToAgent();
chrome.alarms.create('keepAlive', { periodInMinutes: 0.4 });
@@ -108,10 +126,8 @@ chrome.webRequest.onBeforeSendHeaders.addListener(
chrome.storage.local.set({ flowKey, metrics });
console.log('[Flow Agent] Bearer token captured');
- // Notify agent
- if (ws?.readyState === WebSocket.OPEN) {
- ws.send(JSON.stringify({ type: 'token_captured', flowKey }));
- }
+ // Notify agent (HTTP callback preferred; WS fallback)
+ sendToAgent({ type: 'token_captured', flowKey, session_id: httpSessionId });
},
{ urls: ['https://aisandbox-pa.googleapis.com/*', 'https://labs.google/*'] },
['requestHeaders', 'extraHeaders'],
@@ -163,9 +179,217 @@ async function captureTokenFromFlowTab() {
}
}
-// ─── WebSocket to Agent ─────────────────────────────────────
+// ─── Agent Transport (HTTP preferred, WS fallback) ──────────
+
+function isAgentConnected() {
+ return httpConnected || ws?.readyState === WebSocket.OPEN;
+}
+
+async function ensureSessionId() {
+ if (httpSessionId) return httpSessionId;
+ const data = await chrome.storage.local.get(['httpSessionId']);
+ httpSessionId = data.httpSessionId || crypto.randomUUID();
+ await chrome.storage.local.set({ httpSessionId });
+ return httpSessionId;
+}
+
+async function connectViaHttp() {
+ if (manualDisconnect) return false;
+ try {
+ const sessionId = await ensureSessionId();
+ const resp = await fetch(AGENT_HELLO_URL, {
+ method: 'POST',
+ headers: { 'Content-Type': 'application/json' },
+ body: JSON.stringify({
+ type: 'hello',
+ session_id: sessionId,
+ extension_version: chrome.runtime.getManifest().version,
+ flowKeyPresent: !!flowKey,
+ flowKey: flowKey || undefined,
+ capabilities: ['api_request', 'trpc_request', 'upload_video', 'solve_captcha', 'get_status'],
+ }),
+ });
+ if (!resp.ok) {
+ console.warn('[Flow Agent] HTTP hello failed:', resp.status);
+ return false;
+ }
+ const body = await resp.json();
+ if (!body?.ok) return false;
+
+ callbackSecret = body.secret || callbackSecret;
+ callbackUrl = body.callback_url || AGENT_CALLBACK_URL;
+ httpPollIntervalMs = Number(body.poll_interval_ms) || 1000;
+ httpConnected = true;
+ activeTransport = 'http';
+ chrome.storage.local.set({
+ callbackSecret,
+ callbackUrl,
+ httpSessionId: sessionId,
+ });
+
+ chrome.alarms.clear('reconnect');
+ chrome.alarms.create('token-refresh', { periodInMinutes: 45 });
+ // MV3 SW may sleep; alarm + soft timer both keep polling alive.
+ chrome.alarms.create('http-poll', { periodInMinutes: Math.max(0.05, httpPollIntervalMs / 60000) });
+ scheduleHttpPoll();
+ setState('idle');
+ console.log('[Flow Agent] Connected to agent via HTTP');
+
+ await sendToAgent({
+ type: 'extension_ready',
+ flowKeyPresent: !!flowKey,
+ session_id: sessionId,
+ tokenAge: flowKey && metrics.tokenCapturedAt ? Date.now() - metrics.tokenCapturedAt : null,
+ });
+ if (flowKey) {
+ await sendToAgent({ type: 'token_captured', flowKey, session_id: sessionId });
+ }
+ flushOutbox();
+ return true;
+ } catch (e) {
+ console.warn('[Flow Agent] HTTP connect error:', e);
+ httpConnected = false;
+ if (activeTransport === 'http') activeTransport = 'none';
+ return false;
+ }
+}
+
+function scheduleHttpPoll() {
+ if (httpPollTimer) clearTimeout(httpPollTimer);
+ if (!httpConnected || manualDisconnect) return;
+ httpPollTimer = setTimeout(() => {
+ pollAgentCommands().finally(scheduleHttpPoll);
+ }, httpPollIntervalMs);
+}
+
+async function pollAgentCommands() {
+ if (manualDisconnect || !httpConnected) return;
+ try {
+ const sessionId = await ensureSessionId();
+ const headers = {};
+ if (callbackSecret) headers.Authorization = `Bearer ${callbackSecret}`;
+ const resp = await fetch(
+ `${AGENT_POLL_URL}?session_id=${encodeURIComponent(sessionId)}`,
+ { method: 'GET', headers },
+ );
+ if (resp.status === 401 || resp.status === 403 || resp.status === 404) {
+ httpConnected = false;
+ activeTransport = ws?.readyState === WebSocket.OPEN ? 'ws' : 'none';
+ return;
+ }
+ if (!resp.ok) return;
+ const body = await resp.json();
+ httpConnected = true;
+ activeTransport = 'http';
+ const commands = Array.isArray(body?.commands) ? body.commands : [];
+ for (const cmd of commands) {
+ await handleAgentMessage(cmd);
+ }
+ } catch (e) {
+ // Transient network error — keep session; next poll/hello will recover.
+ console.debug('[Flow Agent] HTTP poll error:', e);
+ }
+}
-function connectToAgent() {
+async function handleAgentMessage(msg) {
+ try {
+ if (msg.method === 'api_request') {
+ await handleApiRequest(msg);
+ } else if (msg.method === 'trpc_request') {
+ await handleTrpcRequest(msg);
+ } else if (msg.method === 'upload_video') {
+ await handleUploadVideo(msg);
+ } else if (msg.method === 'solve_captcha') {
+ await handleSolveCaptcha(msg);
+ } else if (msg.method === 'get_status') {
+ sendToAgent({
+ id: msg.id,
+ result: {
+ state,
+ flowKeyPresent: !!flowKey,
+ manualDisconnect,
+ transport: activeTransport,
+ httpConnected,
+ tokenAge: metrics.tokenCapturedAt ? Date.now() - metrics.tokenCapturedAt : null,
+ metrics,
+ },
+ });
+ } else if (msg.method === 'open_flow_tab') {
+ console.log('[Flow Agent] Agent requested: open Flow tab');
+ const tabs = await chrome.tabs.query({
+ url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'],
+ });
+ if (tabs.length) {
+ await chrome.tabs.reload(tabs[0].id);
+ console.log('[Flow Agent] Refreshed existing Flow tab');
+ } else {
+ await chrome.tabs.create({ url: 'https://labs.google/fx/tools/flow', active: true });
+ console.log('[Flow Agent] Opened new Flow tab');
+ }
+ await sleep(5000);
+ if (flowKey) {
+ sendToAgent({ type: 'token_captured', flowKey, session_id: httpSessionId });
+ console.log('[Flow Agent] Sent stored token after tab open');
+ } else {
+ const data = await chrome.storage.local.get(['flowKey']);
+ if (data.flowKey) {
+ flowKey = data.flowKey;
+ sendToAgent({ type: 'token_captured', flowKey, session_id: httpSessionId });
+ console.log('[Flow Agent] Sent token from storage after tab open');
+ }
+ }
+ } else if (msg.method === 'refresh_flow_tab') {
+ console.log('[Flow Agent] Agent requested: refresh token');
+ await captureTokenFromFlowTab();
+ await sleep(3000);
+ if (flowKey) {
+ sendToAgent({ type: 'token_captured', flowKey, session_id: httpSessionId });
+ console.log('[Flow Agent] Sent token after refresh');
+ } else {
+ const data = await chrome.storage.local.get(['flowKey']);
+ if (data.flowKey) {
+ flowKey = data.flowKey;
+ sendToAgent({ type: 'token_captured', flowKey, session_id: httpSessionId });
+ console.log('[Flow Agent] Sent token from storage after refresh');
+ }
+ }
+ } else if (msg.type === 'callback_config') {
+ callbackSecret = msg.secret;
+ callbackUrl = msg.callback_url || AGENT_CALLBACK_URL;
+ chrome.storage.local.set({ callbackSecret: msg.secret, callbackUrl });
+ console.log('[Flow Agent] Received callback config:', callbackUrl);
+ } else if (msg.type === 'callback_secret') {
+ callbackSecret = msg.secret;
+ chrome.storage.local.set({ callbackSecret: msg.secret });
+ console.log('[Flow Agent] Received callback secret');
+ } else if (msg.type === 'pong') {
+ // keepalive response
+ }
+ } catch (e) {
+ console.error('[Flow Agent] Message error:', e);
+ }
+}
+
+async function connectToAgent() {
+ if (manualDisconnect) return;
+
+ if (TRANSPORT_MODE === 'http' || TRANSPORT_MODE === 'auto') {
+ const ok = await connectViaHttp();
+ if (ok) {
+ // HTTP success: do not force WS in auto mode.
+ if (TRANSPORT_MODE === 'auto') return;
+ } else if (TRANSPORT_MODE === 'http') {
+ scheduleReconnect();
+ return;
+ }
+ }
+
+ if (TRANSPORT_MODE === 'ws' || TRANSPORT_MODE === 'auto') {
+ connectViaWebSocket();
+ }
+}
+
+function connectViaWebSocket() {
if (manualDisconnect) return;
if (ws?.readyState === WebSocket.CONNECTING) return;
if (ws?.readyState === WebSocket.OPEN) return;
@@ -179,14 +403,12 @@ function connectToAgent() {
}
ws.onopen = () => {
- console.log('[Flow Agent] Connected to agent');
+ console.log('[Flow Agent] Connected to agent via WebSocket');
+ if (!httpConnected) activeTransport = 'ws';
chrome.alarms.clear('reconnect');
setState('idle');
-
- // Token refresh alarm — 45 min gives buffer before ~60 min expiry
chrome.alarms.create('token-refresh', { periodInMinutes: 45 });
- // Send current state + resend token if we have one
ws.send(JSON.stringify({
type: 'extension_ready',
flowKeyPresent: !!flowKey,
@@ -195,105 +417,25 @@ function connectToAgent() {
if (flowKey) {
ws.send(JSON.stringify({ type: 'token_captured', flowKey }));
}
- // Backend is reachable again — push any responses queued while it was down.
flushOutbox();
};
ws.onmessage = async ({ data }) => {
try {
const msg = JSON.parse(data);
-
- if (msg.method === 'api_request') {
- await handleApiRequest(msg);
- } else if (msg.method === 'trpc_request') {
- await handleTrpcRequest(msg);
- } else if (msg.method === 'upload_video') {
- await handleUploadVideo(msg);
- } else if (msg.method === 'solve_captcha') {
- await handleSolveCaptcha(msg);
- } else if (msg.method === 'get_status') {
- sendToAgent({
- id: msg.id,
- result: {
- state,
- flowKeyPresent: !!flowKey,
- manualDisconnect,
- tokenAge: metrics.tokenCapturedAt ? Date.now() - metrics.tokenCapturedAt : null,
- metrics,
- },
- });
- } else if (msg.method === 'open_flow_tab') {
- // Python bridge asks us to open/focus a Flow tab
- console.log('[Flow Agent] Agent requested: open Flow tab');
- const tabs = await chrome.tabs.query({
- url: ['https://labs.google/fx/tools/flow*', 'https://labs.google/fx/*/tools/flow*'],
- });
- if (tabs.length) {
- // Tab exists — refresh it to trigger fresh API calls → token capture
- await chrome.tabs.reload(tabs[0].id);
- console.log('[Flow Agent] Refreshed existing Flow tab');
- } else {
- // No tab — open one (active so it loads properly)
- await chrome.tabs.create({ url: 'https://labs.google/fx/tools/flow', active: true });
- console.log('[Flow Agent] Opened new Flow tab');
- }
- // Wait for page to load and make API calls that trigger token capture
- await sleep(5000);
- // If token was captured by webRequest during page load, send it
- if (flowKey && ws?.readyState === WebSocket.OPEN) {
- ws.send(JSON.stringify({ type: 'token_captured', flowKey }));
- console.log('[Flow Agent] Sent stored token after tab open');
- } else {
- // Try reading from storage as fallback
- const data = await chrome.storage.local.get(['flowKey']);
- if (data.flowKey) {
- flowKey = data.flowKey;
- if (ws?.readyState === WebSocket.OPEN) {
- ws.send(JSON.stringify({ type: 'token_captured', flowKey }));
- console.log('[Flow Agent] Sent token from storage after tab open');
- }
- }
- }
- } else if (msg.method === 'refresh_flow_tab') {
- // Python bridge asks us to refresh token
- console.log('[Flow Agent] Agent requested: refresh token');
- await captureTokenFromFlowTab();
- await sleep(3000);
- // Actively send token if we have one
- if (flowKey && ws?.readyState === WebSocket.OPEN) {
- ws.send(JSON.stringify({ type: 'token_captured', flowKey }));
- console.log('[Flow Agent] Sent token after refresh');
- } else {
- const data = await chrome.storage.local.get(['flowKey']);
- if (data.flowKey) {
- flowKey = data.flowKey;
- if (ws?.readyState === WebSocket.OPEN) {
- ws.send(JSON.stringify({ type: 'token_captured', flowKey }));
- console.log('[Flow Agent] Sent token from storage after refresh');
- }
- }
- }
- } else if (msg.type === 'callback_config') {
- callbackSecret = msg.secret;
- callbackUrl = msg.callback_url;
- chrome.storage.local.set({ callbackSecret: msg.secret, callbackUrl: msg.callback_url });
- console.log('[Flow Agent] Received callback config:', callbackUrl);
- } else if (msg.type === 'callback_secret') {
- callbackSecret = msg.secret;
- chrome.storage.local.set({ callbackSecret: msg.secret });
- console.log('[Flow Agent] Received callback secret');
- } else if (msg.type === 'pong') {
- // keepalive response
- }
+ await handleAgentMessage(msg);
} catch (e) {
console.error('[Flow Agent] Message error:', e);
}
};
ws.onclose = () => {
- setState('off');
+ if (activeTransport === 'ws') {
+ setState(httpConnected ? 'idle' : 'off');
+ activeTransport = httpConnected ? 'http' : 'none';
+ }
chrome.alarms.clear('token-refresh');
- if (!manualDisconnect) scheduleReconnect();
+ if (!manualDisconnect && !httpConnected) scheduleReconnect();
};
ws.onerror = (e) => {
@@ -308,6 +450,12 @@ function scheduleReconnect() {
}
function keepAlive() {
+ if (httpConnected) {
+ // Soft hello refresh keeps session last_seen fresh.
+ connectViaHttp().catch(() => {});
+ pollAgentCommands().catch(() => {});
+ return;
+ }
if (ws?.readyState === WebSocket.OPEN) {
ws.send(JSON.stringify({ type: 'ping' }));
} else {
@@ -315,16 +463,39 @@ function keepAlive() {
}
}
-function sendToAgent(msg) {
+async function sendToAgent(msg) {
// API responses (with msg.id) go through a durable outbox so a generated
// result is never lost — persisted and retried until the agent acks it.
if (msg.id) {
enqueueResponse(msg);
return;
}
- // Non-response messages (ping, status, token) — best-effort over WS.
+
+ const payload = { ...msg };
+ if (httpSessionId && !payload.session_id && !payload.sessionId) {
+ payload.session_id = httpSessionId;
+ }
+
+ // Prefer authenticated HTTP callback for control messages.
+ if (callbackSecret || httpConnected) {
+ try {
+ const headers = { 'Content-Type': 'application/json' };
+ if (callbackSecret) headers.Authorization = `Bearer ${callbackSecret}`;
+ const target = callbackUrl || AGENT_CALLBACK_URL;
+ const resp = await fetch(target, {
+ method: 'POST',
+ headers,
+ body: JSON.stringify(payload),
+ });
+ if (resp.ok || resp.status === 404) return;
+ } catch (e) {
+ console.debug('[Flow Agent] HTTP send failed, falling back:', e);
+ }
+ }
+
+ // Non-response messages — best-effort over WS.
if (ws?.readyState === WebSocket.OPEN) {
- ws.send(JSON.stringify(msg));
+ ws.send(JSON.stringify(payload));
}
}
@@ -356,10 +527,17 @@ function enqueueResponse(msg) {
async function deliverOnce(entry) {
try {
- const resp = await fetch(callbackUrl, {
+ const headers = { 'Content-Type': 'application/json' };
+ if (callbackSecret) headers.Authorization = `Bearer ${callbackSecret}`;
+ const target = callbackUrl || AGENT_CALLBACK_URL;
+ const payload = { ...entry.msg };
+ if (httpSessionId && !payload.session_id && !payload.sessionId) {
+ payload.session_id = httpSessionId;
+ }
+ const resp = await fetch(target, {
method: 'POST',
- headers: { 'Content-Type': 'application/json' },
- body: JSON.stringify(entry.msg),
+ headers,
+ body: JSON.stringify(payload),
});
// Any HTTP reply means the backend is reachable and has taken the response
// (ok:true = matched a request, ok:false = unknown id / already handled).
@@ -732,8 +910,10 @@ function broadcastStatus() {
chrome.runtime.onMessage.addListener((msg, _, reply) => {
if (msg.type === 'STATUS') {
reply({
- connected: ws?.readyState === WebSocket.OPEN,
- agentConnected: ws?.readyState === WebSocket.OPEN,
+ connected: isAgentConnected(),
+ agentConnected: isAgentConnected(),
+ transport: activeTransport,
+ httpConnected,
flowKeyPresent: !!flowKey,
manualDisconnect,
tokenAge: metrics.tokenCapturedAt ? Date.now() - metrics.tokenCapturedAt : null,
@@ -749,6 +929,13 @@ chrome.runtime.onMessage.addListener((msg, _, reply) => {
if (msg.type === 'DISCONNECT') {
manualDisconnect = true;
+ httpConnected = false;
+ activeTransport = 'none';
+ if (httpPollTimer) {
+ clearTimeout(httpPollTimer);
+ httpPollTimer = null;
+ }
+ chrome.alarms.clear('http-poll');
if (ws) ws.close();
reply({ ok: true });
return true;
@@ -804,17 +991,22 @@ chrome.runtime.onMessage.addListener((msg, _, reply) => {
if (msg.type === 'SNIFFED_AISANDBOX_REQUEST') {
console.log('[Flow Agent] SNIFFED aisandbox request:', msg.url);
- fetch('http://127.0.0.1:8100/api/ext/callback', {
- method: 'POST',
- headers: { 'Content-Type': 'application/json' },
- body: JSON.stringify({
- type: 'sniffed_video_request',
- url: msg.url,
- method: msg.method,
- payload: msg.payload,
- timestamp: msg.timestamp,
- }),
- }).catch((e) => console.error('[Flow Agent] Failed to forward sniffed request:', e));
+ {
+ const headers = { 'Content-Type': 'application/json' };
+ if (callbackSecret) headers.Authorization = `Bearer ${callbackSecret}`;
+ fetch((callbackUrl || AGENT_CALLBACK_URL), {
+ method: 'POST',
+ headers,
+ body: JSON.stringify({
+ type: 'sniffed_video_request',
+ url: msg.url,
+ method: msg.method,
+ payload: msg.payload,
+ timestamp: msg.timestamp,
+ session_id: httpSessionId,
+ }),
+ }).catch((e) => console.error('[Flow Agent] Failed to forward sniffed request:', e));
+ }
reply({ ok: true });
return true;
}
diff --git a/flow-chrome-extension/manifest.json b/flow-chrome-extension/manifest.json
index 5a9d1cf..770b765 100644
--- a/flow-chrome-extension/manifest.json
+++ b/flow-chrome-extension/manifest.json
@@ -2,18 +2,28 @@
"manifest_version": 3,
"name": "Flow Agent",
"version": "1.0.0",
- "description": "Automate Google Flow — T2V, V2V, I2V video + T2I, I2I unlimited image generation from terminal",
+ "description": "Automate Google Flow \u2014 T2V, V2V, I2V video + T2I, I2I unlimited image generation from terminal",
"icons": {
"16": "icon16.png",
"48": "icon48.png",
"128": "icon128.png"
},
- "permissions": ["storage", "alarms", "tabs", "webRequest", "scripting", "declarativeNetRequest", "sidePanel"],
+ "permissions": [
+ "storage",
+ "alarms",
+ "tabs",
+ "webRequest",
+ "scripting",
+ "declarativeNetRequest",
+ "sidePanel"
+ ],
"host_permissions": [
"https://labs.google/*",
"https://aisandbox-pa.googleapis.com/*",
"https://aisandbox-pa.sandbox.googleapis.com/*",
"https://storage.googleapis.com/*",
+ "http://127.0.0.1:8001/*",
+ "http://localhost:8001/*",
"http://127.0.0.1:8100/*"
],
"background": {
@@ -25,14 +35,20 @@
"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://labs.google/*"
+ ]
}
],
"declarative_net_request": {
diff --git a/flow-chrome-extension/side_panel.html b/flow-chrome-extension/side_panel.html
index e10fe92..84b218f 100644
--- a/flow-chrome-extension/side_panel.html
+++ b/flow-chrome-extension/side_panel.html
@@ -135,7 +135,32 @@
flex-shrink: 0;
}
- #conn-dot.on {
+
+#transport-status.transport {
+ font-size: 10px;
+ font-weight: 700;
+ letter-spacing: 0.04em;
+ padding: 3px 7px;
+ border-radius: 999px;
+ background: rgba(255,255,255,0.08);
+ color: rgba(255,255,255,0.45);
+ margin-right: 6px;
+ flex-shrink: 0;
+}
+#transport-status.transport.on {
+ color: #9ef0c2;
+ background: rgba(46, 204, 113, 0.14);
+}
+#transport-status.transport.http {
+ color: #8fd3ff;
+ background: rgba(52, 152, 219, 0.16);
+}
+#transport-status.transport.ws {
+ color: #ffd27f;
+ background: rgba(241, 196, 15, 0.16);
+}
+
+#conn-dot.on {
background: var(--green);
box-shadow: 0 0 8px 2px rgba(16, 185, 129, 0.4);
animation: pulse-glow 2.5s ease-in-out infinite;
@@ -1023,6 +1048,7 @@
+ OFF
OFF