"""9Router / subscription management endpoints. Moved from agents.py and main.py to a dedicated sub-app. """ from __future__ import annotations import asyncio import logging import httpx from contextlib import asynccontextmanager from fastapi import HTTPException, Request from fastapi.responses import JSONResponse, HTMLResponse from backend.config.Apps import SubApp from backend.apps.nine_router import poll_oauth from backend.apps.analytics.collector import record as _analytics from backend.apps.nine_router import ( is_running, ensure_running, get_providers, get_models, start_oauth, exchange_oauth, NINE_ROUTER_API, ) logger = logging.getLogger(__name__) _pending_oauth: dict[str, dict] = {} _ensure_task: asyncio.Task | None = None @asynccontextmanager async def subscriptions_lifespan(): logger.info("Subscriptions sub-app starting") yield global _ensure_task if _ensure_task and not _ensure_task.done(): _ensure_task.cancel() logger.info("Subscriptions sub-app shutting down") subscriptions = SubApp("subscriptions", subscriptions_lifespan) @subscriptions.router.get("/status") async def subscriptions_status(): """Check if 9Router is running and list connected providers.""" global _ensure_task if not is_running(): if _ensure_task is None or _ensure_task.done(): _ensure_task = asyncio.create_task(ensure_running()) return {"running": False, "providers": [], "models": []} providers = await get_providers() models = await get_models() return {"running": True, "providers": providers, "models": models} @subscriptions.router.post("/connect") async def subscriptions_connect(body: dict): """Start OAuth flow for a subscription provider.""" provider = body.get("provider", "") if not provider: raise HTTPException(status_code=400, detail="provider required") if not is_running(): await ensure_running() if not is_running(): raise HTTPException(status_code=503, detail="9Router not available. Please install Node.js.") try: result = await start_oauth(provider) if result.get("flow") == "authorization_code" and result.get("state"): _pending_oauth[result["state"]] = { "provider": provider, "code_verifier": result.get("code_verifier", ""), "redirect_uri": result.get("redirect_uri", ""), } return result except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @subscriptions.router.post("/poll") async def subscriptions_poll(body: dict): """Poll for OAuth completion.""" provider = body.get("provider", "") device_code = body.get("device_code", "") if not provider or not device_code: raise HTTPException(status_code=400, detail="provider and device_code required") try: result = await poll_oauth( provider, device_code, code_verifier=body.get("code_verifier"), extra_data=body.get("extra_data"), ) if result.get("success"): _analytics("subscription.connected", {"provider": provider}) return result except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @subscriptions.router.post("/disconnect") async def subscriptions_disconnect(body: dict): """Disconnect a subscription provider via 9Router.""" provider = body.get("provider", "") if not provider: raise HTTPException(status_code=400, detail="provider required") try: providers_data = await get_providers() connections = providers_data.get("connections", []) if isinstance(providers_data, dict) else [] conn = next((c for c in connections if c.get("provider") == provider), None) if conn and conn.get("id"): async with httpx.AsyncClient(timeout=10.0) as client: await client.delete(f"{NINE_ROUTER_API}/providers/{conn['id']}") _analytics("subscription.disconnected", {"provider": provider}) return {"ok": True} return {"ok": False, "error": "Connection not found"} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @subscriptions.router.get("/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": "*"}) @subscriptions.router.get("/callback") async def subscriptions_callback(request: Request): """Catch OAuth redirect from provider, exchange code via 9Router, close window.""" code = request.query_params.get("code", "") state = request.query_params.get("state", "") error = request.query_params.get("error", "") if error: desc = request.query_params.get("error_description", error) return HTMLResponse(f'
{desc}
Please try connecting again.
{e}
You can close this window
' '