mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-11 20:27:44 +02:00
[haik]: refactor: rename ws_manager singleton to WS_MANAGER and await_reconnect to await_reconnect across 8 files. Module-level ConnectionManager instance follows P/uppercase constant convention (ws_manager -> WS_MANAGER in ws_manager.py). Cross-module helper _await_reconnect drops leading underscore since browser_agent.py imports it directly. Updated all ~80 call sites in agent_manager.py, agents.py, browser_agent.py, main.py, and ws_manager.py; updated monkeypatch targets in test_browser_agent_loop.py, test_browser_fast_path.py, and test_disconnect_resilience.py. 176 insertions, 176 deletions, no behavioral changes.
This commit is contained in:
@@ -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"),
|
||||
})
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -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": (
|
||||
|
||||
@@ -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()
|
||||
|
||||
+15
-15
@@ -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",
|
||||
})
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user