From b3e1fe224e2d55491e8a9420f8893dc014a4bca8 Mon Sep 17 00:00:00 2001 From: ciregenz Date: Tue, 23 Jun 2026 22:37:38 -0700 Subject: [PATCH] [eric] agents: extract run-error classification (long-context/capacity/free-trial/auth/unknown-model cards) into manager/run/error_cards.py; agent_manager 1166->932 --- backend/apps/agents/agent_manager.py | 237 +--------------- backend/apps/agents/manager/run/__init__.py | 0 .../apps/agents/manager/run/error_cards.py | 252 ++++++++++++++++++ 3 files changed, 254 insertions(+), 235 deletions(-) create mode 100644 backend/apps/agents/manager/run/__init__.py create mode 100644 backend/apps/agents/manager/run/error_cards.py diff --git a/backend/apps/agents/agent_manager.py b/backend/apps/agents/agent_manager.py index 71752421..6c7f848e 100644 --- a/backend/apps/agents/agent_manager.py +++ b/backend/apps/agents/agent_manager.py @@ -20,13 +20,6 @@ from backend.apps.tools_lib.tools_lib import ( from backend.apps.agents.core.error_classify import ( CAPACITY_BACKOFFS, capacity_retry_wait, - is_auth_error, - is_free_trial_exhausted, - is_long_context_error, - is_transient_capacity_error, - is_unknown_model_error, - parse_retry_after, - redact_for_telemetry, ) # SESSIONS_DIR is re-exported on purpose: session_store reads agent_manager.SESSIONS_DIR at # call time (dodging a circular import), and the disk-resilience test monkeypatches it here. @@ -57,6 +50,7 @@ from backend.apps.agents.manager.AgentLaunchMixin import AgentLaunchMixin from backend.apps.agents.manager.MockAgentMixin import MockAgentMixin from backend.apps.agents.manager.RunSupportMixin import RunSupportMixin from backend.apps.agents.manager.permissions import gate_hooks +from backend.apps.agents.manager.run.error_cards import handle_run_error from backend.apps.agents.manager.session.workspace_git import ensure_cwd_git_repo from backend.apps.agents.manager.prompt.tool_catalog import ( get_all_tool_names, @@ -862,234 +856,7 @@ class AgentManager(SessionLifecycleMixin, SessionPersistenceMixin, MessagingMixi turn.stream_text_msg_id = None turn.stream_text_accum = "" except Exception as e: - logger.exception(f"Agent {session_id} error: {e}") - session.status = "error" - - # Long-context-required 429 fork: surface a friendly overflow event - # so the frontend can render an actionable card ("Switch to Chat - # mode" / "Start a fresh chat") instead of a raw error blob. The - # user can't recover by waiting, this is a tier-gate, not a rate - # limit, so the UX matters. - try: - p_stderr_tail = "\n".join(p_stderr_buffer[-50:]) - except Exception: - p_stderr_tail = "" - # If we already streamed a substantive assistant response this - # turn, the user got their answer; the error fired on a - # subsequent step (title gen, follow-up tool turn, etc.). - # Don't blast a "context exceeded" card over a completed reply. - p_streamed_substantive = bool(turn.stream_text_msg_id) and turn.current_turn_emitted - if p_streamed_substantive and is_long_context_error(e, extra_text=p_stderr_tail): - # Mark the session completed (not error), keep the assistant - # reply visible, and skip the overflow card. The next user - # turn will properly hit the pre-send guard if the chat is - # still over cap. - session.status = "completed" - if turn.stream_text_msg_id: - try: - await ws_manager.send_to_session(session_id, "agent:stream_end", { - "session_id": session_id, - "message_id": turn.stream_text_msg_id, - }) - except Exception: - pass - return - if is_long_context_error(e, extra_text=p_stderr_tail): - friendly_msg = ( - "This conversation has grown too large for your account's " - "standard context window. Long-context requests require an " - "upgraded tier, switch to Chat mode or start a fresh chat " - "to continue." - ) - error_msg = Message(role="system", content=friendly_msg, branch_id=session.active_branch_id) - session.messages.append(error_msg) - p_ovf_payload = { - "session_id": session_id, - "reason": "long_context_required", - "message": friendly_msg, - "model": session.model, - "provider": session.provider, - "context_window": session.context_window, - "framework_overhead_tokens": session.framework_overhead_tokens, - "input_tokens": session.tokens.get("input", 0), - "active_mcps": list(session.active_mcps), - "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", p_ovf_payload) - await ws_manager.send_to_session(session_id, "agent:message", { - "session_id": session_id, - "message": error_msg.model_dump(mode="json"), - }) - try: - from backend.apps.service.client import submit_diagnostic - submit_diagnostic({ - "kind": "context_overflow", - "where": "agent_manager.p_run_streaming_turn", - "session_id": session_id, - "model": session.model, - "provider": session.provider, - "context_window": session.context_window, - "input_tokens": session.tokens.get("input", 0), - "framework_overhead_tokens": session.framework_overhead_tokens, - "active_mcps_count": len(session.active_mcps), - "messages_count": len(session.messages), - "error_preview": redact_for_telemetry(str(e), limit=500), - }) - except Exception: - logger.debug("submit_diagnostic for context_overflow failed", exc_info=True) - elif is_transient_capacity_error(e, extra_text=p_stderr_tail): - # A genuine throttle (429/overload/capacity) that already burned - # the whole silent-backoff budget (the only way one reaches here). - # It's a limit, not a failure, so don't append a system-message - # card; emit a transient signal for the muted pill and mark the - # turn completed so it doesn't read as an error. - session.status = "completed" - if turn.stream_text_msg_id: - try: - await ws_manager.send_to_session(session_id, "agent:stream_end", { - "session_id": session_id, - "message_id": turn.stream_text_msg_id, - }) - except Exception: - pass - await ws_manager.send_to_session(session_id, "agent:rate_limited", { - "session_id": session_id, - "retry_after_s": parse_retry_after(e, p_stderr_tail), - }) - elif is_free_trial_exhausted(e, extra_text=p_stderr_tail): - # Free runs spent. Flip back to own_key and show a friendly - # "connect a model" upsell instead of a raw 402. - try: - from backend.apps.subscription.free_trial import clear_free_trial - await clear_free_trial(load_settings()) - except Exception: - logger.debug("clear_free_trial after exhaustion failed", exc_info=True) - friendly_msg = ( - "You've used your free runs. Connect a model to keep going: " - "your own API key, an AI subscription you already pay for, or " - "OpenSwarm Pro." - ) - 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", { - "session_id": session_id, - "message": friendly_msg, - }) - await ws_manager.send_to_session(session_id, "agent:message", { - "session_id": session_id, - "message": error_msg.model_dump(mode="json"), - }) - elif is_auth_error(e, extra_text=p_stderr_tail): - # Three sub-cases the user can hit, with distinct fixes: - # 1. "No credentials for provider: claude", user picked a - # -cc route but doesn't have Claude Pro/Max connected - # via 9Router. Tell them to either connect Claude - # Pro/Max OR pick a non--cc model. - # 2. OpenSwarm Pro 401, bearer expired. Reconnect. - # 3. Anthropic API key 401, wrong key. Re-enter. - p_model = (session.model or "").lower() - p_combined = f"{e!s}\n{p_stderr_tail}".lower() - # Codex/OpenAI subscription tokens rotate every ~2-3 - # minutes, the user sees the rotation window as a 401 - # with "reset after 1m 59s" or similar. Don't ask them to - # reconnect; just tell them to wait it out and retry. - if ( - ("codex/" in p_combined or "[codex/" in p_combined or p_model.startswith(("cx/", "gpt-"))) - and ("authentication token is expired" in p_combined or "authentication token has expired" in p_combined or "401" in p_combined) - ): - friendly_msg = ( - "GPT subscription token just rotated, this is " - "automatic and resets every couple minutes. Send " - "your message again in ~1 minute and it'll go " - "through. (No need to reconnect anything.)" - ) - reason = "codex_token_rotating" - elif "no credentials for provider" in p_combined: - friendly_msg = ( - "Selected route requires Claude Pro / Max, but it's " - "not connected. Open Settings → Models and either " - "connect Claude Pro / Max, or switch the model to a " - "non-`-cc` variant (e.g. Claude Sonnet 4.6 instead " - "of Sonnet 4.6 -cc)." - ) - reason = "claude_sub_not_connected" - elif ( - "-cc" not in p_model - and getattr(load_settings(), "connection_mode", "own_key") == "openswarm-pro" - ): - friendly_msg = ( - "OpenSwarm Pro authentication failed. Your subscription " - "token may have expired even though the connection still " - "shows green. Open Settings → Models and click " - "Disconnect / Reconnect on Claude Pro / Max to refresh " - "the token." - ) - reason = "openswarm_pro_auth_expired" - else: - friendly_msg = ( - "Anthropic authentication failed. The API key or " - "subscription token for this model is invalid. Open " - "Settings → Models and re-enter the API key, or " - "reconnect Claude Pro / Max." - ) - 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", { - "session_id": session_id, - "reason": reason, - "message": friendly_msg, - "model": session.model, - }) - await ws_manager.send_to_session(session_id, "agent:message", { - "session_id": session_id, - "message": error_msg.model_dump(mode="json"), - }) - elif is_unknown_model_error(e, extra_text=p_stderr_tail): - # Upstream rejected the model code itself (e.g. Codex 1211 on a - # ChatGPT plan that lacks our GPT ids). Track it; the friendly - # "add an API key / pick another model" card is rendered frontend-side. - try: - from backend.apps.service.client import submit_diagnostic - submit_diagnostic({ - "kind": "model_error", - "subkind": "unknown_model", - "model": session.model, - "provider": session.provider, - "connection_mode": getattr(load_settings(), "connection_mode", "own_key"), - "error_preview": redact_for_telemetry(str(e), limit=400), - "stderr_tail": redact_for_telemetry(p_stderr_tail), - }) - except Exception: - 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", { - "session_id": session_id, - "message": error_msg.model_dump(mode="json"), - }) - else: - # Track unclassified agent failures too so we stop flying blind on them. - try: - from backend.apps.service.client import submit_diagnostic - submit_diagnostic({ - "kind": "model_error", - "subkind": "unclassified", - "model": session.model, - "provider": session.provider, - "connection_mode": getattr(load_settings(), "connection_mode", "own_key"), - "error_preview": redact_for_telemetry(str(e), limit=400), - "stderr_tail": redact_for_telemetry(p_stderr_tail), - }) - except Exception: - 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", { - "session_id": session_id, - "message": error_msg.model_dump(mode="json"), - }) + await handle_run_error(e, session, session_id, turn, p_stderr_buffer) except BaseException as e: # Catch BaseExceptionGroup from anyio task groups (e.g. concurrent # CLI crash + pending approval cancellation) so it doesn't escape diff --git a/backend/apps/agents/manager/run/__init__.py b/backend/apps/agents/manager/run/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/backend/apps/agents/manager/run/error_cards.py b/backend/apps/agents/manager/run/error_cards.py new file mode 100644 index 00000000..f7e23ad2 --- /dev/null +++ b/backend/apps/agents/manager/run/error_cards.py @@ -0,0 +1,252 @@ +"""Friendly error cards for a failed agent run. The run_agent_loop except-handler classifies +the exception (long-context / capacity / free-trial / auth / unknown-model / unclassified) and +emits the matching system message + WS event. Pulled out of agent_manager so the loop stays under +the file ceiling; pure relocation, no self (operates on the passed run state).""" + +import logging + +from backend.apps.agents.core.models import Message +from backend.apps.agents.core.ws_manager import ws_manager +from backend.apps.settings.settings import load_settings +from backend.apps.agents.core.error_classify import ( + is_long_context_error, + is_transient_capacity_error, + is_free_trial_exhausted, + is_auth_error, + is_unknown_model_error, + parse_retry_after, + redact_for_telemetry, +) + +logger = logging.getLogger(__name__) + + +async def handle_run_error(e, session, session_id, turn, p_stderr_buffer) -> None: + logger.exception(f"Agent {session_id} error: {e}") + session.status = "error" + + # Long-context-required 429 fork: surface a friendly overflow event + # so the frontend can render an actionable card ("Switch to Chat + # mode" / "Start a fresh chat") instead of a raw error blob. The + # user can't recover by waiting, this is a tier-gate, not a rate + # limit, so the UX matters. + try: + p_stderr_tail = "\n".join(p_stderr_buffer[-50:]) + except Exception: + p_stderr_tail = "" + # If we already streamed a substantive assistant response this + # turn, the user got their answer; the error fired on a + # subsequent step (title gen, follow-up tool turn, etc.). + # Don't blast a "context exceeded" card over a completed reply. + p_streamed_substantive = bool(turn.stream_text_msg_id) and turn.current_turn_emitted + if p_streamed_substantive and is_long_context_error(e, extra_text=p_stderr_tail): + # Mark the session completed (not error), keep the assistant + # reply visible, and skip the overflow card. The next user + # turn will properly hit the pre-send guard if the chat is + # still over cap. + session.status = "completed" + if turn.stream_text_msg_id: + try: + await ws_manager.send_to_session(session_id, "agent:stream_end", { + "session_id": session_id, + "message_id": turn.stream_text_msg_id, + }) + except Exception: + pass + return + if is_long_context_error(e, extra_text=p_stderr_tail): + friendly_msg = ( + "This conversation has grown too large for your account's " + "standard context window. Long-context requests require an " + "upgraded tier, switch to Chat mode or start a fresh chat " + "to continue." + ) + error_msg = Message(role="system", content=friendly_msg, branch_id=session.active_branch_id) + session.messages.append(error_msg) + p_ovf_payload = { + "session_id": session_id, + "reason": "long_context_required", + "message": friendly_msg, + "model": session.model, + "provider": session.provider, + "context_window": session.context_window, + "framework_overhead_tokens": session.framework_overhead_tokens, + "input_tokens": session.tokens.get("input", 0), + "active_mcps": list(session.active_mcps), + "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", p_ovf_payload) + await ws_manager.send_to_session(session_id, "agent:message", { + "session_id": session_id, + "message": error_msg.model_dump(mode="json"), + }) + try: + from backend.apps.service.client import submit_diagnostic + submit_diagnostic({ + "kind": "context_overflow", + "where": "agent_manager.p_run_streaming_turn", + "session_id": session_id, + "model": session.model, + "provider": session.provider, + "context_window": session.context_window, + "input_tokens": session.tokens.get("input", 0), + "framework_overhead_tokens": session.framework_overhead_tokens, + "active_mcps_count": len(session.active_mcps), + "messages_count": len(session.messages), + "error_preview": redact_for_telemetry(str(e), limit=500), + }) + except Exception: + logger.debug("submit_diagnostic for context_overflow failed", exc_info=True) + elif is_transient_capacity_error(e, extra_text=p_stderr_tail): + # A genuine throttle (429/overload/capacity) that already burned + # the whole silent-backoff budget (the only way one reaches here). + # It's a limit, not a failure, so don't append a system-message + # card; emit a transient signal for the muted pill and mark the + # turn completed so it doesn't read as an error. + session.status = "completed" + if turn.stream_text_msg_id: + try: + await ws_manager.send_to_session(session_id, "agent:stream_end", { + "session_id": session_id, + "message_id": turn.stream_text_msg_id, + }) + except Exception: + pass + await ws_manager.send_to_session(session_id, "agent:rate_limited", { + "session_id": session_id, + "retry_after_s": parse_retry_after(e, p_stderr_tail), + }) + elif is_free_trial_exhausted(e, extra_text=p_stderr_tail): + # Free runs spent. Flip back to own_key and show a friendly + # "connect a model" upsell instead of a raw 402. + try: + from backend.apps.subscription.free_trial import clear_free_trial + await clear_free_trial(load_settings()) + except Exception: + logger.debug("clear_free_trial after exhaustion failed", exc_info=True) + friendly_msg = ( + "You've used your free runs. Connect a model to keep going: " + "your own API key, an AI subscription you already pay for, or " + "OpenSwarm Pro." + ) + 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", { + "session_id": session_id, + "message": friendly_msg, + }) + await ws_manager.send_to_session(session_id, "agent:message", { + "session_id": session_id, + "message": error_msg.model_dump(mode="json"), + }) + elif is_auth_error(e, extra_text=p_stderr_tail): + # Three sub-cases the user can hit, with distinct fixes: + # 1. "No credentials for provider: claude", user picked a + # -cc route but doesn't have Claude Pro/Max connected + # via 9Router. Tell them to either connect Claude + # Pro/Max OR pick a non--cc model. + # 2. OpenSwarm Pro 401, bearer expired. Reconnect. + # 3. Anthropic API key 401, wrong key. Re-enter. + p_model = (session.model or "").lower() + p_combined = f"{e!s}\n{p_stderr_tail}".lower() + # Codex/OpenAI subscription tokens rotate every ~2-3 + # minutes, the user sees the rotation window as a 401 + # with "reset after 1m 59s" or similar. Don't ask them to + # reconnect; just tell them to wait it out and retry. + if ( + ("codex/" in p_combined or "[codex/" in p_combined or p_model.startswith(("cx/", "gpt-"))) + and ("authentication token is expired" in p_combined or "authentication token has expired" in p_combined or "401" in p_combined) + ): + friendly_msg = ( + "GPT subscription token just rotated, this is " + "automatic and resets every couple minutes. Send " + "your message again in ~1 minute and it'll go " + "through. (No need to reconnect anything.)" + ) + reason = "codex_token_rotating" + elif "no credentials for provider" in p_combined: + friendly_msg = ( + "Selected route requires Claude Pro / Max, but it's " + "not connected. Open Settings → Models and either " + "connect Claude Pro / Max, or switch the model to a " + "non-`-cc` variant (e.g. Claude Sonnet 4.6 instead " + "of Sonnet 4.6 -cc)." + ) + reason = "claude_sub_not_connected" + elif ( + "-cc" not in p_model + and getattr(load_settings(), "connection_mode", "own_key") == "openswarm-pro" + ): + friendly_msg = ( + "OpenSwarm Pro authentication failed. Your subscription " + "token may have expired even though the connection still " + "shows green. Open Settings → Models and click " + "Disconnect / Reconnect on Claude Pro / Max to refresh " + "the token." + ) + reason = "openswarm_pro_auth_expired" + else: + friendly_msg = ( + "Anthropic authentication failed. The API key or " + "subscription token for this model is invalid. Open " + "Settings → Models and re-enter the API key, or " + "reconnect Claude Pro / Max." + ) + 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", { + "session_id": session_id, + "reason": reason, + "message": friendly_msg, + "model": session.model, + }) + await ws_manager.send_to_session(session_id, "agent:message", { + "session_id": session_id, + "message": error_msg.model_dump(mode="json"), + }) + elif is_unknown_model_error(e, extra_text=p_stderr_tail): + # Upstream rejected the model code itself (e.g. Codex 1211 on a + # ChatGPT plan that lacks our GPT ids). Track it; the friendly + # "add an API key / pick another model" card is rendered frontend-side. + try: + from backend.apps.service.client import submit_diagnostic + submit_diagnostic({ + "kind": "model_error", + "subkind": "unknown_model", + "model": session.model, + "provider": session.provider, + "connection_mode": getattr(load_settings(), "connection_mode", "own_key"), + "error_preview": redact_for_telemetry(str(e), limit=400), + "stderr_tail": redact_for_telemetry(p_stderr_tail), + }) + except Exception: + 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", { + "session_id": session_id, + "message": error_msg.model_dump(mode="json"), + }) + else: + # Track unclassified agent failures too so we stop flying blind on them. + try: + from backend.apps.service.client import submit_diagnostic + submit_diagnostic({ + "kind": "model_error", + "subkind": "unclassified", + "model": session.model, + "provider": session.provider, + "connection_mode": getattr(load_settings(), "connection_mode", "own_key"), + "error_preview": redact_for_telemetry(str(e), limit=400), + "stderr_tail": redact_for_telemetry(p_stderr_tail), + }) + except Exception: + 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", { + "session_id": session_id, + "message": error_msg.model_dump(mode="json"), + })