mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-28 12:34:50 +02:00
[eric] context: lossless tool-report spill (browser+gws caps point at the full file) + visible recovery pill
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<AgentChatProps> = ({ sessionId: sessionIdProp, onClose
|
||||
)}
|
||||
|
||||
<RateLimitPill sessionId={session.id} />
|
||||
<ContextRecoveredPill sessionId={session.id} />
|
||||
|
||||
{isGlowing ? (
|
||||
<Box
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
import React, { useEffect } from 'react';
|
||||
import Box from '@mui/material/Box';
|
||||
import Fade from '@mui/material/Fade';
|
||||
import Typography from '@mui/material/Typography';
|
||||
import RestartAltIcon from '@mui/icons-material/RestartAlt';
|
||||
import { useAppDispatch, useAppSelector } from '@/shared/hooks';
|
||||
import { clearContextRecovered } from '@/shared/state/agentsSlice';
|
||||
import { useClaudeTokens } from '@/shared/styles/ThemeContext';
|
||||
|
||||
// Muted, transient pill shown when the backend self-healed a context-overflow crash mid-turn (rebuilt the chat from its local copy and retried). Visible so the recovery isn't silent, calm so it doesn't read as an error; the "why" lives in the hover.
|
||||
export const ContextRecoveredPill: React.FC<{ sessionId: string }> = ({ 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 (
|
||||
<Fade in={!!cr} timeout={{ enter: 200, exit: 220 }} unmountOnExit>
|
||||
<Box
|
||||
title="This chat's memory overflowed mid-reply. OpenSwarm recovered it and retried automatically; nothing was lost."
|
||||
sx={{
|
||||
display: 'inline-flex',
|
||||
alignItems: 'center',
|
||||
gap: 0.6,
|
||||
alignSelf: 'flex-start',
|
||||
mx: 2,
|
||||
mb: 1,
|
||||
px: 1.25,
|
||||
py: 0.5,
|
||||
borderRadius: 999,
|
||||
bgcolor: c.bg.secondary,
|
||||
color: c.text.tertiary,
|
||||
}}
|
||||
>
|
||||
<RestartAltIcon sx={{ fontSize: 14 }} />
|
||||
<Typography sx={{ fontSize: '0.75rem', fontWeight: 500 }}>Recovered and retried</Typography>
|
||||
</Box>
|
||||
</Fade>
|
||||
);
|
||||
};
|
||||
@@ -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,
|
||||
|
||||
@@ -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') {
|
||||
|
||||
Reference in New Issue
Block a user