import logging import os from uuid import uuid4 logger = logging.getLogger(__name__) from fastapi.responses import JSONResponse, HTMLResponse from fastapi import Request # In-memory store for pending OAuth flows (state -> {provider, code_verifier, redirect_uri}) _pending_oauth: dict[str, dict] = {} # Recently-completed OAuth states so the /api/subscriptions/callback handler # can distinguish a legitimate duplicate callback (browser prefetch, refresh, # or Google redirect retry after a slow first response) from a truly stale # request. Bounded FIFO — drops the oldest entries once it grows past # _MAX_COMPLETED_OAUTH so it can't leak memory. _completed_oauth: list[str] = [] _MAX_COMPLETED_OAUTH = 64 def _mark_oauth_completed(state: str) -> None: if state in _completed_oauth: return _completed_oauth.append(state) # Trim head if we've outgrown the bound while len(_completed_oauth) > _MAX_COMPLETED_OAUTH: _completed_oauth.pop(0) from backend.config.Apps import MainApp from backend.apps.health.health import health from backend.apps.agents.agents import agents from backend.apps.agents.ws_manager import ws_manager from backend.apps.skills.skills import skills from backend.apps.tools_lib.tools_lib import tools_lib from backend.apps.modes.modes import modes from backend.apps.settings.settings import settings from backend.apps.mcp_registry.mcp_registry import mcp_registry from backend.apps.skill_registry.skill_registry import skill_registry from backend.apps.outputs.outputs import outputs from backend.apps.dashboards.dashboards import dashboards from backend.apps.analytics.analytics import analytics from backend.apps.subscription.router import subscription from fastapi.middleware.cors import CORSMiddleware from fastapi import WebSocket, WebSocketDisconnect import json main_app = MainApp([health, agents, skills, tools_lib, modes, settings, mcp_registry, skill_registry, outputs, dashboards, analytics, subscription]) app = main_app.app app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) # Chrome's Private Network Access check: a page on https://api.openswarm.com # POSTing to http://127.0.0.1:8324 triggers a preflight that requires this # header. Without it the request is blocked and the post-checkout activation # flow silently fails. Harmless on every other request — it's only read when # the browser is crossing from a public origin into a private network. @app.middleware("http") async def _allow_private_network(request, call_next): response = await call_next(request) response.headers["Access-Control-Allow-Private-Network"] = "true" return response @app.websocket("/ws/agents/{session_id}") async def websocket_session(websocket: WebSocket, session_id: str): await ws_manager.connect_session(session_id, websocket) try: while True: data = await websocket.receive_text() msg = json.loads(data) event = msg.get("event") payload = msg.get("data", {}) if event == "agent:send_message": from backend.apps.agents.agent_manager import agent_manager await agent_manager.send_message( session_id, payload.get("prompt", ""), mode=payload.get("mode"), model=payload.get("model"), provider=payload.get("provider"), images=payload.get("images"), ) elif event == "agent:approval_response": from backend.apps.agents.agent_manager import agent_manager agent_manager.handle_approval(payload.get("request_id"), { "behavior": payload.get("behavior", "deny"), "message": payload.get("message"), "updated_input": payload.get("updated_input"), }) elif event == "agent:edit_message": from backend.apps.agents.agent_manager import agent_manager await agent_manager.edit_message( session_id, payload.get("message_id", ""), payload.get("content", ""), ) elif event == "agent:stop": from backend.apps.agents.agent_manager import agent_manager await agent_manager.stop_agent(session_id) except WebSocketDisconnect: ws_manager.disconnect_session(session_id, websocket) @app.websocket("/ws/dashboard") async def websocket_dashboard(websocket: WebSocket): await ws_manager.connect_global(websocket) try: while True: data = await websocket.receive_text() msg = json.loads(data) event = msg.get("event") payload = msg.get("data", {}) if event == "agent:approval_response": from backend.apps.agents.agent_manager import agent_manager agent_manager.handle_approval(payload.get("request_id"), { "behavior": payload.get("behavior", "deny"), "message": payload.get("message"), "updated_input": payload.get("updated_input"), }) elif event == "browser:result": ws_manager.resolve_browser_command( payload.get("request_id", ""), payload, ) except WebSocketDisconnect: ws_manager.disconnect_global(websocket) @app.post("/api/browser/command") async def browser_command(request: Request): """HTTP endpoint called by the browser MCP server subprocess. Proxies commands to the frontend via WebSocket and waits for results.""" body = await request.json() action = body.get("action", "") browser_id = body.get("browser_id", "") tab_id = body.get("tab_id", "") params = body.get("params", {}) if not action or not browser_id: return JSONResponse({"error": "action and browser_id are required"}, status_code=400) request_id = uuid4().hex result = await ws_manager.send_browser_command(request_id, action, browser_id, params, tab_id=tab_id) return JSONResponse(result) @app.get("/api/subscriptions/pending/{state}") async def subscriptions_pending(state: str): """Return pending OAuth data for a state param. Called by 9Router's callback page.""" pending = _pending_oauth.get(state) if not pending: return JSONResponse({"error": "not found"}, status_code=404, headers={"Access-Control-Allow-Origin": "*"}) return JSONResponse({ "provider": pending["provider"], "code_verifier": pending["code_verifier"], "redirect_uri": pending["redirect_uri"], }, headers={"Access-Control-Allow-Origin": "*"}) _SUCCESS_HTML = ( '
' 'You can close this window
' '{desc}
Please try connecting again.
{e}