mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-13 13:17:40 +02:00
[Haik]: Made a COMMS_MANAGER singelton which handles all high level frontend comm concerns used by the agents subapp
This commit is contained in:
@@ -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()
|
||||
@@ -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(
|
||||
|
||||
@@ -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
|
||||
@@ -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,
|
||||
})
|
||||
@@ -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,
|
||||
)
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user