from backend.config.Apps import SubApp from backend.apps.agents.manager.agent_manager import agent_manager from backend.apps.agents.models import AgentConfig, ApprovalResponse from backend.apps.agents.browser.runner import run_browser_agents from contextlib import asynccontextmanager from fastapi import HTTPException, Request from fastapi.responses import JSONResponse import logging logger = logging.getLogger(__name__) @asynccontextmanager async def agents_lifespan(): logger.info("Agents sub-app starting") await agent_manager.reconcile_on_startup() await agent_manager.restore_all_sessions() yield logger.info("Agents sub-app shutting down") for session_id in list(agent_manager.tasks.keys()): await agent_manager.stop_agent(session_id) await agent_manager.persist_all_sessions() agents = SubApp("agents", agents_lifespan) # REST Endpoints @agents.router.get("/sessions") async def list_sessions(dashboard_id: str = ""): sessions = agent_manager.get_all_sessions(dashboard_id=dashboard_id or None) return {"sessions": [s.model_dump(mode="json") for s in sessions]} @agents.router.get("/sessions/{session_id}") async def get_session(session_id: str): session = agent_manager.get_session(session_id) if not session: raise HTTPException(status_code=404, detail="Session not found") return session.model_dump(mode="json") @agents.router.post("/launch") async def launch_agent(config: AgentConfig): session = await agent_manager.launch_agent(config) return {"session_id": session.id, "session": session.model_dump(mode="json")} @agents.router.post("/sessions/{session_id}/message") async def send_message(session_id: str, body: dict): prompt = body.get("prompt", "") if not prompt: raise HTTPException(status_code=400, detail="prompt is required") await agent_manager.send_message( session_id, prompt, mode=body.get("mode"), model=body.get("model"), images=body.get("images"), context_paths=body.get("context_paths"), forced_tools=body.get("forced_tools"), attached_skills=body.get("attached_skills"), hidden=body.get("hidden", False), selected_browser_ids=body.get("selected_browser_ids"), ) return {"ok": True} @agents.router.post("/sessions/{session_id}/stop") async def stop_agent(session_id: str): await agent_manager.stop_agent(session_id) return {"ok": True} @agents.router.post("/approval") async def handle_approval(response: ApprovalResponse): agent_manager.handle_approval(response.request_id, { "behavior": response.behavior, "message": response.message, "updated_input": response.updated_input, }) return {"ok": True} @agents.router.post("/sessions/{session_id}/edit_message") async def edit_message(session_id: str, body: dict): message_id = body.get("message_id") new_content = body.get("content", "") if not message_id or not new_content: raise HTTPException(status_code=400, detail="message_id and content are required") await agent_manager.edit_message(session_id, message_id, new_content) return {"ok": True} @agents.router.post("/sessions/{session_id}/switch_branch") async def switch_branch(session_id: str, body: dict): branch_id = body.get("branch_id", "") if not branch_id: raise HTTPException(status_code=400, detail="branch_id is required") await agent_manager.switch_branch(session_id, branch_id) return {"ok": True} @agents.router.post("/sessions/{session_id}/generate-title") async def generate_title(session_id: str, body: dict): prompt = body.get("prompt", "") if not prompt: raise HTTPException(status_code=400, detail="prompt is required") title = await agent_manager.generate_title(session_id, prompt) return {"title": title} @agents.router.post("/sessions/{session_id}/generate-group-meta") async def generate_group_meta(session_id: str, body: dict): group_id = body.get("group_id", "") tool_calls = body.get("tool_calls", []) if not group_id or not tool_calls: raise HTTPException(status_code=400, detail="group_id and tool_calls are required") result = await agent_manager.generate_group_meta( session_id, group_id, tool_calls, results_summary=body.get("results_summary"), is_refinement=body.get("is_refinement", False), ) return result @agents.router.patch("/sessions/{session_id}") async def update_session(session_id: str, body: dict): session = agent_manager.get_session(session_id) if not session: raise HTTPException(status_code=404, detail="Session not found") await agent_manager.update_session(session_id, **body) return {"ok": True} @agents.router.post("/sessions/{session_id}/duplicate") async def duplicate_session(session_id: str, body: dict = {}): try: session = await agent_manager.duplicate_session( session_id, dashboard_id=body.get("dashboard_id"), up_to_message_id=body.get("up_to_message_id"), ) except ValueError as e: raise HTTPException(status_code=404, detail=str(e)) return {"session": session.model_dump(mode="json")} @agents.router.post("/sessions/{session_id}/close") async def close_session(session_id: str): try: await agent_manager.close_session(session_id) except ValueError as e: raise HTTPException(status_code=404, detail=str(e)) return {"ok": True} @agents.router.delete("/sessions/{session_id}") async def delete_session(session_id: str): await agent_manager.delete_session(session_id) return {"ok": True} @agents.router.get("/history") async def get_history(q: str = "", limit: int = 20, offset: int = 0, dashboard_id: str = ""): return agent_manager.get_history( q=q, limit=limit, offset=offset, dashboard_id=dashboard_id or None, ) @agents.router.get("/sessions/{session_id}/browser-agents") async def get_browser_agent_children(session_id: str): children = agent_manager.get_browser_agent_children(session_id) return {"sessions": children} @agents.router.post("/sessions/{session_id}/resume") async def resume_session(session_id: str): try: session = await agent_manager.resume_session(session_id) except ValueError as e: raise HTTPException(status_code=404, detail=str(e)) return {"session": session.model_dump(mode="json")} @agents.router.post("/browser-agent/run") async def browser_agent_run(request: Request): """Run one or more browser sub-agents in parallel.""" body = await request.json() tasks = body.get("tasks", []) if not tasks: return JSONResponse({"error": "tasks array is required"}, status_code=400) results = await run_browser_agents( tasks=tasks, model=body.get("model", "sonnet"), dashboard_id=body.get("dashboard_id", "") or None, pre_selected_browser_ids=body.get("pre_selected_browser_ids", []), parent_session_id=body.get("parent_session_id", "") or None, ) return JSONResponse({"results": results}) @agents.router.post("/invoke-agent/run") async def invoke_agent_run(request: Request): """Fork an existing agent session and send it a new message.""" body = await request.json() session_id = body.get("session_id", "") message = body.get("message", "") if not session_id: return JSONResponse({"error": "session_id is required"}, status_code=400) if not message: return JSONResponse({"error": "message is required"}, status_code=400) try: result = await agent_manager.invoke_agent( source_session_id=session_id, message=message, parent_session_id=body.get("parent_session_id", "") or None, dashboard_id=body.get("dashboard_id", "") or None, ) return JSONResponse(result) except ValueError as e: return JSONResponse({"error": str(e)}, status_code=404) except Exception as e: return JSONResponse({"error": str(e)}, status_code=500)