diff --git a/backend/apps/agents/agent_manager.py b/backend/apps/agents/agent_manager.py index 1f822bcf..3c06f3f2 100644 --- a/backend/apps/agents/agent_manager.py +++ b/backend/apps/agents/agent_manager.py @@ -12,7 +12,7 @@ from typing import Optional from backend.apps.agents.core.models import ( AgentConfig, AgentSession, Message, MessageBranch, ApprovalRequest, ToolGroupMeta, ) -from backend.apps.agents.core.ws_manager import ws_manager +from backend.apps.agents.core.ws_manager import WS_MANAGER from backend.apps.settings.store import load_settings from backend.apps.tools_lib.oauth_tokens import ( refresh_google_token, @@ -269,7 +269,7 @@ class AgentManager: # written meta.json. try: new_output = _load(output_id) - await ws_manager.broadcast_global("agent:output_upserted", { + await WS_MANAGER.broadcast_global("agent:output_upserted", { "output": new_output.model_dump(mode="json"), }) except Exception: @@ -314,7 +314,7 @@ class AgentManager: _apply_context_window(session, global_settings) self.sessions[session_id] = session - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "running", "session": session.model_dump(mode="json"), @@ -696,12 +696,12 @@ class AgentManager: session.status = "waiting_approval" - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "waiting_approval", }) - decision = await ws_manager.send_approval_request( + decision = await WS_MANAGER.send_approval_request( session_id, request_id, tool_name, safe_input, sensitive_pattern=sensitive_pattern, sensitive_label=label, @@ -740,7 +740,7 @@ class AgentManager: a for a in session.pending_approvals if a.id != request_id ] session.status = "running" - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "running", }) @@ -981,7 +981,7 @@ class AgentManager: ) _apply_context_window(sub_session) self.sessions[sub_session_id] = sub_session - await ws_manager.broadcast_global("agent:status", { + await WS_MANAGER.broadcast_global("agent:status", { "session_id": sub_session_id, "status": sub_session.status, "session": sub_session.model_dump(mode="json"), @@ -1005,7 +1005,7 @@ class AgentManager: except Exception: logger.exception("Tool result truncation failed; keeping inline body") session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) @@ -1046,7 +1046,7 @@ class AgentManager: if _stale: session.active_mcps = [s for s in session.active_mcps if s in _enabled] session.needs_fork = True - await ws_manager.send_to_session(session_id, "agent:context_status", { + await WS_MANAGER.send_to_session(session_id, "agent:context_status", { "session_id": session_id, "reason": "mcp_disabled_externally", "deactivated": _stale, @@ -1801,7 +1801,7 @@ class AgentManager: # zero latency on the user's turn. try: if self._maybe_compact(session): - await ws_manager.send_to_session(session_id, "agent:context_status", { + await WS_MANAGER.send_to_session(session_id, "agent:context_status", { "session_id": session_id, "reason": "compacted", "compacted_through_msg_id": session.compacted_through_msg_id, @@ -1830,7 +1830,7 @@ class AgentManager: trimmed.append(f"mcp:{session.active_mcps.pop(0)}") _est_tokens -= 8_000 # rough per-MCP schema cost if trimmed: - await ws_manager.send_to_session(session_id, "agent:context_status", { + await WS_MANAGER.send_to_session(session_id, "agent:context_status", { "session_id": session_id, "reason": "trimmed", "trimmed": trimmed, @@ -1853,7 +1853,7 @@ class AgentManager: branch_id=session.active_branch_id, ) session.messages.append(_trim_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": _trim_msg.model_dump(mode="json"), }) @@ -2156,7 +2156,7 @@ class AgentManager: else: session.messages.append(consolidated) try: - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": consolidated.model_dump(mode="json"), }) @@ -2274,7 +2274,7 @@ class AgentManager: if block_type == "text": if stream_text_msg_id is None: stream_text_msg_id = uuid4().hex - await ws_manager.send_to_session(session_id, "agent:stream_start", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_start", { "session_id": session_id, "message_id": stream_text_msg_id, "role": "assistant", @@ -2296,7 +2296,7 @@ class AgentManager: # thinking blocks (think → tool → think # → answer turns sum correctly). _thinking_block_starts[index] = time.time() - await ws_manager.send_to_session(session_id, "agent:stream_start", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_start", { "session_id": session_id, "message_id": thinking_msg_id, "role": "thinking", @@ -2322,7 +2322,7 @@ class AgentManager: # the dedupe at the AssistantMessage # block below. _turn_tool_count += 1 - await ws_manager.send_to_session(session_id, "agent:stream_start", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_start", { "session_id": session_id, "message_id": tool_msg_id, "role": "tool_call", @@ -2338,7 +2338,7 @@ class AgentManager: if msg_id and delta_type == "text_delta": _text_chunk = delta.get("text", "") _turn_assistant_text_chars += len(_text_chunk) - await ws_manager.send_to_session(session_id, "agent:stream_delta", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_delta", { "session_id": session_id, "message_id": msg_id, "delta": _text_chunk, @@ -2348,7 +2348,7 @@ class AgentManager: # with a "thinking" field (not "text") _think_chunk = delta.get("thinking", "") _thinking_total_chars += len(_think_chunk) - await ws_manager.send_to_session(session_id, "agent:stream_delta", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_delta", { "session_id": session_id, "message_id": msg_id, "delta": _think_chunk, @@ -2356,7 +2356,7 @@ class AgentManager: elif msg_id and delta_type == "input_json_delta": _json_chunk = delta.get("partial_json", "") _turn_tool_input_chars += len(_json_chunk) - await ws_manager.send_to_session(session_id, "agent:stream_delta", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_delta", { "session_id": session_id, "message_id": msg_id, "delta": _json_chunk, @@ -2376,14 +2376,14 @@ class AgentManager: (time.time() - _thinking_block_starts.pop(index)) * 1000 ) if msg_id and msg_id != stream_text_msg_id: - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": msg_id, }) elif event_type == "message_stop": if stream_text_msg_id: - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": stream_text_msg_id, }) @@ -2526,13 +2526,13 @@ class AgentManager: branch_id=session.active_branch_id, ) session.messages.append(_err_msg) - await ws_manager.send_to_session(session_id, "agent:auth_error", { + await WS_MANAGER.send_to_session(session_id, "agent:auth_error", { "session_id": session_id, "reason": reason, "message": friendly, "model": session.model, }) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": _err_msg.model_dump(mode="json"), }) @@ -2544,7 +2544,7 @@ class AgentManager: branch_id=session.active_branch_id, ) session.messages.append(asst_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": asst_msg.model_dump(mode="json"), }) @@ -2553,7 +2553,7 @@ class AgentManager: msg_id = stream_tool_msg_ids_ordered[i] if i < len(stream_tool_msg_ids_ordered) else uuid4().hex tool_msg = Message(id=msg_id, role="tool_call", content=tu, branch_id=session.active_branch_id) session.messages.append(tool_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": tool_msg.model_dump(mode="json"), }) @@ -2739,7 +2739,7 @@ class AgentManager: cost = 0.0 session.cost_usd = cost - await ws_manager.send_to_session(session_id, "agent:cost_update", { + await WS_MANAGER.send_to_session(session_id, "agent:cost_update", { "session_id": session_id, "cost_usd": session.cost_usd, }) @@ -2757,7 +2757,7 @@ class AgentManager: ctx_used_pct = round(total_input / _ctx_window, 4) if total_input else 0.0 cache_read_pct = round(cache_read / total_input, 4) if total_input else 0.0 try: - await ws_manager.send_to_session(session_id, "agent:context_update", { + await WS_MANAGER.send_to_session(session_id, "agent:context_update", { "session_id": session_id, "input_tokens": total_input, "output_tokens": out, @@ -2811,13 +2811,13 @@ class AgentManager: # it with stream_end and start the fresh turn under a # new message id. if stream_text_msg_id: - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": stream_text_msg_id, }) stream_text_msg_id = None for _tool_msg_id in stream_tool_msg_ids_ordered: - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": _tool_msg_id, }) @@ -2884,7 +2884,7 @@ class AgentManager: session.status = "completed" if stream_text_msg_id: try: - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": stream_text_msg_id, }) @@ -2913,8 +2913,8 @@ class AgentManager: "compact_threshold_pct": session.compact_threshold_pct, "context_soft_cap_pct": session.context_soft_cap_pct, } - await ws_manager.send_to_session(session_id, "agent:context_overflow", _ovf_payload) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:context_overflow", _ovf_payload) + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) @@ -2950,11 +2950,11 @@ class AgentManager: ) error_msg = Message(role="system", content=friendly_msg, branch_id=session.active_branch_id) session.messages.append(error_msg) - await ws_manager.send_to_session(session_id, "agent:free_trial_exhausted", { + await WS_MANAGER.send_to_session(session_id, "agent:free_trial_exhausted", { "session_id": session_id, "message": friendly_msg, }) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) @@ -3014,13 +3014,13 @@ class AgentManager: reason = "anthropic_auth_invalid" error_msg = Message(role="system", content=friendly_msg, branch_id=session.active_branch_id) session.messages.append(error_msg) - await ws_manager.send_to_session(session_id, "agent:auth_error", { + await WS_MANAGER.send_to_session(session_id, "agent:auth_error", { "session_id": session_id, "reason": reason, "message": friendly_msg, "model": session.model, }) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) @@ -3042,7 +3042,7 @@ class AgentManager: logger.debug("submit_diagnostic model_error failed", exc_info=True) error_msg = Message(role="system", content=f"Error: {str(e)}", branch_id=session.active_branch_id) session.messages.append(error_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) @@ -3062,7 +3062,7 @@ class AgentManager: logger.debug("submit_diagnostic model_error failed", exc_info=True) error_msg = Message(role="system", content=f"Error: {str(e)}", branch_id=session.active_branch_id) session.messages.append(error_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) @@ -3074,7 +3074,7 @@ class AgentManager: session.status = "error" error_msg = Message(role="system", content=f"Error: {str(e)}", branch_id=session.active_branch_id) session.messages.append(error_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) @@ -3098,14 +3098,14 @@ class AgentManager: try: matching = [o for o in _load_all() if o.workspace_id == session_id] if matching: - await ws_manager.broadcast_global("agent:output_upserted", { + await WS_MANAGER.broadcast_global("agent:output_upserted", { "output": matching[0].model_dump(mode="json"), }) except Exception: logger.exception("post-sync output_upserted broadcast failed") except Exception: logger.exception("post-session meta sync failed") - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), @@ -3117,7 +3117,7 @@ class AgentManager: async def _stream_text(self, session_id: str, msg_id: str, text: str, delay: float = 0.03): """Emit stream_start, word-by-word deltas, and stream_end for a text message.""" - await ws_manager.send_to_session(session_id, "agent:stream_start", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_start", { "session_id": session_id, "message_id": msg_id, "role": "assistant", @@ -3125,20 +3125,20 @@ class AgentManager: words = text.split(" ") for i, word in enumerate(words): chunk = word if i == 0 else " " + word - await ws_manager.send_to_session(session_id, "agent:stream_delta", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_delta", { "session_id": session_id, "message_id": msg_id, "delta": chunk, }) await asyncio.sleep(delay) - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": msg_id, }) async def _stream_tool_input(self, session_id: str, msg_id: str, tool_name: str, input_json: str, delay: float = 0.02): """Emit stream_start, chunked deltas, and stream_end for a tool_call input.""" - await ws_manager.send_to_session(session_id, "agent:stream_start", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_start", { "session_id": session_id, "message_id": msg_id, "role": "tool_call", @@ -3146,13 +3146,13 @@ class AgentManager: }) chunk_size = 12 for i in range(0, len(input_json), chunk_size): - await ws_manager.send_to_session(session_id, "agent:stream_delta", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_delta", { "session_id": session_id, "message_id": msg_id, "delta": input_json[i:i + chunk_size], }) await asyncio.sleep(delay) - await ws_manager.send_to_session(session_id, "agent:stream_end", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_end", { "session_id": session_id, "message_id": msg_id, }) @@ -3174,19 +3174,19 @@ class AgentManager: ) session.pending_approvals.append(approval_req) session.status = "waiting_approval" - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "waiting_approval", }) - decision = await ws_manager.send_approval_request( + decision = await WS_MANAGER.send_approval_request( session_id, request_id, "Bash", {"command": f"echo 'Processing: {prompt}'", "description": "Echo the user prompt"} ) session.pending_approvals = [a for a in session.pending_approvals if a.id != request_id] session.status = "running" - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "running", }) @@ -3200,7 +3200,7 @@ class AgentManager: ) tool_msg = Message(id=tool_msg_id, role="tool_call", content=tool_input_content, branch_id=session.active_branch_id) session.messages.append(tool_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": tool_msg.model_dump(mode="json"), }) @@ -3210,7 +3210,7 @@ class AgentManager: if decision.get("behavior") == "allow": tool_result = Message(role="tool_result", content=f"Processing: {prompt}", branch_id=session.active_branch_id) session.messages.append(tool_result) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": tool_result.model_dump(mode="json"), }) @@ -3228,7 +3228,7 @@ class AgentManager: asst_msg = Message(id=asst_msg_id, role="assistant", content=asst_text, branch_id=session.active_branch_id) session.messages.append(asst_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": asst_msg.model_dump(mode="json"), }) @@ -3241,12 +3241,12 @@ class AgentManager: # `_mock_run` flag is read by the close path so a mock session # doesn't get reported to the cloud as a real one. setattr(session, "_mock_run", True) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "completed", "session": session.model_dump(mode="json"), }) - await ws_manager.send_to_session(session_id, "agent:cost_update", { + await WS_MANAGER.send_to_session(session_id, "agent:cost_update", { "session_id": session_id, "cost_usd": session.cost_usd, }) @@ -3306,7 +3306,7 @@ class AgentManager: session.allowed_tools = mode_tools session_changed = True if session_changed: - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), @@ -3326,7 +3326,7 @@ class AgentManager: client_message_id=client_message_id, ) session.messages.append(user_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": user_msg.model_dump(mode="json"), }) @@ -3360,7 +3360,7 @@ class AgentManager: pass session.status = "running" - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "running", "session": session.model_dump(mode="json"), @@ -3455,14 +3455,14 @@ class AgentManager: _tc = Message(role="tool_call", branch_id=session.active_branch_id, content={"id": _bubble_tid, "tool": _BROWSER_TOOL, "input": {"task": prompt}}) session.messages.append(_tc) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": _tc.model_dump(mode="json")}) first = await _dispatch(browser_fast_path.compose_task(prompt, brief)) text = _summary(first) if browser_fast_path.dispatch_failed(first): # Retry only transient failures; a dead dashboard fails the # retry identically, so skip it and tell the user instead. - if not ws_manager.global_connections: + if not WS_MANAGER.global_connections: _fp_path += "+no-dashboard" text = browser_fast_path.NO_DASHBOARD_REPLY else: @@ -3506,17 +3506,17 @@ class AgentManager: _tr = Message(role="tool_result", branch_id=session.active_branch_id, content={"tool_use_id": _bubble_tid, "tool": _BROWSER_TOOL, "text": "done"}) session.messages.append(_tr) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": _tr.model_dump(mode="json")}) asst_msg = Message(role="assistant", content=text, branch_id=session.active_branch_id) session.messages.append(asst_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": asst_msg.model_dump(mode="json"), }) session.status = "completed" session.closed_at = datetime.now() - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "completed", "session": session.model_dump(mode="json"), @@ -3544,13 +3544,13 @@ class AgentManager: session._cancel_event.set() for req in list(session.pending_approvals): - ws_manager.resolve_approval(req.id, {"behavior": "deny", "message": "Agent stopped"}) + WS_MANAGER.resolve_approval(req.id, {"behavior": "deny", "message": "Agent stopped"}) session.pending_approvals = [] session.status = "stopped" if not session.closed_at: session.closed_at = datetime.now() - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "stopped", "session": session.model_dump(mode="json"), @@ -3566,7 +3566,7 @@ class AgentManager: def handle_approval(self, request_id: str, decision: dict): """Resolve a pending HITL approval.""" - ws_manager.resolve_approval(request_id, decision) + WS_MANAGER.resolve_approval(request_id, decision) async def edit_message(self, session_id: str, message_id: str, new_content: str): """Edit a prior user message, creating a new branch (fork).""" @@ -3627,18 +3627,18 @@ class AgentManager: ) session.messages.append(edited_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": edited_msg.model_dump(mode="json"), }) - await ws_manager.send_to_session(session_id, "agent:branch_created", { + await WS_MANAGER.send_to_session(session_id, "agent:branch_created", { "session_id": session_id, "branch": new_branch.model_dump(mode="json"), "active_branch_id": new_branch_id, }) session.status = "running" - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "running", "session": session.model_dump(mode="json"), @@ -3662,7 +3662,7 @@ class AgentManager: raise ValueError(f"Branch {branch_id} not found") session.active_branch_id = branch_id session.needs_fresh_session = True - await ws_manager.send_to_session(session_id, "agent:branch_switched", { + await WS_MANAGER.send_to_session(session_id, "agent:branch_switched", { "session_id": session_id, "active_branch_id": branch_id, }) @@ -3720,7 +3720,7 @@ class AgentManager: logger.warning(f"Title generation failed, using fallback: {e}") session.name = title - await ws_manager.send_to_session(session_id, "agent:name_updated", { + await WS_MANAGER.send_to_session(session_id, "agent:name_updated", { "session_id": session_id, "name": title, }) @@ -3791,7 +3791,7 @@ class AgentManager: if not label: return - await ws_manager.send_to_session(session_id, "agent:turn_label", { + await WS_MANAGER.send_to_session(session_id, "agent:turn_label", { "session_id": session_id, "turn_id": turn_id, "label": label, @@ -3926,7 +3926,7 @@ class AgentManager: meta = ToolGroupMeta(id=group_id, name=name, svg=svg, is_refined=is_refinement) session.tool_group_meta[group_id] = meta - await ws_manager.send_to_session(session_id, "agent:group_meta_updated", { + await WS_MANAGER.send_to_session(session_id, "agent:group_meta_updated", { "session_id": session_id, "group_id": group_id, "name": name, @@ -3950,7 +3950,7 @@ class AgentManager: continue setattr(session, key, value) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), @@ -3990,7 +3990,7 @@ class AgentManager: session.closed_at = datetime.now() for req in list(session.pending_approvals): - ws_manager.resolve_approval(req.id, {"behavior": "deny", "message": "Session closed"}) + WS_MANAGER.resolve_approval(req.id, {"behavior": "deny", "message": "Session closed"}) session.pending_approvals = [] if hasattr(session, '_cancel_event'): @@ -4003,7 +4003,7 @@ class AgentManager: save_session(session_id, doc_data) - await ws_manager.send_to_session(session_id, "agent:closed", { + await WS_MANAGER.send_to_session(session_id, "agent:closed", { "session_id": session_id, "status": session.status, "name": session.name, @@ -4065,7 +4065,7 @@ class AgentManager: # completions and close_session calls overwrite it via # save_session, so memory and disk stay in sync. - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), @@ -4138,7 +4138,7 @@ class AgentManager: session.status = "stopped" session.closed_at = None for req in list(session.pending_approvals): - ws_manager.resolve_approval(req.id, {"behavior": "deny", "message": "Server shutting down"}) + WS_MANAGER.resolve_approval(req.id, {"behavior": "deny", "message": "Server shutting down"}) session.pending_approvals = [] # Tag this close as "shutdown" so the cloud can tell it apart # from a user-initiated close. The desktop doesn't care; the @@ -4243,7 +4243,7 @@ class AgentManager: self.sessions[new_session.id] = new_session - await ws_manager.send_to_session(new_session.id, "agent:status", { + await WS_MANAGER.send_to_session(new_session.id, "agent:status", { "session_id": new_session.id, "status": new_session.status, "session": new_session.model_dump(mode="json"), @@ -4329,7 +4329,7 @@ class AgentManager: self.sessions[fork.id] = fork - await ws_manager.broadcast_global("agent:status", { + await WS_MANAGER.broadcast_global("agent:status", { "session_id": fork.id, "status": fork.status, "session": fork.model_dump(mode="json"), @@ -4341,7 +4341,7 @@ class AgentManager: branch_id=fork.active_branch_id, ) fork.messages.append(user_msg) - await ws_manager.send_to_session(fork.id, "agent:message", { + await WS_MANAGER.send_to_session(fork.id, "agent:message", { "session_id": fork.id, "message": user_msg.model_dump(mode="json"), }) diff --git a/backend/apps/agents/agents.py b/backend/apps/agents/agents.py index 8570534d..b705d4dc 100644 --- a/backend/apps/agents/agents.py +++ b/backend/apps/agents/agents.py @@ -65,13 +65,13 @@ async def send_message(session_id: str, body: dict): # Run MCP-suggestion classifier in parallel with the agent launch; fails open. try: from backend.apps.agents.core.mcp_preflight import run_preflight - from backend.apps.agents.core.ws_manager import ws_manager as _ws + from backend.apps.agents.core.ws_manager import WS_MANAGER async def _emit_preflight(): try: result = await run_preflight(prompt) if result.get("suggestions") or result.get("is_vague"): - await _ws.send_to_session(session_id, "agent:mcp_suggestions", { + await WS_MANAGER.send_to_session(session_id, "agent:mcp_suggestions", { "session_id": session_id, "suggestions": result.get("suggestions", []), "is_vague": bool(result.get("is_vague")), @@ -297,9 +297,9 @@ async def compact_session(session_id: str): raise HTTPException(status_code=404, detail="session not found") fired = agent_manager._maybe_compact(session, force=True) if fired: - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.ws_manager import WS_MANAGER try: - await ws_manager.send_to_session(session_id, "agent:context_status", { + await WS_MANAGER.send_to_session(session_id, "agent:context_status", { "session_id": session_id, "reason": "compacted", "compacted_through_msg_id": session.compacted_through_msg_id, @@ -322,9 +322,9 @@ async def clear_session(session_id: str): session.compacted_through_msg_id = None session.tokens = {"input": 0, "output": 0} session.needs_fresh_session = True - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.ws_manager import WS_MANAGER try: - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), diff --git a/backend/apps/agents/browser/browser_agent.py b/backend/apps/agents/browser/browser_agent.py index 3e468057..1f70c3e0 100644 --- a/backend/apps/agents/browser/browser_agent.py +++ b/backend/apps/agents/browser/browser_agent.py @@ -2,7 +2,7 @@ Browser sub-agent runner. Provides a lightweight Anthropic API tool-use loop that drives browser -interactions directly through ws_manager (no MCP subprocess needed). +interactions directly through WS_MANAGER (no MCP subprocess needed). Sub-agents appear as visible AgentSession cards on the dashboard. """ @@ -83,7 +83,7 @@ from backend.apps.agents.browser.browser_schema import ( SYSTEM_PROMPT, ) from backend.apps.agents.core.models import AgentSession, ApprovalRequest, Message -from backend.apps.agents.core.ws_manager import ws_manager, _await_reconnect +from backend.apps.agents.core.ws_manager import WS_MANAGER, await_reconnect from backend.apps.tools_lib.tools_lib import load_builtin_permissions logger = logging.getLogger(__name__) @@ -99,14 +99,14 @@ _CONFIRM_TOOLS = { async def execute_browser_tool( tool_name: str, tool_input: dict, browser_id: str, tab_id: str = "", ) -> dict: - """Execute a browser tool via ws_manager directly (no MCP/HTTP round-trip).""" + """Execute a browser tool via WS_MANAGER directly (no MCP/HTTP round-trip).""" action = ACTION_MAP.get(tool_name) if not action: return {"error": f"Unknown browser tool: {tool_name}"} params = {k: v for k, v in tool_input.items()} request_id = uuid4().hex - result = await ws_manager.send_browser_command( + result = await WS_MANAGER.send_browser_command( request_id, action, browser_id, params, tab_id=tab_id, ) return result @@ -323,14 +323,14 @@ async def _request_browser_approval( session.pending_approvals.append(approval_req) session.status = "waiting_approval" - await ws_manager.send_to_session(session.id, "agent:status", { + await WS_MANAGER.send_to_session(session.id, "agent:status", { "session_id": session.id, "status": "waiting_approval", }) try: decision = await asyncio.wait_for( - ws_manager.send_approval_request( + WS_MANAGER.send_approval_request( session.id, request_id, tool_name, tool_input, ), timeout=300.0, @@ -342,7 +342,7 @@ async def _request_browser_approval( a for a in session.pending_approvals if a.id != request_id ] session.status = "running" - await ws_manager.send_to_session(session.id, "agent:status", { + await WS_MANAGER.send_to_session(session.id, "agent:status", { "session_id": session.id, "status": "running", }) @@ -390,7 +390,7 @@ async def run_browser_agent( if parent and parent.status == "stopped": cancel_event.set() - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "running", "session": session.model_dump(mode="json"), @@ -480,11 +480,11 @@ async def run_browser_agent( ) err_msg = Message(role="system", content=f"Error: {error_text}") session.messages.append(err_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": err_msg.model_dump(mode="json"), }) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "error", "session": session.model_dump(mode="json"), @@ -645,7 +645,7 @@ async def run_browser_agent( user_msg = Message(role="user", content=task) session.messages.append(user_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": user_msg.model_dump(mode="json"), }) @@ -658,12 +658,12 @@ async def run_browser_agent( _recall_msg = Message(role="assistant", content=f"Picking up what I learned about {_pb_host} from a previous visit.") session.messages.append(_recall_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": _recall_msg.model_dump(mode="json"), }) # Push the session so the "Remembered" chip shows WHILE it works (the # high-value moment), not just on the finished card. - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), }) @@ -855,7 +855,7 @@ async def run_browser_agent( pass session.status = "completed" agent_manager._sync_session_close(session) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "completed", "session": session.model_dump(mode="json"), }) @@ -1060,7 +1060,7 @@ async def run_browser_agent( content="\n".join(text_parts), ) session.messages.append(asst_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": asst_msg.model_dump(mode="json"), }) @@ -1071,7 +1071,7 @@ async def run_browser_agent( content={"id": tu.id, "tool": tu.name, "input": tu.input}, ) session.messages.append(tool_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": tool_msg.model_dump(mode="json"), }) @@ -1233,7 +1233,7 @@ async def run_browser_agent( ) brain_msg = Message(role="assistant", content=brain_text) session.messages.append(brain_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": brain_msg.model_dump(mode="json"), }) @@ -1273,7 +1273,7 @@ async def run_browser_agent( tool_results.append({"type": "tool_result", "tool_use_id": tu.id, "content": [{"type": "text", "text": meta_text}]}) result_msg = Message(role="tool_result", content={"text": meta_text, "tool_name": tu.name, "elapsed_ms": 0}) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) continue @@ -1321,7 +1321,7 @@ async def run_browser_agent( tool_results.append({"type": "tool_result", "tool_use_id": tu.id, "content": [{"type": "text", "text": ex_text}]}) result_msg = Message(role="tool_result", content={"text": ex_text, "tool_name": tu.name, "elapsed_ms": 0}) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) continue @@ -1358,7 +1358,7 @@ async def run_browser_agent( tool_results.append({"type": "tool_result", "tool_use_id": tu.id, "content": [{"type": "text", "text": sv_text}]}) result_msg = Message(role="tool_result", content={"text": sv_text, "tool_name": tu.name, "elapsed_ms": int((time.time() - st) * 1000)}) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) continue @@ -1419,7 +1419,7 @@ async def run_browser_agent( tool_results.append({"type": "tool_result", "tool_use_id": tu.id, "content": [{"type": "text", "text": bf_text}]}) result_msg = Message(role="tool_result", content={"text": bf_text, "tool_name": tu.name, "elapsed_ms": 0}) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) continue @@ -1464,7 +1464,7 @@ async def run_browser_agent( content={"text": result_text, "tool_name": tu.name, "elapsed_ms": 0}, ) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) @@ -1484,7 +1484,7 @@ async def run_browser_agent( content={"text": denied_text, "tool_name": tu.name, "elapsed_ms": 0}, ) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) @@ -1506,7 +1506,7 @@ async def run_browser_agent( content={"text": denied_text, "tool_name": tu.name, "elapsed_ms": 0}, ) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) @@ -1911,7 +1911,7 @@ async def run_browser_agent( content={"text": result_text, "tool_name": tu.name, "elapsed_ms": elapsed_ms}, ) session.messages.append(result_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": result_msg.model_dump(mode="json"), }) @@ -1969,7 +1969,7 @@ async def run_browser_agent( session.status = "stopped" record_task(session_id, browser_id, task, "stopped", metrics_started_at, turn + 1, action_log, session.tokens) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "stopped", "session": session.model_dump(mode="json"), @@ -2121,13 +2121,13 @@ async def run_browser_agent( _learn_msg = Message(role="assistant", content=f"Noted what worked on {pb_host} so I'm faster here next time.") session.messages.append(_learn_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": _learn_msg.model_dump(mode="json"), }) except Exception as e: logger.debug(f"[browser-playbook] distill skipped: {e}") agent_manager._sync_session_close(session) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": final_status, "session": session.model_dump(mode="json"), @@ -2158,11 +2158,11 @@ async def run_browser_agent( action_log, session.tokens) error_msg = Message(role="system", content=f"Error: {str(e)}") session.messages.append(error_msg) - await ws_manager.send_to_session(session_id, "agent:message", { + await WS_MANAGER.send_to_session(session_id, "agent:message", { "session_id": session_id, "message": error_msg.model_dump(mode="json"), }) - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": "error", "session": session.model_dump(mode="json"), @@ -2240,7 +2240,7 @@ async def _create_browser_card(dashboard_id: str, url: str, parent_session_id: s dashboard.updated_at = datetime.now() save_dashboard(dashboard) - await ws_manager.broadcast_global("dashboard:browser_card_added", { + await WS_MANAGER.broadcast_global("dashboard:browser_card_added", { "dashboard_id": dashboard_id, "browser_card": card.model_dump(mode="json"), "parent_session_id": parent_session_id or "", @@ -2267,7 +2267,7 @@ async def run_browser_agents( # corpse before card-gone detection trips. But a CPU-starved renderer can # briefly drop its WS then auto-reconnect, so wait (capped) for it to come # back before refusing, turning a load blip into a pause, not a failed run. - if not ws_manager.global_connections and not await _await_reconnect(lambda: bool(ws_manager.global_connections)): + if not WS_MANAGER.global_connections and not await await_reconnect(lambda: bool(WS_MANAGER.global_connections)): logger.warning("[browser-agent] dispatch refused: no dashboard after reconnect wait") return [{ "summary": ( diff --git a/backend/apps/agents/core/ws_manager.py b/backend/apps/agents/core/ws_manager.py index 4abd8f85..d9f4f98d 100644 --- a/backend/apps/agents/core/ws_manager.py +++ b/backend/apps/agents/core/ws_manager.py @@ -26,8 +26,8 @@ _BROWSER_CMD_REBROADCAST_S = 3.0 # reconnect even on a loaded machine. _WS_RECONNECT_WAIT_S = 8.0 - -async def _await_reconnect(has_conn) -> bool: +# Public - called by browser_agent.py +async def await_reconnect(has_conn) -> bool: """Poll up to _WS_RECONNECT_WAIT_S for a dashboard socket to (re)appear. `has_conn` is a 0-arg callable returning truthy when connected.""" if has_conn(): @@ -260,7 +260,7 @@ class ConnectionManager: self, request_id: str, action: str, browser_id: str, params: dict, tab_id: str = "" ) -> dict: """Send a browser command to the frontend and wait for the result.""" - if not self.global_connections and not await _await_reconnect(lambda: bool(self.global_connections)): + if not self.global_connections and not await await_reconnect(lambda: bool(self.global_connections)): return {"error": "No dashboard is connected. Open the dashboard to use browser tools."} loop = asyncio.get_event_loop() @@ -308,4 +308,4 @@ class ConnectionManager: future.set_result(result) -ws_manager = ConnectionManager() +WS_MANAGER = ConnectionManager() diff --git a/backend/main.py b/backend/main.py index f9bf9cff..2840f755 100644 --- a/backend/main.py +++ b/backend/main.py @@ -29,7 +29,7 @@ from backend.apps.oauth_state import ( from backend.config.Apps import MainApp from backend.apps.health.health import health from backend.apps.agents.agents import agents -from backend.apps.agents.core.ws_manager import ws_manager +from backend.apps.agents.core.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 @@ -186,7 +186,7 @@ async def websocket_session(websocket: WebSocket, session_id: str): """ if not p_ws_auth_ok(websocket): return - await ws_manager.connect_session(session_id, websocket) + await WS_MANAGER.connect_session(session_id, websocket) try: while True: data = await websocket.receive_text() @@ -203,7 +203,7 @@ async def websocket_session(websocket: WebSocket, session_id: str): # terminal event for already-finished sessions. last_seq = int(payload.get("last_seq") or 0) connection_uuid = payload.get("connection_uuid") or "" - ack = await ws_manager.replay_to(session_id, websocket, last_seq) + ack = await WS_MANAGER.replay_to(session_id, websocket, last_seq) from backend.apps.agents.core.seq_log import SEQ_LOG await websocket.send_text(json.dumps({ "event": "server:hello", @@ -255,7 +255,7 @@ async def websocket_session(websocket: WebSocket, session_id: str): except WebSocketDisconnect: # Drops the socket from the connection list. Does NOT cancel # the agent task, that's intentional. See module docstring. - ws_manager.disconnect_session(session_id, websocket) + WS_MANAGER.disconnect_session(session_id, websocket) def p_ws_auth_ok(websocket: WebSocket) -> bool: """Validate token + origin before accepting a WS. Returns True if OK. @@ -371,7 +371,7 @@ async def websocket_runtime_logs(websocket: WebSocket, workspace_id: str): async def websocket_dashboard(websocket: WebSocket): if not p_ws_auth_ok(websocket): return - await ws_manager.connect_global(websocket) + await WS_MANAGER.connect_global(websocket) try: while True: data = await websocket.receive_text() @@ -394,12 +394,12 @@ async def websocket_dashboard(websocket: WebSocket): "trust_pattern": bool(payload.get("trust_pattern")), }) elif event == "browser:result": - ws_manager.resolve_browser_command( + WS_MANAGER.resolve_browser_command( payload.get("request_id", ""), payload, ) except WebSocketDisconnect: - ws_manager.disconnect_global(websocket) + WS_MANAGER.disconnect_global(websocket) @app.get("/api/dev/token") @@ -427,7 +427,7 @@ async def browser_command(request: Request): 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) + result = await WS_MANAGER.send_browser_command(request_id, action, browser_id, params, tab_id=tab_id) return JSONResponse(result) @@ -680,8 +680,8 @@ async def mcp_meta(action: str, request: Request): if session.sdk_session_id: session.needs_fresh_session = True try: - from backend.apps.agents.core.ws_manager import ws_manager - await ws_manager.send_to_session(parent_session_id, "agent:status", { + from backend.apps.agents.core.WS_MANAGER import WS_MANAGER + await WS_MANAGER.send_to_session(parent_session_id, "agent:status", { "session_id": parent_session_id, "status": session.status, "session": session.model_dump(mode="json"), @@ -748,14 +748,14 @@ async def session_compact(session_id: str): only sets the marker; the button is the user opting into the cost). """ from backend.apps.agents.agent_manager import agent_manager - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.WS_MANAGER import WS_MANAGER session = agent_manager.sessions.get(session_id) if not session: return JSONResponse({"error": "session not found"}, status_code=404) did_compact = agent_manager._maybe_compact(session, force=True) if did_compact: session.needs_fresh_session = True - await ws_manager.send_to_session(session_id, "agent:context_status", { + await WS_MANAGER.send_to_session(session_id, "agent:context_status", { "session_id": session_id, "reason": "compacted_manual" if did_compact else "noop", "compacted_through_msg_id": session.compacted_through_msg_id, @@ -767,7 +767,7 @@ async def session_compact(session_id: str): async def session_clear(session_id: str): """Wipe the session's UI history AND its SDK convo state (/clear slash cmd, Reset history button).""" from backend.apps.agents.agent_manager import agent_manager - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.WS_MANAGER import WS_MANAGER from backend.apps.agents.core.models import MessageBranch session = agent_manager.sessions.get(session_id) if not session: @@ -783,12 +783,12 @@ async def session_clear(session_id: str): session.branches = {"main": MessageBranch(id="main")} session.active_branch_id = "main" session.tool_group_meta = {} - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": session.status, "session": session.model_dump(mode="json"), }) - await ws_manager.send_to_session(session_id, "agent:context_status", { + await WS_MANAGER.send_to_session(session_id, "agent:context_status", { "session_id": session_id, "reason": "cleared", }) diff --git a/backend/tests/test_browser_agent_loop.py b/backend/tests/test_browser_agent_loop.py index 40cd69df..96ec27bc 100644 --- a/backend/tests/test_browser_agent_loop.py +++ b/backend/tests/test_browser_agent_loop.py @@ -137,8 +137,8 @@ def _install(mocker, monkeypatch: pytest.MonkeyPatch, primary, aux): return {"text": f"GET {params.get('url')} -> HTTP 200\n{{\"docs\": []}}", "status": 200, "url": DOC_URL} return {"text": "ok", "url": DOC_URL} - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _send_browser_command, raising=False) - monkeypatch.setattr(BA.ws_manager, "send_to_session", AsyncMock(return_value=None), raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _send_browser_command, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", AsyncMock(return_value=None), raising=False) return sent @@ -491,13 +491,13 @@ def test_replay_falls_back_to_full_agent_when_a_step_fails(monkeypatch, mocker): aux = FakeAux() sent = _install(mocker, monkeypatch, primary, aux) # make click_by_name FAIL (target gone) so replay must fall back - orig = BA.ws_manager.send_browser_command + orig = BA.WS_MANAGER.send_browser_command async def _fail_cbn(request_id, action, browser_id, params, tab_id=""): if action == "click_by_name": sent.append({"action": action, "params": params}) return {"error": 'No element matching name="Save" on this page.'} return await orig(request_id, action, browser_id, params, tab_id) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _fail_cbn, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _fail_cbn, raising=False) r = asyncio.run(BA.run_browser_agent( task="click the Save button", browser_id="b1", model="sonnet", initial_url=DOC_URL, @@ -527,7 +527,7 @@ def test_deferred_replay_fires_after_navigating_to_the_right_host(monkeypatch, m ]) sent = _install(mocker, monkeypatch, primary, FakeAux()) GOOGLE = "https://www.google.com/" - orig = BA.ws_manager.send_browser_command + orig = BA.WS_MANAGER.send_browser_command async def _cmd(request_id, action, browser_id, params, tab_id=""): # perception + reads report GOOGLE (so the DISPATCH replay misses there), @@ -535,7 +535,7 @@ def test_deferred_replay_fires_after_navigating_to_the_right_host(monkeypatch, m if action in ("list_interactives", "get_text"): return {"text": "stuff", "url": GOOGLE} return await orig(request_id, action, browser_id, params, tab_id) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _cmd, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _cmd, raising=False) # NO initial_url -> dispatch perceives google -> dispatch replay misses. r = asyncio.run(BA.run_browser_agent( @@ -567,13 +567,13 @@ def test_deferred_replay_does_not_fire_after_the_page_was_dirtied(monkeypatch, m ]) sent = _install(mocker, monkeypatch, primary, FakeAux()) GOOGLE = "https://www.google.com/" - orig = BA.ws_manager.send_browser_command + orig = BA.WS_MANAGER.send_browser_command async def _cmd(request_id, action, browser_id, params, tab_id=""): if action in ("list_interactives", "get_text"): return {"text": "stuff", "url": GOOGLE} return await orig(request_id, action, browser_id, params, tab_id) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _cmd, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _cmd, raising=False) r = asyncio.run(BA.run_browser_agent( task="Please click the Search button", browser_id="b1", model="sonnet", @@ -746,14 +746,14 @@ def test_unproven_skill_that_fails_is_quarantined_and_never_retried(monkeypatch, primary = FakeLLM([Resp([Blk("text", "full agent handled it")], stop_reason="end_turn")]) aux = FakeAux() sent = _install(mocker, monkeypatch, primary, aux) - orig = BA.ws_manager.send_browser_command + orig = BA.WS_MANAGER.send_browser_command async def _fail_cbn(request_id, action, browser_id, params, tab_id=""): if action == "click_by_name": sent.append({"action": action, "params": params}) return {"error": 'No element matching name="Save" on this page.'} return await orig(request_id, action, browser_id, params, tab_id) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _fail_cbn, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _fail_cbn, raising=False) # Run 1: replay is attempted, the step fails -> skill is quarantined. asyncio.run(BA.run_browser_agent( @@ -807,13 +807,13 @@ def test_read_answered_from_frontloaded_perception_is_not_a_ghost(monkeypatch, m ]) captured = {} _install(mocker, monkeypatch, primary, FakeAux()) - orig = BA.ws_manager.send_to_session + orig = BA.WS_MANAGER.send_to_session async def _cap(session_id, event, payload): if event == "agent:status": captured["status"] = payload.get("status") return await orig(session_id, event, payload) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) r = asyncio.run(BA.run_browser_agent( task="read me the first sentence", browser_id="b1", model="sonnet", initial_url=DOC_URL, @@ -839,13 +839,13 @@ def test_ghost_completion_is_reported_as_error_not_completed(monkeypatch, mocker _install(mocker, monkeypatch, primary, aux) # every click errors (the fake returns an error for action 'click') captured = {} - orig_send = BA.ws_manager.send_to_session + orig_send = BA.WS_MANAGER.send_to_session async def _cap(session_id, event, payload): if event == "agent:status": captured["status"] = payload.get("status") return await orig_send(session_id, event, payload) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) r = asyncio.run(BA.run_browser_agent( task="Submit the form", browser_id="b1", model="sonnet", initial_url=DOC_URL, @@ -877,15 +877,15 @@ def test_dead_browser_card_aborts_fast_without_spinning(monkeypatch, mocker): card_gone = AsyncMock(return_value={ "error": "Browser card 'b1' not found or not an Electron webview", }) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", card_gone, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", card_gone, raising=False) captured = {} - orig = BA.ws_manager.send_to_session + orig = BA.WS_MANAGER.send_to_session async def _cap(session_id, event, payload): if event == "agent:status": captured["status"] = payload.get("status") return await orig(session_id, event, payload) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) r = asyncio.run(BA.run_browser_agent( task="Click submit", browser_id="b1", model="sonnet", initial_url=DOC_URL, @@ -911,15 +911,15 @@ def test_hung_browser_card_aborts_fast_not_a_20_minute_loop(monkeypatch, mocker) # a wedged tab returns the same timeout error to every command hung = AsyncMock(return_value={"error": "Browser command timed out"}) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", hung, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", hung, raising=False) captured = {} - orig = BA.ws_manager.send_to_session + orig = BA.WS_MANAGER.send_to_session async def _cap(session_id, event, payload): if event == "agent:status": captured["status"] = payload.get("status") return await orig(session_id, event, payload) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) r = asyncio.run(BA.run_browser_agent( task="Read the page", browser_id="b1", model="sonnet", initial_url=DOC_URL, @@ -1070,14 +1070,14 @@ def test_ambient_memory_signals_fire_calmly(monkeypatch, mocker): return Resp([Blk("text", _json.dumps({"playbook": ["search company+React, not generic"]}))], stop_reason="end_turn") msgs = [] - orig = BA.ws_manager.send_to_session + orig = BA.WS_MANAGER.send_to_session async def _cap(session_id, event, payload): if event == "agent:message": c = payload.get("message", {}).get("content") msgs.append(c if isinstance(c, str) else (c or {}).get("text", "")) return await orig(session_id, event, payload) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) def _run(): return FakeLLM([ @@ -1090,7 +1090,7 @@ def test_ambient_memory_signals_fire_calmly(monkeypatch, mocker): # Run 1: nothing learned yet -> NO recall line, but it learns -> closing line. _install(mocker, monkeypatch, _run(), PBAux()) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) asyncio.run(BA.run_browser_agent(task="find engineers", browser_id="b1", model="sonnet", initial_url=DOC_URL)) joined1 = " ".join(msgs) assert "Picking up what I learned" not in joined1, "no recall on the first-ever visit" @@ -1099,7 +1099,7 @@ def test_ambient_memory_signals_fire_calmly(monkeypatch, mocker): # Run 2: now there's a playbook -> recall line fires. msgs.clear() _install(mocker, monkeypatch, _run(), PBAux()) - monkeypatch.setattr(BA.ws_manager, "send_to_session", _cap, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_to_session", _cap, raising=False) asyncio.run(BA.run_browser_agent(task="find more", browser_id="b2", model="sonnet", initial_url=DOC_URL)) assert any("Picking up what I learned about docs.google.com" in m for m in msgs), "recall line on a return visit" @@ -1162,7 +1162,7 @@ def test_batch_replay_runs_a_read_loop_for_all_values(monkeypatch, mocker): if action == "navigate": return {"text": "Navigated", "url": params.get("url")} return {"text": "ok", "url": DOC_URL} - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _data, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _data, raising=False) asyncio.run(BA.run_browser_agent(task="read three profiles", browser_id="b1", model="sonnet", initial_url=DOC_URL)) navs = [c for c in sent if c["action"] == "navigate" and "/in/" in c["params"].get("url", "")] @@ -1196,7 +1196,7 @@ def test_batch_replay_is_ghost_proof_when_an_item_does_not_match(monkeypatch, mo if action == "navigate": return {"text": "Navigated", "url": params.get("url")} return {"text": "profile data", "url": DOC_URL} - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _vary, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _vary, raising=False) asyncio.run(BA.run_browser_agent(task="read three", browser_id="b1", model="sonnet", initial_url=DOC_URL)) all_msgs = json.dumps([c["messages"] for c in primary.calls]) @@ -1256,13 +1256,13 @@ def test_captured_routes_are_surfaced_once_per_host(monkeypatch, mocker): Resp([Blk("text", "done")], stop_reason="end_turn"), ]) _install(mocker, monkeypatch, primary, FakeAux()) - orig = BA.ws_manager.send_browser_command + orig = BA.WS_MANAGER.send_browser_command async def _with_routes(request_id, action, browser_id, params, tab_id=""): if action == "evaluate": return {"text": "Reddit Programming", "url": DOC_URL, "routes_available": 4} return await orig(request_id, action, browser_id, params, tab_id) - monkeypatch.setattr(BA.ws_manager, "send_browser_command", _with_routes, raising=False) + monkeypatch.setattr(BA.WS_MANAGER, "send_browser_command", _with_routes, raising=False) asyncio.run(BA.run_browser_agent(task="browse", browser_id="b1", model="sonnet", initial_url=DOC_URL)) # messages are cumulative across calls, so count within ONE call's full diff --git a/backend/tests/test_browser_fast_path.py b/backend/tests/test_browser_fast_path.py index 500b156a..545a993a 100644 --- a/backend/tests/test_browser_fast_path.py +++ b/backend/tests/test_browser_fast_path.py @@ -103,7 +103,7 @@ def test_dispatch_refused_when_no_dashboard_connected(monkeypatch): # genuinely-closed window that wait just elapses and it still refuses without # dispatching an agent or burning a turn. Zero the wait so the test is instant. monkeypatch.setattr(wsm, "_WS_RECONNECT_WAIT_S", 0.0) - assert not wsm.ws_manager.global_connections + assert not wsm.WS_MANAGER.global_connections results = asyncio.run(run_browser_agents(tasks=[{"task": "go to example.com"}], model="sonnet")) assert len(results) == 1 assert results[0]["summary"].startswith("Error: no dashboard window is connected") diff --git a/backend/tests/test_disconnect_resilience.py b/backend/tests/test_disconnect_resilience.py index e00db6e8..b65f0ef7 100644 --- a/backend/tests/test_disconnect_resilience.py +++ b/backend/tests/test_disconnect_resilience.py @@ -82,13 +82,13 @@ def _build_app(seq_log): so the test thread can drive event emission through the same event loop as the WS handler, avoiding the cross-loop hazards of `asyncio.run()` mid-test.""" - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.ws_manager import WS_MANAGER app = FastAPI() @app.websocket("/ws/agents/{session_id}") async def ws_session(websocket: WebSocket, session_id: str): - await ws_manager.connect_session(session_id, websocket) + await WS_MANAGER.connect_session(session_id, websocket) try: while True: data = await websocket.receive_text() @@ -97,7 +97,7 @@ def _build_app(seq_log): payload = msg.get("data", {}) if event == "client:hello": last_seq = int(payload.get("last_seq") or 0) - ack = await ws_manager.replay_to(session_id, websocket, last_seq) + ack = await WS_MANAGER.replay_to(session_id, websocket, last_seq) await websocket.send_text(json.dumps({ "event": "server:hello", "session_id": session_id, @@ -114,7 +114,7 @@ def _build_app(seq_log): "data": {"nonce": payload.get("nonce")}, })) except WebSocketDisconnect: - ws_manager.disconnect_session(session_id, websocket) + WS_MANAGER.disconnect_session(session_id, websocket) @app.post("/test/emit/{session_id}") async def emit_events(session_id: str, body: dict): @@ -148,11 +148,11 @@ async def _emit_run(session_id: str, n_events: int, terminate: str | None = "com fanning out the broadcast across multiple coroutines. The seq log must still order them strictly. """ - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.ws_manager import WS_MANAGER async def emit_chunk(start: int, count: int): for i in range(count): - await ws_manager.send_to_session(session_id, "agent:stream_delta", { + await WS_MANAGER.send_to_session(session_id, "agent:stream_delta", { "session_id": session_id, "message_id": "m1", "delta": f"chunk-{start + i}", @@ -176,7 +176,7 @@ async def _emit_run(session_id: str, n_events: int, terminate: str | None = "com await emit_chunk(per * concurrent_tasks, rem) if terminate is not None: - await ws_manager.send_to_session(session_id, "agent:status", { + await WS_MANAGER.send_to_session(session_id, "agent:status", { "session_id": session_id, "status": terminate, }) @@ -516,13 +516,13 @@ def test_disconnect_does_not_touch_agent_task(_patch_persist_dir): this test will catch it. We import agent_manager lazily so the `tasks` dict starts empty; we register a sentinel task and confirm disconnect_session doesn't poke it.""" - from backend.apps.agents.core.ws_manager import ws_manager + from backend.apps.agents.core.ws_manager import WS_MANAGER # Insert a real Future into a parallel registry to mimic # `agent_manager.tasks[session_id]` and confirm ws_manager # never reaches into it. We don't import agent_manager (heavy); # we just inspect the source. import inspect - src = inspect.getsource(ws_manager.disconnect_session) + src = inspect.getsource(WS_MANAGER.disconnect_session) assert "cancel" not in src.lower() assert "agent_manager" not in src assert "tasks" not in src