diff --git a/CHANGELOG.md b/CHANGELOG.md index 1c94c2130..3d55f45f3 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,6 +29,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Bug Fixes +* **transforms/content_router:** stop replacing `role="tool"` output with a lossy-unrecoverable summary on the live compression path (refs [#1307](https://github.com/chopratejas/headroom/issues/1307)). `ContentRouter.apply()` routed OpenAI-style `role="tool"` string messages — `Bash`/`grep`/`ls`/`cat` output — through the ML/word-drop summarizers; when the result carried no CCR retrieve marker (CCR off, ratio >= 0.8, or the size-gate fallback) the original was unrecoverable and the agent acted on a fabricated summary. Tool-role string content is now kept verbatim unless the compressed form is CCR-recoverable. Assistant/user text is unaffected, and structurally-lossless passes (SmartCrusher/Log/Search) still apply. The Anthropic `tool_result` block path is tracked separately. * **rtk:** stop `rtk` hook registration from spuriously timing out during `headroom wrap`. Output is captured to a temp file instead of pipes, and `stdin` is closed, so a background process forked by `rtk init` can no longer hold the pipe open and block `subprocess.run` past its 10s timeout after the hooks were already registered. * **ccr:** stop re-compressing `headroom_retrieve` output, which created an infinite retrieval loop, and stop emitting retrieval markers when the `headroom_retrieve` tool is not injected, which silently dropped data ([#1077](https://github.com/chopratejas/headroom/issues/1077), [#1006](https://github.com/chopratejas/headroom/issues/1006)). * **dashboard:** include RTK stats in the Historical tab; `/stats-history` now attaches live RTK/CLI-filtering stats the same way the Session tab does, so they survive a proxy restart ([#1177](https://github.com/chopratejas/headroom/issues/1177)). diff --git a/headroom/transforms/content_router.py b/headroom/transforms/content_router.py index 7105be997..53c03be75 100644 --- a/headroom/transforms/content_router.py +++ b/headroom/transforms/content_router.py @@ -54,6 +54,7 @@ from ..config import ( TransformResult, is_tool_excluded, ) +from ..parser import CCR_RETRIEVAL_MARKER_RE from ..tokenizer import Tokenizer from .base import Transform from .content_detector import ContentType, DetectionResult @@ -944,6 +945,17 @@ class ContentRouter(Transform): name: str = "content_router" + # Lossy summarizers that emit a CCR retrieve marker only when they store the + # original — a marker-less result from one of these is unrecoverable. Tool + # ground truth (role="tool") must not be replaced by such a result (#1307). + LOSSY_UNMARKED_STRATEGIES = frozenset( + { + CompressionStrategy.KOMPRESS, + CompressionStrategy.TEXT, + CompressionStrategy.CODE_AWARE, + } + ) + def __init__( self, config: ContentRouterConfig | None = None, @@ -2594,7 +2606,7 @@ class ContentRouter(Transform): netcost_p_alive_override = max(0.0, 1.0 - idle_f / ttl) # Tasks: list of (slot_index, content, context, bias, content_key) - _PendingTask = tuple[int, str, str, float, int] + _PendingTask = tuple[int, str, str, float, int, bool] pending_tasks: list[_PendingTask] = [] # #856 P2b (flag-gated, default off): net-cost frozen-floor unlock. @@ -2768,6 +2780,12 @@ class ContentRouter(Transform): # Key on the runtime target_ratio too: the same content compressed at # a different ratio is a different result, so it must not alias. content_key = hash((content, getattr(self, "_runtime_target_ratio", None))) + # Tool ground truth is gated against lossy-unrecoverable results below + # (#1307). Partition its cache namespace so a gated tool entry is never + # served from — or poisons — an ungated entry for byte-identical content. + enforce_reversibility = role == "tool" + if enforce_reversibility: + content_key = hash((content_key, True)) # Tier 1: skip set — instant rejection if self._cache.is_skipped(content_key): @@ -2816,7 +2834,9 @@ class ContentRouter(Transform): # Cache miss — defer to parallel compression pass route_counts.setdefault("cache_miss", 0) route_counts["cache_miss"] += 1 - pending_tasks.append((i, content, context, msg_bias, content_key)) + pending_tasks.append( + (i, content, context, msg_bias, content_key, enforce_reversibility) + ) # --- Pass 2: Parallel compression of all cache-miss messages --- if pending_tasks: @@ -2828,7 +2848,7 @@ class ContentRouter(Transform): if max_workers <= 1 or len(pending_tasks) == 1: # Single task or parallelism disabled — compress inline task_results = [] - for _, task_content, task_ctx, task_bias, _ in pending_tasks: + for _, task_content, task_ctx, task_bias, _, _ in pending_tasks: t0 = time.perf_counter() r = self.compress(task_content, context=task_ctx, bias=task_bias) task_results.append((r, (time.perf_counter() - t0) * 1000)) @@ -2836,7 +2856,7 @@ class ContentRouter(Transform): # Parallel compression via thread pool with ThreadPoolExecutor(max_workers=max_workers) as executor: futures = [] - for _, task_content, task_ctx, task_bias, _ in pending_tasks: + for _, task_content, task_ctx, task_bias, _, _ in pending_tasks: futures.append( executor.submit(self._timed_compress, task_content, task_ctx, task_bias) ) @@ -2846,9 +2866,10 @@ class ContentRouter(Transform): compressor_timing["parallel_compress_total"] = parallel_ms # --- Pass 3: Merge results back (sequential, updates caches) --- - for (slot_idx, task_content, _, _, content_key), (result, compress_ms) in zip( - pending_tasks, task_results - ): + for (slot_idx, task_content, _, _, content_key, enforce_rev), ( + result, + compress_ms, + ) in zip(pending_tasks, task_results): message = messages[slot_idx] strategy_key = f"compressor:{result.strategy_used.value}" compressor_timing[strategy_key] = ( @@ -2856,6 +2877,21 @@ class ContentRouter(Transform): ) if result.compression_ratio < min_ratio: + # tool ground truth must stay reversible — a lossy summarizer + # (kompress/text/code) that emitted no CCR retrieve marker is + # unrecoverable, so the agent would act on a fabricated summary + # (#1307). Keep the original verbatim instead. + if ( + enforce_rev + and result.strategy_used in self.LOSSY_UNMARKED_STRATEGIES + and not CCR_RETRIEVAL_MARKER_RE.search(result.compressed) + ): + self._cache.mark_skip(content_key) + result_slots[slot_idx] = message + route_counts["lossy_unrecoverable_skipped"] = ( + route_counts.get("lossy_unrecoverable_skipped", 0) + 1 + ) + continue # Compressed — store in result cache. The cache is still # warmed when the net-cost gate blocks the slot: the # gate's verdict is contextual (suffix size), the diff --git a/tests/test_canonical_pipeline.py b/tests/test_canonical_pipeline.py index 5608338e5..f15d29780 100644 --- a/tests/test_canonical_pipeline.py +++ b/tests/test_canonical_pipeline.py @@ -168,7 +168,9 @@ def test_content_router_protects_instruction_roles_but_compresses_tool_outputs() def fake_compress(text: str, **kwargs: Any) -> SimpleNamespace: calls.append(text) return SimpleNamespace( - compressed="COMPRESSED", + # CCR marker -> recoverable, so the #1307 gate keeps this lossy tool + # compression (tool output still compresses *because* it can be retrieved). + compressed="COMPRESSED <>", compression_ratio=0.1, strategy_used=CompressionStrategy.KOMPRESS, ) @@ -187,7 +189,7 @@ def test_content_router_protects_instruction_roles_but_compresses_tool_outputs() assert result.messages[0]["content"] == messages[0]["content"] assert result.messages[1]["content"] == messages[1]["content"] assert result.messages[2]["content"] == messages[2]["content"] - assert result.messages[3]["content"] == "COMPRESSED" + assert result.messages[3]["content"] == "COMPRESSED <>" assert calls == [tool_text] diff --git a/tests/test_content_router_tool_role_reversibility.py b/tests/test_content_router_tool_role_reversibility.py new file mode 100644 index 000000000..b5f9c4f5b --- /dev/null +++ b/tests/test_content_router_tool_role_reversibility.py @@ -0,0 +1,142 @@ +"""Tests that ContentRouter.apply() gates the role="tool" STRING path against +lossy-unrecoverable compression (#1307). + +This exercises the REAL proxy path: ContentRouter.apply() is what the pipeline +runs, and a role="tool" string message routes through Pass-1 -> pending_tasks -> +self.compress() (Pass-2) -> result merge (Pass-3). The fix gates that merge so a +lossy summarizer (kompress/text/code) that did not store the original (no CCR +retrieve marker) cannot replace verbatim tool output. + +The Kompress ML model is unavailable offline (it falls back to passthrough), so +the compression *result* is forced via monkeypatch — the seam is self.compress(), +the method apply() actually calls. apply() itself runs unmocked, so this proves +the live path, not an isolated unit (the gap PR #1363's apply()-direct tests had). +""" + +from __future__ import annotations + +from types import SimpleNamespace + +import pytest + +from headroom.transforms.content_router import ( + CompressionStrategy, + ContentRouter, +) + +# Realistic grep ground truth — comfortably > 50 word-tokens so it clears the +# min_tokens small-skip and reaches the compression path (otherwise apply() +# skips it before compress() is ever called). Every file:line is factual; a +# lossy reconstruction would fabricate paths the agent then acts on as fact. +GREP_OUTPUT = ( + 'headroom/transforms/content_router.py:1375: if role in ("tool", "assistant"):\n' + 'headroom/transforms/smart_crusher.py:1010: if msg.get("role") == "tool":\n' + 'headroom/transforms/content_router.py:2653: if role == "tool":\n' + 'headroom/proxy/handlers/anthropic.py:44: elif block.get("type") == "tool_result":\n' + 'headroom/cache/prefix_tracker.py:88: if message.get("role") == "tool":\n' + 'headroom/proxy/helpers.py:102: if msg.get("role") == "tool":\n' + "headroom/transforms/pipeline.py:132: transforms.append(ContentRouter())\n" + "headroom/transforms/kompress_compressor.py:1376: result = self.compress(content)\n" + 'headroom/transforms/content_router.py:1483: compressor_name = "KompressCompressor"\n' + 'headroom/transforms/content_router.py:1568: compressor_name = "KompressCompressor"\n' + "headroom/transforms/content_router.py:2667: bias = self._get_tool_bias(tool_name)\n" + "headroom/transforms/content_router.py:3317: result = self.compress(content, context=context)\n" + "headroom/transforms/content_router.py:3331: and not CCR_RETRIEVAL_MARKER_RE.search(result.compressed)\n" + "headroom/proxy/handlers/openai.py:697: def _compress_openai_responses_live_text_units(self)\n" + "headroom/transforms/compression_units.py:204: def compress_unit_with_router(self, unit)\n" + "headroom/config.py:676:class TransformResult: # messages, tokens_before, tokens_after\n" +) + +LOSSY_SUMMARY = "grep found 8 matches across config and proxy modules (kompressed)." +CCR_MARKER_SUMMARY = "grep matches (kompressed) <>" + +LOSSY_STRATEGIES = [ + CompressionStrategy.KOMPRESS, + CompressionStrategy.TEXT, + CompressionStrategy.CODE_AWARE, +] +STRUCTURED_STRATEGIES = [ + CompressionStrategy.SMART_CRUSHER, + CompressionStrategy.LOG, + CompressionStrategy.SEARCH, + CompressionStrategy.DIFF, +] + + +class _WordTokenizer: + """Word-count tokenizer stub — no model, deterministic, offline-safe.""" + + def count_text(self, text: object) -> int: + return len(str(text).split()) + + def count_messages(self, messages: list[dict]) -> int: + return sum(self.count_text(m.get("content", "")) for m in messages) + + +def _force_result(strategy: CompressionStrategy, compressed: str) -> SimpleNamespace: + """A RouterCompressionResult stand-in: apply() reads strategy_used, compressed, + and compression_ratio. ratio 0.3 < any min_ratio so it takes the 'compressed' + branch and reaches the reversibility gate.""" + return SimpleNamespace( + compressed=compressed, + original="", + strategy_used=strategy, + compression_ratio=0.3, + ) + + +def _tool_msg(content: str) -> dict: + # tool_call_id with no matching assistant tool_calls -> not in the exclude + # map -> not protected by the Read/Glob/Grep/Write/Edit window, so it reaches + # compression (matches Bash/shell output, which is never excluded). + return {"role": "tool", "tool_call_id": "call_bash_1", "content": content} + + +def _run(monkeypatch, message: dict, strategy: CompressionStrategy, compressed: str): + router = ContentRouter() + monkeypatch.setattr(router, "compress", lambda *a, **k: _force_result(strategy, compressed)) + # protect_recent / analysis protections are orthogonal to the reversibility + # gate and would preempt compression for recent code-like content. Disabling + # them isolates the gate and mirrors the real "aged-out tool output reaches + # compression" case that #1307 is about. + return router.apply( + [message], _WordTokenizer(), protect_recent=0, protect_analysis_context=False + ) + + +def test_tool_role_lossy_unmarked_kept_verbatim(monkeypatch) -> None: + """role=tool + lossy strategy + no CCR marker -> original preserved bit-for-bit.""" + result = _run(monkeypatch, _tool_msg(GREP_OUTPUT), CompressionStrategy.KOMPRESS, LOSSY_SUMMARY) + assert result.messages[0]["content"] == GREP_OUTPUT + + +def test_tool_role_lossy_with_ccr_marker_accepted(monkeypatch) -> None: + """role=tool + lossy strategy WITH a CCR marker -> compressed accepted (recoverable).""" + result = _run( + monkeypatch, _tool_msg(GREP_OUTPUT), CompressionStrategy.KOMPRESS, CCR_MARKER_SUMMARY + ) + assert result.messages[0]["content"] == CCR_MARKER_SUMMARY + + +def test_assistant_role_lossy_still_compressed(monkeypatch) -> None: + """Same lossy-unmarked result on role=assistant -> still compressed. + + Proves the gate is scoped to tool ground truth and does not regress + assistant-text compression effectiveness.""" + msg = {"role": "assistant", "content": GREP_OUTPUT} + result = _run(monkeypatch, msg, CompressionStrategy.KOMPRESS, LOSSY_SUMMARY) + assert result.messages[0]["content"] == LOSSY_SUMMARY + + +@pytest.mark.parametrize("strategy", LOSSY_STRATEGIES, ids=lambda s: s.value) +def test_tool_role_lossy_strategies_all_gated(monkeypatch, strategy) -> None: + """Every lossy-unmarked strategy is gated for tool role.""" + result = _run(monkeypatch, _tool_msg(GREP_OUTPUT), strategy, LOSSY_SUMMARY) + assert result.messages[0]["content"] == GREP_OUTPUT + + +@pytest.mark.parametrize("strategy", STRUCTURED_STRATEGIES, ids=lambda s: s.value) +def test_tool_role_structured_strategies_accepted(monkeypatch, strategy) -> None: + """Structured strategies are lossless/self-marking -> not gated, compressed kept.""" + result = _run(monkeypatch, _tool_msg(GREP_OUTPUT), strategy, LOSSY_SUMMARY) + assert result.messages[0]["content"] == LOSSY_SUMMARY diff --git a/tests/test_transforms_content_router.py b/tests/test_transforms_content_router.py index 9a60d7dc7..838eabbe4 100644 --- a/tests/test_transforms_content_router.py +++ b/tests/test_transforms_content_router.py @@ -354,7 +354,9 @@ def test_force_kompress_apply_uses_lightweight_detection( router, "compress", lambda content, context="", bias=1.0: RouterCompressionResult( - compressed="compressed", + # CCR marker -> the original was stored and is retrievable, so the + # #1307 reversibility gate accepts this lossy KOMPRESS tool result. + compressed="compressed <>", original=content, strategy_used=CompressionStrategy.KOMPRESS, routing_log=[ @@ -376,7 +378,7 @@ def test_force_kompress_apply_uses_lightweight_detection( protect_recent=2, ) - assert result.messages[0]["content"] == "compressed" + assert result.messages[0]["content"] == "compressed <>" def test_force_kompress_apply_lightweight_detection_protects_recent_code(