mirror of
https://github.com/headroomlabs-ai/headroom.git
synced 2026-08-27 14:17:10 -04:00
## 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.
321 lines
9.9 KiB
Python
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,
|
|
]
|