From 1b96bc597c6325b19e9f7f4cf31d0757872847e3 Mon Sep 17 00:00:00 2001 From: ciregenz Date: Mon, 6 Jul 2026 14:16:31 -0700 Subject: [PATCH] [eric] context: lossless tool-report spill (browser+gws caps point at the full file) + visible recovery pill --- backend/apps/agents/agent_manager.py | 7 +++ .../apps/agents/browser_agent_mcp_server.py | 59 ++++++++++++++++--- .../cap_tool_result.py | 44 ++++++++++++-- .../tests/test_browser_agent_mcp_format.py | 26 +++++--- backend/tests/test_gws_cap_tool_result.py | 18 ++++-- .../src/app/pages/AgentChat/AgentChat.tsx | 2 + .../AgentChat/shell/ContextRecoveredPill.tsx | 45 ++++++++++++++ frontend/src/shared/state/agentsSlice.ts | 13 ++++ frontend/src/shared/ws/WebSocketManager.ts | 8 +++ 9 files changed, 198 insertions(+), 24 deletions(-) create mode 100644 frontend/src/app/pages/AgentChat/shell/ContextRecoveredPill.tsx diff --git a/backend/apps/agents/agent_manager.py b/backend/apps/agents/agent_manager.py index 6f642ac5..e69b44ad 100644 --- a/backend/apps/agents/agent_manager.py +++ b/backend/apps/agents/agent_manager.py @@ -151,6 +151,13 @@ class AgentManager(SessionLifecycle, SessionPersistence, Messaging, SessionContr "message_id": p_tool_msg_id, }) self.live_partial.pop(session_id, None) + # Tell the user we self-healed instead of retrying in silence: the frontend renders this as a muted transient pill (same language as the rate-limit pill), not an error card. + try: + await ws_manager.send_to_session(session_id, "agent:context_recovered", { + "session_id": session_id, + }) + except Exception: + logger.debug("context_recovered broadcast failed", exc_info=True) try: from backend.apps.service.client import submit_diagnostic from backend.apps.agents.core.error_classify import redact_for_telemetry diff --git a/backend/apps/agents/browser_agent_mcp_server.py b/backend/apps/agents/browser_agent_mcp_server.py index d2c6df06..39021d23 100644 --- a/backend/apps/agents/browser_agent_mcp_server.py +++ b/backend/apps/agents/browser_agent_mcp_server.py @@ -5,6 +5,7 @@ import base64 import json import sys import os +import time import urllib.request import urllib.error from io import BytesIO @@ -192,16 +193,43 @@ def call_backend(tasks: list[dict]) -> dict: MAX_IMAGE_B64_BYTES = 400_000 MAX_SUMMARY_CHARS = 16_000 MAX_ACTION_LOG_ENTRIES = 40 +REPORT_DIR = os.environ.get( + "OPENSWARM_TOOL_REPORT_DIR", + os.path.join(os.path.expanduser("~"), ".openswarm", "tool-reports"), +) -def p_cap_summary(text: str) -> str: - """Head+tail split: the CLI hard-rejects tool results past ~25K tokens, and a vanished report is worse than a trimmed one.""" +def spill_full_report(text: str, prefix: str) -> str: + """Write the unabridged report to disk so trimming is lossless: the agent can Read + the file (with offset/limit) whenever the capped version isn't enough. Empty string + when the write fails; callers degrade to cap-only.""" + try: + os.makedirs(REPORT_DIR, exist_ok=True) + # Reports are point-in-time working files, not archives; prune week-old ones so the folder can't grow forever. + cutoff = time.time() - 7 * 86400 + for old in os.listdir(REPORT_DIR): + p = os.path.join(REPORT_DIR, old) + try: + if os.path.getmtime(p) < cutoff: + os.remove(p) + except OSError: + pass + path = os.path.join(REPORT_DIR, f"{prefix}-{os.getpid()}-{int(time.time()*1000)}.md") + with open(path, "w", encoding="utf-8") as f: + f.write(text) + return path + except Exception: + return "" + + +def p_cap_summary(text: str) -> tuple[str, bool]: + """Head+tail split, plus a truncated? flag so the caller can spill the full text: the CLI hard-rejects tool results past ~25K tokens, and a vanished report is worse than a trimmed one.""" if len(text) <= MAX_SUMMARY_CHARS: - return text + return text, False head = text[: MAX_SUMMARY_CHARS - 4_000] tail = text[-3_500:] omitted = len(text) - len(head) - len(tail) - return f"{head}\n\n[... {omitted} chars of the report omitted ...]\n\n{tail}" + return f"{head}\n\n[... {omitted} chars of the report omitted ...]\n\n{tail}", True def p_sniff_image_mime(b64: str) -> str: @@ -244,23 +272,36 @@ def format_result(result: dict) -> dict: browser_id = result.get("browser_id", "") action_log = result.get("action_log", []) + capped_summary, summary_truncated = p_cap_summary(summary) lines = [f"**Browser Agent Result** (browser: {browser_id}, session: {session_id})", ""] - lines.append(f"**Summary:** {p_cap_summary(summary)}") + lines.append(f"**Summary:** {capped_summary}") + actions_omitted = 0 if action_log: lines.append("") lines.append("**Actions taken:**") entries = action_log[-MAX_ACTION_LOG_ENTRIES:] - omitted = len(action_log) - len(entries) - if omitted > 0: - lines.append(f" (... {omitted} earlier actions omitted ...)") - for i, entry in enumerate(entries, omitted + 1): + actions_omitted = len(action_log) - len(entries) + if actions_omitted > 0: + lines.append(f" (... {actions_omitted} earlier actions omitted ...)") + for i, entry in enumerate(entries, actions_omitted + 1): tool = entry.get("tool", "?") inp = entry.get("input", {}) ms = entry.get("elapsed_ms", 0) brief = json.dumps(inp)[:120] lines.append(f" {i}. {tool}({brief}) [{ms}ms]") + if summary_truncated or actions_omitted > 0: + full_lines = [f"# Browser Agent Full Report (browser: {browser_id}, session: {session_id})", "", summary, ""] + if action_log: + full_lines.append("## Actions") + for i, entry in enumerate(action_log, 1): + full_lines.append(f"{i}. {entry.get('tool', '?')}({json.dumps(entry.get('input', {}))}) [{entry.get('elapsed_ms', 0)}ms]") + report_path = spill_full_report("\n".join(full_lines), "browser-report") + if report_path: + lines.append("") + lines.append(f"Full unabridged report saved to: {report_path} (use Read with offset/limit for the omitted parts)") + content.append({"type": "text", "text": "\n".join(lines)}) screenshot = result.get("final_screenshot") diff --git a/backend/apps/google_workspace_mcp_shim/cap_tool_result.py b/backend/apps/google_workspace_mcp_shim/cap_tool_result.py index b8e1e8a8..8ce75f21 100644 --- a/backend/apps/google_workspace_mcp_shim/cap_tool_result.py +++ b/backend/apps/google_workspace_mcp_shim/cap_tool_result.py @@ -5,18 +5,46 @@ and unit-testable outside the shim's ephemeral uv env. The bundled Claude CLI hard-rejects any MCP result over ~25K tokens and spills it to a file, which the model then re-reads back in, refilling the context and tripping the CLI's autocompact-thrash. Capping under that spill threshold keeps the result inline and the -model out of the re-read loop.""" +model out of the re-read loop. Lossless: the full text is saved to a report file the +model can Read selectively, and the truncation note points at it.""" +import os +import time from typing import Any MAX_RESULT_CHARS = 48_000 +REPORT_DIR = os.environ.get( + "OPENSWARM_TOOL_REPORT_DIR", + os.path.join(os.path.expanduser("~"), ".openswarm", "tool-reports"), +) P_TRUNCATION_NOTE = ( "\n\n[Truncated: this tool returned more than {cap} characters, too much to fit " - "in context at once. Narrow the request (add a search filter, a date range, or a " - "smaller max_results / page size) or fetch the next page.]" + "in context at once.{saved} Narrow the request (add a search filter, a date range, " + "or a smaller max_results / page size) or fetch the next page.]" ) +def p_spill(text: str) -> str: + """Write the full result to disk so the cap is lossless; empty string on failure.""" + try: + os.makedirs(REPORT_DIR, exist_ok=True) + # Reports are point-in-time working files, not archives; prune week-old ones so the folder can't grow forever. + cutoff = time.time() - 7 * 86400 + for old in os.listdir(REPORT_DIR): + p = os.path.join(REPORT_DIR, old) + try: + if os.path.getmtime(p) < cutoff: + os.remove(p) + except OSError: + pass + path = os.path.join(REPORT_DIR, f"gws-result-{os.getpid()}-{int(time.time()*1000)}.txt") + with open(path, "w", encoding="utf-8") as f: + f.write(text) + return path + except Exception: + return "" + + def cap_tool_result(result: Any, max_chars: int = MAX_RESULT_CHARS) -> Any: """Cap the text content blocks of a call_tool return in place. Duck-typed and fail-open: any shape we don't recognize passes through unchanged, so an upstream @@ -25,6 +53,14 @@ def cap_tool_result(result: Any, max_chars: int = MAX_RESULT_CHARS) -> Any: blocks = result[0] if isinstance(result, tuple) else result if not isinstance(blocks, list): return result + texts = [ + b.text for b in blocks + if getattr(b, "type", None) == "text" and getattr(b, "text", None) is not None + ] + if sum(len(t) for t in texts) <= max_chars: + return result + full_path = p_spill("\n".join(texts)) + saved = f" The complete result was saved to {full_path}; Read it with offset/limit if you truly need the rest." if full_path else "" used = 0 truncated = False for b in blocks: @@ -37,7 +73,7 @@ def cap_tool_result(result: Any, max_chars: int = MAX_RESULT_CHARS) -> Any: if used + len(text) <= max_chars: used += len(text) continue - b.text = text[: max(0, max_chars - used)] + P_TRUNCATION_NOTE.format(cap=max_chars) + b.text = text[: max(0, max_chars - used)] + P_TRUNCATION_NOTE.format(cap=max_chars, saved=saved) truncated = True return result except Exception: diff --git a/backend/tests/test_browser_agent_mcp_format.py b/backend/tests/test_browser_agent_mcp_format.py index 6800c577..1439379c 100644 --- a/backend/tests/test_browser_agent_mcp_format.py +++ b/backend/tests/test_browser_agent_mcp_format.py @@ -29,22 +29,32 @@ def test_small_summary_passes_through_unchanged() -> None: assert "omitted" not in text -def test_giant_summary_keeps_head_and_tail() -> None: +def test_giant_summary_keeps_head_and_tail_and_spills_full_report(tmp_path, monkeypatch) -> None: + import backend.apps.agents.browser_agent_mcp_server as srv + monkeypatch.setattr(srv, "REPORT_DIR", str(tmp_path)) summary = "HEADSTART " + ("x" * 60_000) + " TAILEND" text = result_text(format_result({"summary": summary})) - assert len(text) < MAX_SUMMARY_CHARS + 300 + assert len(text) < MAX_SUMMARY_CHARS + 500 assert "HEADSTART" in text - assert text.endswith("TAILEND") assert "omitted" in text + assert "Full unabridged report saved to:" in text + reports = list(tmp_path.iterdir()) + assert len(reports) == 1 + assert summary in reports[0].read_text() -def test_action_log_keeps_last_entries_with_original_numbering() -> None: +def test_action_log_keeps_last_entries_with_original_numbering(tmp_path, monkeypatch) -> None: + import backend.apps.agents.browser_agent_mcp_server as srv + monkeypatch.setattr(srv, "REPORT_DIR", str(tmp_path)) log = [{"tool": f"Act{i}", "input": {}, "elapsed_ms": i} for i in range(100)] text = result_text(format_result({"summary": "ok", "action_log": log})) assert "(... 60 earlier actions omitted ...)" in text - assert "Act59" not in text assert "61. Act60(" in text assert "100. Act99(" in text + # The full log (including the 60 omitted entries) lands in the spilled report. + reports = list(tmp_path.iterdir()) + assert len(reports) == 1 + assert "Act59" in reports[0].read_text() def test_short_action_log_has_no_omission_line() -> None: @@ -54,11 +64,13 @@ def test_short_action_log_has_no_omission_line() -> None: assert "1. Click(" in text -def test_pathological_result_stays_far_under_cli_rejection_cap() -> None: +def test_pathological_result_stays_far_under_cli_rejection_cap(tmp_path, monkeypatch) -> None: + import backend.apps.agents.browser_agent_mcp_server as srv + monkeypatch.setattr(srv, "REPORT_DIR", str(tmp_path)) log = [{"tool": "T", "input": {"v": "y" * 500}, "elapsed_ms": 1} for i in range(500)] out = format_result({"summary": "z" * 200_000, "action_log": log}) total = len(result_text(out)) - assert total < MAX_SUMMARY_CHARS + MAX_ACTION_LOG_ENTRIES * 160 + 500 + assert total < MAX_SUMMARY_CHARS + MAX_ACTION_LOG_ENTRIES * 160 + 800 def test_error_result_untouched() -> None: diff --git a/backend/tests/test_gws_cap_tool_result.py b/backend/tests/test_gws_cap_tool_result.py index 6b185ef7..90d47e36 100644 --- a/backend/tests/test_gws_cap_tool_result.py +++ b/backend/tests/test_gws_cap_tool_result.py @@ -26,16 +26,24 @@ def test_small_result_untouched() -> None: assert b.text == "one short email" -def test_oversized_single_block_capped_with_marker() -> None: +def test_oversized_single_block_capped_with_marker_and_spilled(tmp_path, monkeypatch) -> None: + import backend.apps.google_workspace_mcp_shim.cap_tool_result as capmod + monkeypatch.setattr(capmod, "REPORT_DIR", str(tmp_path)) b = block("E" * 300_000) cap_tool_result(([b], {"result": "E" * 300_000})) - assert len(b.text) < MAX_RESULT_CHARS + 400 + assert len(b.text) < MAX_RESULT_CHARS + 600 assert b.text.startswith("E") assert "Truncated" in b.text + assert "saved to" in b.text assert len(b.text) // 4 < 25_000 + reports = list(tmp_path.iterdir()) + assert len(reports) == 1 + assert reports[0].read_text() == "E" * 300_000 -def test_budget_spans_multiple_blocks() -> None: +def test_budget_spans_multiple_blocks(tmp_path, monkeypatch) -> None: + import backend.apps.google_workspace_mcp_shim.cap_tool_result as capmod + monkeypatch.setattr(capmod, "REPORT_DIR", str(tmp_path)) a, b, c = block("A" * 40_000), block("B" * 40_000), block("C" * 40_000) cap_tool_result([a, b, c]) assert a.text == "A" * 40_000 @@ -51,7 +59,9 @@ def test_non_text_blocks_pass_through() -> None: assert txt.text == "hello" -def test_bare_list_return_shape() -> None: +def test_bare_list_return_shape(tmp_path, monkeypatch) -> None: + import backend.apps.google_workspace_mcp_shim.cap_tool_result as capmod + monkeypatch.setattr(capmod, "REPORT_DIR", str(tmp_path)) b = block("Z" * 100_000) out = cap_tool_result([b]) assert out is not None diff --git a/frontend/src/app/pages/AgentChat/AgentChat.tsx b/frontend/src/app/pages/AgentChat/AgentChat.tsx index ec28d865..097f6dd5 100644 --- a/frontend/src/app/pages/AgentChat/AgentChat.tsx +++ b/frontend/src/app/pages/AgentChat/AgentChat.tsx @@ -61,6 +61,7 @@ import ToolGroupBubble, { RenderItem, ToolGroup, isToolGroup, isToolPair } from import ApprovalBar, { BatchApprovalBar } from './shell/ApprovalBar'; import ForceStopAgentBar from './ForceStopAgentBar'; import { RateLimitPill } from './shell/RateLimitPill'; +import { ContextRecoveredPill } from './shell/ContextRecoveredPill'; import ChatInput, { ChatInputHandle } from './ChatInput'; import ContextDrawer from './shell/ContextDrawer'; import { ErrorSlime } from '@/app/components/feedback/ErrorSlime'; @@ -1841,6 +1842,7 @@ const AgentChat: React.FC = ({ sessionId: sessionIdProp, onClose )} + {isGlowing ? ( = ({ sessionId }) => { + const c = useClaudeTokens(); + const dispatch = useAppDispatch(); + const cr = useAppSelector((s) => s.agents.sessions[sessionId]?.context_recovered); + + useEffect(() => { + if (!cr) return; + const t = setTimeout(() => dispatch(clearContextRecovered({ sessionId })), 12000); + return () => clearTimeout(t); + }, [cr, sessionId, dispatch]); + + return ( + + + + Recovered and retried + + + ); +}; diff --git a/frontend/src/shared/state/agentsSlice.ts b/frontend/src/shared/state/agentsSlice.ts index 66ed34b9..43c257bc 100644 --- a/frontend/src/shared/state/agentsSlice.ts +++ b/frontend/src/shared/state/agentsSlice.ts @@ -109,6 +109,7 @@ export interface AgentSession { framework_overhead_tokens?: number; context_overflow?: { reason: string; message: string; at: string } | null; rate_limited?: { retry_after_s: number | null; at: string } | null; + context_recovered?: { at: string } | null; mcp_suggestions?: Array<{ id: string; title: string; description: string; reason?: string }>; mcp_suggestions_is_vague?: boolean; compacted_through_msg_id?: string | null; @@ -959,6 +960,16 @@ const agentsSlice = createSlice({ if (session) session.rate_limited = null; }, + setContextRecovered(state, action: PayloadAction<{ sessionId: string }>) { + const session = state.sessions[action.payload.sessionId]; + if (session) session.context_recovered = { at: new Date().toISOString() }; + }, + + clearContextRecovered(state, action: PayloadAction<{ sessionId: string }>) { + const session = state.sessions[action.payload.sessionId]; + if (session) session.context_recovered = null; + }, + clearContextOverflow( state, action: PayloadAction<{ sessionId: string }> @@ -1440,6 +1451,8 @@ export const { setContextOverflow, setRateLimited, clearRateLimited, + setContextRecovered, + clearContextRecovered, clearContextOverflow, setMcpSuggestions, clearMcpSuggestions, diff --git a/frontend/src/shared/ws/WebSocketManager.ts b/frontend/src/shared/ws/WebSocketManager.ts index b2a3ffcd..edc169c7 100644 --- a/frontend/src/shared/ws/WebSocketManager.ts +++ b/frontend/src/shared/ws/WebSocketManager.ts @@ -13,6 +13,7 @@ import { updateSessionContext, setContextOverflow, setRateLimited, + setContextRecovered, setMcpSuggestions, addBranch, setActiveBranch, @@ -554,6 +555,13 @@ class WebSocketManager { } break; + case 'agent:context_recovered': + // The backend hit a context-overflow crash mid-turn, rebuilt from its local copy, and retried on its own. Transient muted pill so the recovery is visible without reading like an error. + if (session_id) { + store.dispatch(setContextRecovered({ sessionId: session_id })); + } + break; + case 'agent:context_status': // Auto-compaction collapsed older turns into a summary. Mirror compacted_through_msg_id locally so the renderer can drop a visible "N earlier turns summarized" chip into the transcript. Other reasons (cleared, etc.) flow through this same event but don't currently need a chip, ignore them for now. if (session_id && data.reason === 'compacted') {