diff --git a/backend/apps/agents/COMMS_MANAGER/COMMS_MANAGER.py b/backend/apps/agents/COMMS_MANAGER/COMMS_MANAGER.py new file mode 100644 index 00000000..fd314c83 --- /dev/null +++ b/backend/apps/agents/COMMS_MANAGER/COMMS_MANAGER.py @@ -0,0 +1,81 @@ +from typing import Optional, Literal, Dict, Any +from uuid import uuid4 +from pydantic import BaseModel, Field +from typeguard import typechecked + +from backend.apps.agents.utils.comms_utils.classes.FutureBridge import FutureBridge +from backend.apps.agents.utils.comms_utils.classes.FrontendBroadcaster import FrontendBroadcaster +from backend.core.events.events import AnyEvent, ApprovalRequestEvent, EventCallback + + +class CommsManager(BaseModel): + broadcaster: FrontendBroadcaster = Field(default_factory=FrontendBroadcaster) + approval_bridge: FutureBridge = Field(default_factory=FutureBridge) + browser_bridge: FutureBridge = Field(default_factory=FutureBridge) + + @typechecked + async def resolve_approval( + self, + request_id: str, + behavior: Literal["allow", "deny"], + message: Optional[str] = None, + updated_input: Optional[Dict[str, Any]] = None, + ) -> None: + self.approval_bridge.resolve(request_id, { + "behavior": behavior, + "message": message, + "updated_input": updated_input, + }) + + @typechecked + def make_session_emitter(self, session_id: str) -> EventCallback: + """Create an event callback that routes typed events to the WS pool. + + ApprovalRequestEvents are special-cased: instead of just broadcasting, + the emitter routes through the approval bridge and resolves the + embedded future with the user's decision. + """ + async def emit(event: AnyEvent) -> None: + if isinstance(event, ApprovalRequestEvent): + if not self.broadcaster.has_connections(): + if not event.future.done(): + event.future.set_result({"behavior": "deny", "message": "No dashboard connected for approval."}) + return + result = await self.approval_bridge.request( + request_id=event.request_id, + send_fn=lambda: self.broadcaster.send_to_session(session_id, event.event, { + "request_id": event.request_id, + "session_id": event.session_id, + "tool_name": event.tool_name, + "tool_input": event.tool_input, + }), + timeout=600.0, + ) + if not event.future.done(): + event.future.set_result(result) + return + await self.broadcaster.send_to_session(session_id, event.event, event.model_dump(mode="json")) + return emit + + @typechecked + async def send_browser_command( + self, action: str, browser_id: str, tab_id: str, params: dict, + ) -> dict: + """BrowserCommandFn-compatible method that routes through the browser FutureBridge.""" + request_id: str = uuid4().hex + if not self.broadcaster.has_connections(): + return {"error": "No dashboard connected. Open the dashboard to use browser tools."} + return await self.browser_bridge.request( + request_id=request_id, + send_fn=lambda: self.broadcaster.broadcast("browser:command", { + "request_id": request_id, + "action": action, + "browser_id": browser_id, + "tab_id": tab_id, + "params": params, + }), + timeout=30.0, + ) + + +COMMS_MANAGER = CommsManager() diff --git a/backend/apps/agents/utils/__init__.py b/backend/apps/agents/COMMS_MANAGER/__init__.py similarity index 100% rename from backend/apps/agents/utils/__init__.py rename to backend/apps/agents/COMMS_MANAGER/__init__.py diff --git a/backend/apps/agents/utils/comms_utils/singeltons/classes/FrontendBroadcaster.py b/backend/apps/agents/COMMS_MANAGER/classes/FrontendBroadcaster.py similarity index 100% rename from backend/apps/agents/utils/comms_utils/singeltons/classes/FrontendBroadcaster.py rename to backend/apps/agents/COMMS_MANAGER/classes/FrontendBroadcaster.py diff --git a/backend/apps/agents/utils/comms_utils/singeltons/classes/FutureBridge.py b/backend/apps/agents/COMMS_MANAGER/classes/FutureBridge.py similarity index 100% rename from backend/apps/agents/utils/comms_utils/singeltons/classes/FutureBridge.py rename to backend/apps/agents/COMMS_MANAGER/classes/FutureBridge.py diff --git a/backend/apps/agents/utils/agent_utils/__init__.py b/backend/apps/agents/agent_utils/__init__.py similarity index 100% rename from backend/apps/agents/utils/agent_utils/__init__.py rename to backend/apps/agents/agent_utils/__init__.py diff --git a/backend/apps/agents/utils/comms_utils/build_search_text.py b/backend/apps/agents/agent_utils/build_search_text.py similarity index 100% rename from backend/apps/agents/utils/comms_utils/build_search_text.py rename to backend/apps/agents/agent_utils/build_search_text.py diff --git a/backend/apps/agents/utils/agent_utils/compose_system_prompt.py b/backend/apps/agents/agent_utils/compose_system_prompt.py similarity index 100% rename from backend/apps/agents/utils/agent_utils/compose_system_prompt.py rename to backend/apps/agents/agent_utils/compose_system_prompt.py diff --git a/backend/apps/agents/utils/agent_utils/create_sdk_hooks.py b/backend/apps/agents/agent_utils/create_sdk_hooks.py similarity index 100% rename from backend/apps/agents/utils/agent_utils/create_sdk_hooks.py rename to backend/apps/agents/agent_utils/create_sdk_hooks.py diff --git a/backend/apps/agents/agents.py b/backend/apps/agents/agents.py index 2844a84b..d375899c 100644 --- a/backend/apps/agents/agents.py +++ b/backend/apps/agents/agents.py @@ -23,13 +23,11 @@ from backend.core.Agent.Agent import Agent from backend.core.db.PydanticStore import PydanticStore from backend.core.shared_structs.agent.Message.Message import UserMessage from backend.core.events.events import AgentStatusEvent, AgentClosedEvent, BranchSwitchedEvent -from backend.apps.agents.utils.comms_utils.resolve_approvals import resolve_approval -from backend.apps.agents.utils.agent_utils.compose_system_prompt import compose_system_prompt +from backend.apps.agents.agent_utils.compose_system_prompt import compose_system_prompt from backend.core.tools.make_builtin_toolkit.make_builtin_toolkit import make_builtin_toolkit -from backend.apps.agents.utils.agent_utils.create_sdk_hooks import create_sdk_hooks -from backend.apps.agents.utils.comms_utils.make_session_emitter import make_session_emitter -from backend.apps.agents.utils.comms_utils.send_browser_command import send_browser_command -from backend.apps.agents.utils.comms_utils.build_search_text import build_search_text +from backend.apps.agents.agent_utils.create_sdk_hooks import create_sdk_hooks +from backend.apps.agents.agent_utils.build_search_text import build_search_text +from backend.apps.agents.COMMS_MANAGER.COMMS_MANAGER import COMMS_MANAGER from claude_agent_sdk import ClaudeAgentOptions from claude_agent_sdk.types import HookMatcher, McpServerConfig from backend.core.tools.shared_structs.Toolkit import Toolkit @@ -59,8 +57,8 @@ async def agents_lifespan(): for stored in AGENT_STORE.load_all(): try: stored.status = "stopped" - stored.on_event = make_session_emitter(stored.session_id) - toolkit: Toolkit = make_builtin_toolkit(stored, SESSIONS, send_browser_command) + stored.on_event = COMMS_MANAGER.make_session_emitter(stored.session_id) + toolkit: Toolkit = make_builtin_toolkit(stored, SESSIONS, COMMS_MANAGER.send_browser_command) stored.toolkit = toolkit SESSIONS[stored.session_id] = stored except Exception as e: @@ -109,10 +107,10 @@ async def launch(body: LaunchBody) -> dict: status="stopped", config=ClaudeAgentOptions(max_turns=body.max_turns), ) - agent.on_event = make_session_emitter(agent.session_id) + agent.on_event = COMMS_MANAGER.make_session_emitter(agent.session_id) SESSIONS[agent.session_id] = agent - toolkit: Toolkit = make_builtin_toolkit(agent, SESSIONS, send_browser_command) + toolkit: Toolkit = make_builtin_toolkit(agent, SESSIONS, COMMS_MANAGER.send_browser_command) agent.toolkit = toolkit mcp_servers: Dict[str, McpServerConfig] = toolkit.collect_mcp_servers() allowed_tools, disallowed_tools = toolkit.collect_tool_permissions() @@ -216,7 +214,7 @@ class ApprovalBody(BaseModel): @agents.router.post("/approval") async def handle_approval(body: ApprovalBody) -> dict: - await resolve_approval( + await COMMS_MANAGER.resolve_approval( request_id=body.request_id, behavior=body.behavior, message=body.message, @@ -286,8 +284,8 @@ async def resume_session(session_id: str) -> dict: if not agent: raise HTTPException(status_code=404, detail="Session not found in history") agent.status = "stopped" - agent.on_event = make_session_emitter(agent.session_id) - toolkit: Toolkit = make_builtin_toolkit(agent, SESSIONS, send_browser_command) + agent.on_event = COMMS_MANAGER.make_session_emitter(agent.session_id) + toolkit: Toolkit = make_builtin_toolkit(agent, SESSIONS, COMMS_MANAGER.send_browser_command) agent.toolkit = toolkit SESSIONS[agent.session_id] = agent AGENT_STORE.delete(session_id) @@ -313,8 +311,8 @@ async def duplicate_session(session_id: str, body: dict = {}) -> dict: clone.lock = asyncio.Lock() clone.pending_approvals = [] clone.sub_agents = [] - clone.on_event = make_session_emitter(clone.session_id) - toolkit: Toolkit = make_builtin_toolkit(clone, SESSIONS, send_browser_command) + clone.on_event = COMMS_MANAGER.make_session_emitter(clone.session_id) + toolkit: Toolkit = make_builtin_toolkit(clone, SESSIONS, COMMS_MANAGER.send_browser_command) clone.toolkit = toolkit SESSIONS[clone.session_id] = clone await clone.emit(AgentStatusEvent( diff --git a/backend/apps/agents/utils/comms_utils/__init__.py b/backend/apps/agents/utils/comms_utils/__init__.py deleted file mode 100644 index e69de29b..00000000 diff --git a/backend/apps/agents/utils/comms_utils/make_session_emitter.py b/backend/apps/agents/utils/comms_utils/make_session_emitter.py deleted file mode 100644 index 09c97bbf..00000000 --- a/backend/apps/agents/utils/comms_utils/make_session_emitter.py +++ /dev/null @@ -1,33 +0,0 @@ -from typeguard import typechecked -from backend.core.events.events import AnyEvent, ApprovalRequestEvent, EventCallback -from backend.apps.agents.utils.comms_utils.singeltons.singeltons import APPROVAL_BRIDGE, FRONTEND_BROADCASTER - -@typechecked -def make_session_emitter(session_id: str) -> EventCallback: - """Create an event callback that routes typed events to the WS connection pool. - - ApprovalRequestEvents are special-cased: instead of just broadcasting, - the emitter routes through the APPROVAL_BRIDGE and resolves the - embedded future with the user's decision. - """ - async def emit(event: AnyEvent) -> None: - if isinstance(event, ApprovalRequestEvent): - if not FRONTEND_BROADCASTER.has_connections(): - if not event.future.done(): - event.future.set_result({"behavior": "deny", "message": "No dashboard connected for approval."}) - return - result = await APPROVAL_BRIDGE.request( - request_id=event.request_id, - send_fn=lambda: FRONTEND_BROADCASTER.send_to_session(session_id, event.event, { - "request_id": event.request_id, - "session_id": event.session_id, - "tool_name": event.tool_name, - "tool_input": event.tool_input, - }), - timeout=600.0, - ) - if not event.future.done(): - event.future.set_result(result) - return - await FRONTEND_BROADCASTER.send_to_session(session_id, event.event, event.model_dump(mode="json")) - return emit \ No newline at end of file diff --git a/backend/apps/agents/utils/comms_utils/resolve_approvals.py b/backend/apps/agents/utils/comms_utils/resolve_approvals.py deleted file mode 100644 index fceec64e..00000000 --- a/backend/apps/agents/utils/comms_utils/resolve_approvals.py +++ /dev/null @@ -1,17 +0,0 @@ -from typing import Optional, Literal, Dict, Any -from typeguard import typechecked -from backend.apps.agents.utils.comms_utils.singeltons.singeltons import APPROVAL_BRIDGE - -# TODO: add better type specing for the dict values in the message arg -@typechecked -async def resolve_approval( - request_id: str, - behavior: Literal["allow", "deny"], - message: Optional[str] = None, - updated_input: Optional[Dict[str, Any]] = None, -) -> None: - APPROVAL_BRIDGE.resolve(request_id, { - "behavior": behavior, - "message": message, - "updated_input": updated_input, - }) diff --git a/backend/apps/agents/utils/comms_utils/send_browser_command.py b/backend/apps/agents/utils/comms_utils/send_browser_command.py deleted file mode 100644 index 2085d199..00000000 --- a/backend/apps/agents/utils/comms_utils/send_browser_command.py +++ /dev/null @@ -1,24 +0,0 @@ -from typeguard import typechecked -from backend.apps.agents.utils.comms_utils.singeltons.singeltons import BROWSER_BRIDGE, FRONTEND_BROADCASTER -from uuid import uuid4 - -# TODO: add better type specing for the output of this function -@typechecked -async def send_browser_command( - action: str, browser_id: str, tab_id: str, params: dict, -) -> dict: - """BrowserCommandFn implementation that routes through the browser FutureBridge.""" - request_id: str = uuid4().hex - if not FRONTEND_BROADCASTER.has_connections(): - return {"error": "No dashboard connected. Open the dashboard to use browser tools."} - return await BROWSER_BRIDGE.request( - request_id=request_id, - send_fn=lambda: FRONTEND_BROADCASTER.broadcast("browser:command", { - "request_id": request_id, - "action": action, - "browser_id": browser_id, - "tab_id": tab_id, - "params": params, - }), - timeout=30.0, - ) \ No newline at end of file diff --git a/backend/apps/agents/utils/comms_utils/singeltons/__init__.py b/backend/apps/agents/utils/comms_utils/singeltons/__init__.py deleted file mode 100644 index e69de29b..00000000 diff --git a/backend/apps/agents/utils/comms_utils/singeltons/singeltons.py b/backend/apps/agents/utils/comms_utils/singeltons/singeltons.py deleted file mode 100644 index 6e16234c..00000000 --- a/backend/apps/agents/utils/comms_utils/singeltons/singeltons.py +++ /dev/null @@ -1,6 +0,0 @@ -from backend.apps.agents.utils.comms_utils.singeltons.classes.FutureBridge import FutureBridge -from backend.apps.agents.utils.comms_utils.singeltons.classes.FrontendBroadcaster import FrontendBroadcaster - -APPROVAL_BRIDGE = FutureBridge() -BROWSER_BRIDGE = FutureBridge() -FRONTEND_BROADCASTER = FrontendBroadcaster()