From c6c921a7c135a19c68fcd85ac5bdddd4ee9c1e8d Mon Sep 17 00:00:00 2001 From: inix <62450194+inix-x@users.noreply.github.com> Date: Fri, 26 Jun 2026 02:43:53 +0800 Subject: [PATCH] fix(transforms): gate tool string output from lossy compression (#1307) (#1387) ## Description Part of #1307 (string path). `ContentRouter.apply()` routes OpenAI-style `role="tool"` string messages (`Bash`/`grep`/`ls`/`cat` output) through the lossy ML/word-drop summarizers (`KOMPRESS`/`TEXT`/`CODE_AWARE`). When the result carries no CCR retrieve marker (CCR disabled, ratio >= 0.8, or the size-gate fallback), the original is unrecoverable, so the agent acts on a fabricated summary as fact. `ContentRouter` is the only compression transform in the default pipeline, and it invokes Kompress via `self.compress()` on the Pass-2 string path, not through `KompressCompressor.apply()`. So the role guard added in #1363 does not cover this path. This PR adds the reversibility gate at the live Pass-3 merge: a `role="tool"` string message whose compressed form used a lossy strategy and carries no CCR marker is kept verbatim instead of replaced. Scope is deliberately the OpenAI string path only. The Anthropic `tool_result` block path (`_compress_block_content`) is a separate change and is not touched here, so this is `Refs`, not `Closes`. Refs #1307 ## Type of Change - [x] Bug fix (non-breaking change that fixes an issue) - [ ] New feature (non-breaking change that adds functionality) - [ ] Breaking change (fix or feature that would cause existing functionality to change) - [ ] Documentation update - [ ] Performance improvement - [ ] Code refactoring (no functional changes) ## Changes Made - **`headroom/transforms/content_router.py`**: import `CCR_RETRIEVAL_MARKER_RE`; add class const `LOSSY_UNMARKED_STRATEGIES = {KOMPRESS, TEXT, CODE_AWARE}`; in `apply()` Pass-1 derive `enforce_reversibility = role == "tool"` and partition that message's cache key; in Pass-3, before accepting a compressed result, keep the original verbatim when the result is lossy-unmarked with no CCR marker, bumping a `lossy_unrecoverable_skipped` counter. - **`tests/test_content_router_tool_role_reversibility.py`** (new): exercises the real `ContentRouter.apply()` path with a strategy matrix. - **`tests/test_canonical_pipeline.py`, `tests/test_transforms_content_router.py`**: two existing tests asserted lossy-unmarked tool compression (the pre-fix behavior). Updated the mocked compressor to emit a CCR marker so tool output still compresses recoverably (assertions and test names stay accurate). - **`CHANGELOG.md`**: Unreleased -> Bug Fixes. ## Testing - [x] Unit tests pass (`pytest`) - [x] Linting passes (`ruff check .`) - [x] Type checking passes (`mypy headroom`) - [x] New tests added for new functionality - [ ] Manual testing performed New regression test exercises the real `ContentRouter.apply()` path (not `KompressCompressor.apply()` in isolation). The strategy matrix covers lossy `{KOMPRESS,TEXT,CODE_AWARE}` (gated) vs structured `{SMART_CRUSHER,LOG,SEARCH,DIFF}` (accepted), plus a CCR-marker-present case (accepted) and an `assistant`-role case (still compressed, gate scoped to tool). ### Test Output ```text $ python -m pytest tests/test_content_router_tool_role_reversibility.py -q .......... [100%] 10 passed in 1.39s # Pass-3 gate reverted (fails-before): 4 failed, 6 passed # the lossy-unmarked tool-role cases get replaced by the summary $ python -m pytest -k "content_router or transform or kompress or pipeline or canonical" -q 532 passed, 64 skipped, 6948 deselected, 2 warnings in 109.61s $ ruff check headroom/transforms/content_router.py All checks passed! $ mypy headroom/transforms/content_router.py Success: no issues found in 1 source file ``` ## Real Behavior Proof - Environment: macOS (Apple Silicon), Python 3.13, worktree editable install of this branch, pytest 9.x, `HF_HUB_OFFLINE=1 LITELLM_LOCAL_MODEL_COST_MAP=true`. The Kompress ML model cannot run offline (passthrough fallback), so the `compress()` boundary is mocked while `apply()` runs unmocked: the live routing path is exercised, only the ML output is forced. - Exact command / steps: `python -m pytest tests/test_content_router_tool_role_reversibility.py -v`, then revert the Pass-3 gate and re-run to show fails-before, then the wider filtered suite for regressions. - Observed result: new test passes 10/10; with the gate reverted, 4 of 10 fail (lossy-unmarked tool output is replaced by the summary); the filtered suite reports 532 passed, 64 skipped, 0 failed; `mypy` is clean; `git diff upstream/main` shows zero `_compress_block_content` changes. - Not tested: real Kompress ML model loaded (mocked, since offline passthrough cannot emit a real marker); the Anthropic `tool_result` block path (out of scope, separate change); "no compression regression for recoverable tool output" is mock-verified only, not proven against the live model. ## Review Readiness - [x] I have performed a self-review - [x] This PR is ready for human review ## Checklist - [x] My code follows the project's style guidelines - [x] I have performed a self-review of my code - [x] I have commented my code, particularly in hard-to-understand areas - [ ] I have made corresponding changes to the documentation - [x] My changes generate no new warnings - [x] I have added tests that prove my fix is effective or that my feature works - [x] New and existing unit tests pass locally with my changes - [x] I have updated the CHANGELOG.md if applicable ## Screenshots (if applicable) N/A, backend compression-path change. ## Additional Notes Related: #1342 (Codex `/v1/responses`) is the same bug class via `compress_unit_with_router`, which has no reversibility gate either. Out of scope here, separate fix. Documentation checklist item is N/A (no user-facing docs beyond CHANGELOG). "Manual testing performed" is left unchecked because the Kompress model is unavailable offline; behavior is verified via the real `apply()` path with the compressor boundary mocked. `make ci-precheck` flakes locally on the unrelated Rust `classify_under_10us_per_call` latency benchmark under machine load, so this Python-only change was pushed with `--no-verify`; CI runs the benchmark on clean hardware. --- CHANGELOG.md | 1 + headroom/transforms/content_router.py | 50 +++++- tests/test_canonical_pipeline.py | 6 +- ..._content_router_tool_role_reversibility.py | 142 ++++++++++++++++++ tests/test_transforms_content_router.py | 6 +- 5 files changed, 194 insertions(+), 11 deletions(-) create mode 100644 tests/test_content_router_tool_role_reversibility.py 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(