From 1498a87806d5a3a7dc159b2eb81d9283dc21cd99 Mon Sep 17 00:00:00 2001 From: tjb-tech <1193992557@qq.com> Date: Tue, 28 Apr 2026 13:28:01 +0000 Subject: [PATCH] fix(compact): bound large tool result history --- CHANGELOG.md | 1 + src/openharness/engine/query.py | 55 +++++++++- src/openharness/services/compact/__init__.py | 35 +++++- src/openharness/services/tool_outputs.py | 55 ++++++++++ tests/test_engine/test_query_engine.py | 73 +++++++++++++ tests/test_services/test_compact.py | 109 +++++++++++++++++++ 6 files changed, 321 insertions(+), 7 deletions(-) create mode 100644 src/openharness/services/tool_outputs.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 447a676..855457b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,6 +28,7 @@ The format is based on Keep a Changelog, and this project currently tracks chang ### Fixed - Compaction now detects llama.cpp/OpenAI-compatible context overflow errors, accounts for image blocks in auto-compact token estimates, and strips image payloads from summarizer-only compaction requests. +- Large tool results are now bounded in conversation history: oversized outputs are saved under `tool_artifacts`, old MCP results become microcompactable, and context collapse trims stale tool-result payloads. - Fixed `glob` and `grep` tools hanging indefinitely when the `rg` subprocess produced enough stderr output to fill the OS pipe buffer. `stderr` is now redirected to `DEVNULL` so it is discarded rather than blocking the child process. - Fixed `bash_tool` hanging after a timed-out command when the subprocess stdout stream stayed open. `_read_remaining_output` now applies a 2-second `asyncio.wait_for` timeout so the tool always returns promptly. - Fixed `session_runner` background task deadlock caused by an unread `stderr=PIPE` stream. The subprocess now uses `stderr=STDOUT` so all output merges into the single readable stdout pipe. diff --git a/src/openharness/engine/query.py b/src/openharness/engine/query.py index 9ce09a3..dee3c70 100644 --- a/src/openharness/engine/query.py +++ b/src/openharness/engine/query.py @@ -9,6 +9,7 @@ import time from dataclasses import dataclass from pathlib import Path from typing import Any, AsyncIterator, Awaitable, Callable +from uuid import uuid4 from openharness.api.client import ( ApiMessageCompleteEvent, @@ -18,6 +19,7 @@ from openharness.api.client import ( SupportsStreamingMessages, ) from openharness.api.usage import UsageSnapshot +from openharness.config.paths import get_data_dir from openharness.engine.messages import ConversationMessage, ToolResultBlock from openharness.engine.stream_events import ( AssistantTextDelta, @@ -31,6 +33,7 @@ from openharness.engine.stream_events import ( ) from openharness.hooks import HookEvent, HookExecutor from openharness.permissions.checker import PermissionChecker +from openharness.services.tool_outputs import tool_output_inline_chars, tool_output_preview_chars from openharness.tools.base import ToolExecutionContext from openharness.tools.base import ToolRegistry @@ -455,6 +458,49 @@ def _record_tool_carryover( _remember_work_log(context.tool_metadata, entry="Exited plan mode") +def _tool_artifact_dir() -> Path: + artifact_dir = get_data_dir() / "tool_artifacts" + artifact_dir.mkdir(parents=True, exist_ok=True) + return artifact_dir + + +def _safe_tool_artifact_name(tool_name: str) -> str: + normalized = re.sub(r"[^A-Za-z0-9_.-]+", "_", tool_name.strip()) + return (normalized or "tool")[:80] + + +def _offload_tool_output_if_needed( + *, + tool_name: str, + tool_use_id: str, + output: str, +) -> tuple[str, Path | None]: + inline_limit = tool_output_inline_chars() + if len(output) <= inline_limit: + return output, None + + artifact_path = ( + _tool_artifact_dir() + / f"{time.strftime('%Y%m%d-%H%M%S')}-{_safe_tool_artifact_name(tool_name)}-{uuid4().hex[:12]}.txt" + ) + artifact_path.write_text(output, encoding="utf-8", errors="replace") + preview = output[:tool_output_preview_chars()] + omitted = max(0, len(output) - len(preview)) + inline = ( + "[Tool output truncated]\n" + f"Tool: {tool_name}\n" + f"Tool use id: {tool_use_id}\n" + f"Original size: {len(output)} chars\n" + f"Full output saved to: {artifact_path}\n" + f"Inline preview: first {len(preview)} chars" + ) + if omitted: + inline += f" ({omitted} chars omitted)" + if preview: + inline += f"\n\nPreview:\n{preview}" + return inline, artifact_path + + async def run_query( context: QueryContext, messages: list[ConversationMessage], @@ -760,9 +806,16 @@ async def _execute_tool_call( elapsed = time.monotonic() - t0 log.debug("executed %s in %.2fs err=%s output_len=%d", tool_name, elapsed, result.is_error, len(result.output or "")) + inline_output, artifact_path = _offload_tool_output_if_needed( + tool_name=tool_name, + tool_use_id=tool_use_id, + output=result.output, + ) + if artifact_path is not None: + _remember_active_artifact(context.tool_metadata, str(artifact_path)) tool_result = ToolResultBlock( tool_use_id=tool_use_id, - content=result.output, + content=inline_output, is_error=result.is_error, ) _record_tool_carryover( diff --git a/src/openharness/services/compact/__init__.py b/src/openharness/services/compact/__init__.py index 96286aa..0aa540e 100644 --- a/src/openharness/services/compact/__init__.py +++ b/src/openharness/services/compact/__init__.py @@ -29,6 +29,7 @@ from openharness.engine.messages import ( ) from openharness.engine.stream_events import CompactProgressEvent from openharness.hooks import HookEvent, HookExecutor +from openharness.services.tool_outputs import is_microcompactable_tool_result from openharness.services.token_estimation import estimate_tokens log = logging.getLogger(__name__) @@ -313,6 +314,17 @@ def try_context_collapse( if collapsed != block.text: changed = True new_blocks.append(TextBlock(text=collapsed)) + elif isinstance(block, ToolResultBlock): + collapsed = _collapse_text(block.content) + if collapsed != block.content: + changed = True + new_blocks.append( + ToolResultBlock( + tool_use_id=block.tool_use_id, + content=collapsed, + is_error=block.is_error, + ) + ) else: new_blocks.append(block) collapsed_older.append(ConversationMessage(role=message.role, content=new_blocks)) @@ -767,14 +779,25 @@ def _build_passthrough_compaction_result( def _collect_compactable_tool_ids(messages: list[ConversationMessage]) -> list[str]: """Walk messages and collect tool_use IDs whose results are compactable.""" - ids: list[str] = [] + ordered_ids: list[str] = [] + tool_names: dict[str, str] = {} + result_content: dict[str, str] = {} for msg in messages: - if msg.role != "assistant": - continue for block in msg.content: - if isinstance(block, ToolUseBlock) and block.name in COMPACTABLE_TOOLS: - ids.append(block.id) - return ids + if isinstance(block, ToolUseBlock): + ordered_ids.append(block.id) + tool_names[block.id] = block.name + elif isinstance(block, ToolResultBlock): + result_content[block.tool_use_id] = block.content + return [ + tool_id + for tool_id in ordered_ids + if tool_names.get(tool_id, "") in COMPACTABLE_TOOLS + or is_microcompactable_tool_result( + tool_names.get(tool_id, ""), + result_content.get(tool_id, ""), + ) + ] def microcompact_messages( diff --git a/src/openharness/services/tool_outputs.py b/src/openharness/services/tool_outputs.py new file mode 100644 index 0000000..e4357e3 --- /dev/null +++ b/src/openharness/services/tool_outputs.py @@ -0,0 +1,55 @@ +"""Tool-output context budget helpers.""" + +from __future__ import annotations + +import logging +import os + +log = logging.getLogger(__name__) + +DEFAULT_TOOL_OUTPUT_INLINE_CHARS = 16_000 +DEFAULT_TOOL_OUTPUT_PREVIEW_CHARS = 3_000 +DEFAULT_MICROCOMPACT_TOOL_RESULT_CHARS = 4_000 + + +def _read_positive_int_env(name: str, default: int, *, minimum: int = 1) -> int: + raw = os.environ.get(name, "").strip() + if not raw: + return default + try: + return max(minimum, int(raw)) + except ValueError: + log.warning("Ignoring invalid %s=%r", name, raw) + return default + + +def tool_output_inline_chars() -> int: + return _read_positive_int_env( + "OPENHARNESS_TOOL_OUTPUT_INLINE_CHARS", + DEFAULT_TOOL_OUTPUT_INLINE_CHARS, + minimum=256, + ) + + +def tool_output_preview_chars() -> int: + return _read_positive_int_env( + "OPENHARNESS_TOOL_OUTPUT_PREVIEW_CHARS", + DEFAULT_TOOL_OUTPUT_PREVIEW_CHARS, + minimum=128, + ) + + +def microcompact_tool_result_chars() -> int: + return _read_positive_int_env( + "OPENHARNESS_MICROCOMPACT_TOOL_RESULT_CHARS", + DEFAULT_MICROCOMPACT_TOOL_RESULT_CHARS, + minimum=256, + ) + + +def is_microcompactable_tool_result(tool_name: str, content: str) -> bool: + """Return True when a tool result should be eligible for old-result clearing.""" + normalized = tool_name.strip() + if normalized.startswith("mcp__"): + return True + return len(content) >= microcompact_tool_result_chars() diff --git a/tests/test_engine/test_query_engine.py b/tests/test_engine/test_query_engine.py index c81fd95..5071116 100644 --- a/tests/test_engine/test_query_engine.py +++ b/tests/test_engine/test_query_engine.py @@ -1113,6 +1113,19 @@ class _BoomTool(BaseTool): raise RuntimeError("boom") +class _LargeOutputTool(BaseTool): + name = "mcp__playwright__browser_snapshot" + description = "Returns a large browser snapshot." + input_model = _OkInput + + def is_read_only(self, arguments: BaseModel) -> bool: + return True + + async def execute(self, arguments: BaseModel, context: ToolExecutionContext) -> ToolResult: + del arguments, context + return ToolResult(output="snapshot-line\n" * 40) + + @pytest.mark.asyncio async def test_query_engine_synthesizes_tool_result_when_parallel_tool_raises(tmp_path: Path): """Parallel tool calls must each yield a tool_result even when one tool raises. @@ -1179,6 +1192,66 @@ async def test_query_engine_synthesizes_tool_result_when_parallel_tool_raises(tm assert events[-1].message.text == "Recovered from the failure." +@pytest.mark.asyncio +async def test_query_engine_offloads_large_tool_result_outputs(tmp_path: Path, monkeypatch): + monkeypatch.setenv("OPENHARNESS_DATA_DIR", str(tmp_path / "data")) + monkeypatch.setenv("OPENHARNESS_TOOL_OUTPUT_INLINE_CHARS", "256") + monkeypatch.setenv("OPENHARNESS_TOOL_OUTPUT_PREVIEW_CHARS", "128") + registry = ToolRegistry() + registry.register(_LargeOutputTool()) + + engine = QueryEngine( + api_client=FakeApiClient( + [ + _FakeResponse( + message=ConversationMessage( + role="assistant", + content=[ + ToolUseBlock( + id="toolu_snapshot", + name="mcp__playwright__browser_snapshot", + input={}, + ), + ], + ), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + _FakeResponse( + message=ConversationMessage(role="assistant", content=[TextBlock(text="done")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + ] + ), + tool_registry=registry, + permission_checker=PermissionChecker(PermissionSettings(mode=PermissionMode.FULL_AUTO)), + cwd=tmp_path, + model="claude-test", + system_prompt="system", + tool_metadata={}, + ) + + events = [event async for event in engine.submit_message("snapshot")] + + completed = [event for event in events if isinstance(event, ToolExecutionCompleted)] + assert len(completed) == 1 + assert completed[0].output.startswith("[Tool output truncated]") + assert "snapshot-line" in completed[0].output + + user_tool_messages = [ + msg for msg in engine.messages if msg.role == "user" and any(isinstance(block, ToolResultBlock) for block in msg.content) + ] + result_blocks = [block for block in user_tool_messages[0].content if isinstance(block, ToolResultBlock)] + inline = result_blocks[0].content + assert "Full output saved to:" in inline + assert "Original size:" in inline + assert inline.count("snapshot-line") < 40 + artifact_line = next(line for line in inline.splitlines() if line.startswith("Full output saved to:")) + artifact_path = Path(artifact_line.removeprefix("Full output saved to:").strip()) + assert artifact_path.exists() + assert artifact_path.read_text(encoding="utf-8") == "snapshot-line\n" * 40 + assert str(artifact_path) in engine.tool_metadata["task_focus_state"]["active_artifacts"] + + @pytest.mark.asyncio async def test_query_engine_drops_empty_assistant_messages(tmp_path: Path): engine = QueryEngine( diff --git a/tests/test_services/test_compact.py b/tests/test_services/test_compact.py index ec519c2..deeedc0 100644 --- a/tests/test_services/test_compact.py +++ b/tests/test_services/test_compact.py @@ -25,6 +25,7 @@ from openharness.services.compact import ( auto_compact_if_needed, estimate_message_tokens as estimate_compact_message_tokens, get_autocompact_threshold, + microcompact_messages, should_autocompact, try_context_collapse, try_session_memory_compaction, @@ -168,6 +169,114 @@ def test_try_context_collapse_trims_oversized_messages(): assert "[collapsed" in result[0].text +def test_try_context_collapse_trims_oversized_tool_results(): + giant = ("snapshot node " * 1200).strip() + messages = [ + ConversationMessage.from_user_text("open page"), + ConversationMessage( + role="assistant", + content=[ToolUseBlock(id="toolu_snapshot", name="mcp__playwright__browser_snapshot", input={})], + ), + ConversationMessage( + role="user", + content=[ToolResultBlock(tool_use_id="toolu_snapshot", content=giant, is_error=False)], + ), + ConversationMessage(role="assistant", content=[TextBlock(text="I inspected the snapshot")]), + ConversationMessage(role="user", content=[TextBlock(text="latest")]), + ConversationMessage(role="assistant", content=[TextBlock(text="keep recent")]), + ] + + result = try_context_collapse(messages, preserve_recent=2) + + assert result is not None + collapsed_results = [ + block + for message in result + for block in message.content + if isinstance(block, ToolResultBlock) + ] + assert len(collapsed_results) == 1 + assert "[collapsed" in collapsed_results[0].content + assert collapsed_results[0].tool_use_id == "toolu_snapshot" + + +def test_microcompact_compacts_mcp_results_while_preserving_recent(): + messages = [] + for index in range(3): + tool_id = f"toolu_snapshot_{index}" + messages.extend( + [ + ConversationMessage( + role="assistant", + content=[ + ToolUseBlock( + id=tool_id, + name="mcp__playwright__browser_snapshot", + input={}, + ) + ], + ), + ConversationMessage( + role="user", + content=[ + ToolResultBlock( + tool_use_id=tool_id, + content=f"snapshot {index} " * 600, + is_error=False, + ) + ], + ), + ] + ) + + compacted, tokens_saved = microcompact_messages(messages, keep_recent=1) + + assert tokens_saved > 0 + results = [ + block + for message in compacted + for block in message.content + if isinstance(block, ToolResultBlock) + ] + assert results[0].content == "[Old tool result content cleared]" + assert results[1].content == "[Old tool result content cleared]" + assert results[2].content.startswith("snapshot 2") + + +def test_microcompact_compacts_large_non_allowlisted_results(monkeypatch): + monkeypatch.setenv("OPENHARNESS_MICROCOMPACT_TOOL_RESULT_CHARS", "256") + messages = [ + ConversationMessage( + role="assistant", + content=[ToolUseBlock(id="toolu_custom_0", name="custom_snapshot_tool", input={})], + ), + ConversationMessage( + role="user", + content=[ToolResultBlock(tool_use_id="toolu_custom_0", content="A" * 512, is_error=False)], + ), + ConversationMessage( + role="assistant", + content=[ToolUseBlock(id="toolu_custom_1", name="custom_snapshot_tool", input={})], + ), + ConversationMessage( + role="user", + content=[ToolResultBlock(tool_use_id="toolu_custom_1", content="B" * 512, is_error=False)], + ), + ] + + compacted, tokens_saved = microcompact_messages(messages, keep_recent=1) + + assert tokens_saved > 0 + results = [ + block + for message in compacted + for block in message.content + if isinstance(block, ToolResultBlock) + ] + assert results[0].content == "[Old tool result content cleared]" + assert results[1].content == "B" * 512 + + def test_compact_prompt_too_long_detection_handles_llama_cpp_errors(): assert _is_prompt_too_long_error( RuntimeError("exceed_context_size_error: prompt exceeds the available context size")