headroom/tests/test_canonical_pipeline.py
inix c6c921a7c1
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.
2026-06-25 13:43:53 -05:00

321 lines
9.9 KiB
Python

from __future__ import annotations
import importlib
from types import SimpleNamespace
from typing import Any
from headroom.client import HeadroomClient
from headroom.compress import compress
from headroom.config import HeadroomConfig, HeadroomMode, TransformResult
from headroom.hooks import CompressionHooks
from headroom.pipeline import (
CANONICAL_PIPELINE_STAGES,
PipelineExtensionManager,
PipelineStage,
summarize_routing_markers,
)
from headroom.providers.base import Provider, TokenCounter
from headroom.transforms import ContentRouter, TransformPipeline
from headroom.transforms.content_router import CompressionStrategy
class RecordingExtension:
def __init__(self) -> None:
self.stages: list[PipelineStage] = []
def on_pipeline_event(self, event):
self.stages.append(event.stage)
return None
class MutatingExtension:
def on_pipeline_event(self, event):
if event.stage == PipelineStage.INPUT_RECEIVED:
event.messages = [{"role": "user", "content": "mutated"}]
return event
class ReplacingExtension:
def on_pipeline_event(self, event):
return type(event)(
stage=event.stage,
operation=event.operation,
model=event.model,
messages=[{"role": "user", "content": "replaced"}],
metadata={"replaced": True},
)
class RaisingExtension:
def on_pipeline_event(self, event):
raise RuntimeError("boom")
class RecordingHooks(CompressionHooks):
def __init__(self) -> None:
self.stages: list[PipelineStage] = []
self.post_event = None
def pre_compress(self, messages, ctx):
return messages
def compute_biases(self, messages, ctx):
return {}
def post_compress(self, event):
self.post_event = event
def on_pipeline_event(self, event):
self.stages.append(event.stage)
return None
class StubPipeline:
def apply(self, messages, model, **kwargs):
return TransformResult(
messages=messages,
tokens_before=20,
tokens_after=8,
transforms_applied=["router:text:kompress", "kompress:user:0.40"],
)
def _get_tokenizer(self, model):
return StubTokenCounter()
class StubTokenCounter(TokenCounter):
def count_text(self, text: str) -> int:
return len(text.split())
def count_message(self, message: dict[str, Any]) -> int:
content = message.get("content", "")
if isinstance(content, str):
return len(content.split())
return 1
def count_messages(self, messages: list[dict[str, Any]]) -> int:
return sum(self.count_message(message) for message in messages)
class StubProvider(Provider):
@property
def name(self) -> str:
return "openai"
def get_token_counter(self, model: str) -> TokenCounter:
return StubTokenCounter()
def get_context_limit(self, model: str) -> int:
return 128000
def supports_model(self, model: str) -> bool:
return True
class DummyCompletions:
def __init__(self) -> None:
self.calls: list[dict[str, Any]] = []
def create(self, **kwargs: Any) -> dict[str, Any]:
self.calls.append(kwargs)
return {"id": "resp_123", "messages": kwargs["messages"]}
class DummyOriginalClient:
def __init__(self) -> None:
self.chat = SimpleNamespace(completions=DummyCompletions())
def test_pipeline_extension_manager_uses_canonical_stage_contract():
recorder = RecordingExtension()
manager = PipelineExtensionManager(
extensions=[recorder, MutatingExtension()],
discover=False,
)
event = manager.emit(
PipelineStage.INPUT_RECEIVED,
operation="test",
model="gpt-4o",
messages=[{"role": "user", "content": "hello"}],
)
assert list(CANONICAL_PIPELINE_STAGES)[0] is PipelineStage.SETUP
assert summarize_routing_markers(["router:text:kompress", "smart:kept=3"]) == [
"router:text:kompress"
]
assert recorder.stages == [PipelineStage.INPUT_RECEIVED]
assert event.messages == [{"role": "user", "content": "mutated"}]
def test_default_transform_pipeline_always_uses_content_router() -> None:
config = HeadroomConfig()
pipeline = TransformPipeline(config)
assert any(isinstance(transform, ContentRouter) for transform in pipeline.transforms)
assert not any(type(transform).__name__ == "SmartCrusher" for transform in pipeline.transforms)
def test_content_router_protects_instruction_roles_but_compresses_tool_outputs() -> None:
class Tokenizer:
def count_text(self, text: str) -> int:
return max(1, len(text.split()))
router = ContentRouter()
calls: list[str] = []
def fake_compress(text: str, **kwargs: Any) -> SimpleNamespace:
calls.append(text)
return SimpleNamespace(
# CCR marker -> recoverable, so the #1307 gate keeps this lossy tool
# compression (tool output still compresses *because* it can be retrieved).
compressed="COMPRESSED <<ccr:tool>>",
compression_ratio=0.1,
strategy_used=CompressionStrategy.KOMPRESS,
)
router.compress = fake_compress # type: ignore[method-assign]
tool_text = "tool output " * 120
messages = [
{"role": "system", "content": "system instructions " * 120},
{"role": "developer", "content": "developer instructions " * 120},
{"role": "user", "content": "user prompt " * 120},
{"role": "tool", "tool_call_id": "call_1", "content": tool_text},
]
result = router.apply(messages, Tokenizer())
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 <<ccr:tool>>"
assert calls == [tool_text]
def test_pipeline_extension_manager_replaces_events_and_ignores_failures(caplog):
recorder = RecordingExtension()
manager = PipelineExtensionManager(
extensions=[recorder, RaisingExtension(), ReplacingExtension(), object()],
discover=False,
)
with caplog.at_level("WARNING", logger="headroom.pipeline"):
event = manager.emit(
PipelineStage.PRE_SEND,
operation="test",
model="gpt-4o",
messages=[{"role": "user", "content": "hello"}],
)
assert manager.enabled is True
assert recorder.stages == [PipelineStage.PRE_SEND]
assert event.messages == [{"role": "user", "content": "replaced"}]
assert event.metadata == {"replaced": True}
def test_discover_pipeline_extensions_handles_load_and_init_failures(monkeypatch):
pipeline_module = importlib.import_module("headroom.pipeline")
class Entry:
def __init__(self, name, loader):
self.name = name
self._loader = loader
def load(self):
return self._loader()
class ExtensionClass:
def on_pipeline_event(self, event):
return event
class FailingInit:
def __init__(self):
raise RuntimeError("init failed")
entries = [
Entry("instance", lambda: RecordingExtension()),
Entry("class", lambda: ExtensionClass),
Entry("load-fail", lambda: (_ for _ in ()).throw(RuntimeError("load failed"))),
Entry("init-fail", lambda: FailingInit),
]
monkeypatch.setattr(
pipeline_module.importlib.metadata,
"entry_points",
lambda group: entries if group == pipeline_module.ENTRY_POINT_GROUP else [],
)
discovered = pipeline_module.discover_pipeline_extensions()
assert len(discovered) == 2
assert hasattr(discovered[0], "on_pipeline_event")
assert hasattr(discovered[1], "on_pipeline_event")
def test_discover_pipeline_extensions_returns_empty_when_entrypoint_lookup_fails(monkeypatch):
pipeline_module = importlib.import_module("headroom.pipeline")
monkeypatch.setattr(
pipeline_module.importlib.metadata,
"entry_points",
lambda group: (_ for _ in ()).throw(RuntimeError("lookup failed")),
)
assert pipeline_module.discover_pipeline_extensions() == []
def test_compress_emits_canonical_pipeline_events(monkeypatch):
hooks = RecordingHooks()
compress_module = importlib.import_module("headroom.compress")
monkeypatch.setattr(compress_module, "_get_pipeline", lambda: StubPipeline())
result = compress(
[{"role": "user", "content": "hello world"}],
model="gpt-4o",
hooks=hooks,
)
assert result.tokens_before == 20
assert result.tokens_after == 8
assert hooks.post_event is not None
assert hooks.post_event.tokens_saved == 12
assert hooks.stages == [
PipelineStage.INPUT_RECEIVED,
PipelineStage.INPUT_ROUTED,
PipelineStage.INPUT_COMPRESSED,
]
def test_headroom_client_emits_canonical_pipeline_events(tmp_path):
recorder = RecordingExtension()
original = DummyOriginalClient()
config = HeadroomConfig(
store_url=f"jsonl://{tmp_path / 'headroom.jsonl'}",
default_mode=HeadroomMode.OPTIMIZE,
pipeline_extensions=[recorder],
discover_pipeline_extensions=False,
)
client = HeadroomClient(
original_client=original,
provider=StubProvider(),
store_url=f"jsonl://{tmp_path / 'headroom-client.jsonl'}",
enable_cache_optimizer=False,
config=config,
)
client._pipeline = StubPipeline()
response = client.chat.completions.create(
model="gpt-4o",
messages=[{"role": "user", "content": "hello world"}],
)
assert response["id"] == "resp_123"
assert recorder.stages == [
PipelineStage.SETUP,
PipelineStage.INPUT_RECEIVED,
PipelineStage.INPUT_ROUTED,
PipelineStage.INPUT_COMPRESSED,
PipelineStage.PRE_SEND,
PipelineStage.POST_SEND,
PipelineStage.RESPONSE_RECEIVED,
]