Skip to content

Commit 60af15f

Browse files
chopratejasclaude
andauthored
feat(content-router): lossless-first dispatch, cross-turn dedup, and A7 lossy-after-fold (headroomlabs-ai#1818)
## Description <!-- Briefly explain the change and why it is needed. --> Closes # ## Type of Change - [ ] 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 - ## Testing <!-- Check what you actually ran, then paste the real command output below. --> - [ ] Unit tests pass (`pytest`) - [ ] Linting passes (`ruff check .`) - [ ] Type checking passes (`mypy headroom`) - [ ] New tests added for new functionality - [ ] Manual testing performed ### Test Output ```text # Paste relevant command output or artifact links here ``` ## Real Behavior Proof - Environment: - Exact command / steps: - Observed result: - Not tested: ## Review Readiness - [ ] I have performed a self-review - [ ] This PR is ready for human review ## Checklist - [ ] My code follows the project's style guidelines - [ ] I have performed a self-review of my code - [ ] I have commented my code, particularly in hard-to-understand areas - [ ] I have made corresponding changes to the documentation - [ ] My changes generate no new warnings - [ ] I have added tests that prove my fix is effective or that my feature works - [ ] New and existing unit tests pass locally with my changes - [ ] I have updated the CHANGELOG.md if applicable ## Screenshots (if applicable) Add screenshots to help explain your changes. ## Additional Notes <!-- Mention any N/A checklist items, tradeoffs, follow-ups, or maintainer context. --> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
1 parent da2d8dc commit 60af15f

13 files changed

Lines changed: 1205 additions & 87 deletions

‎headroom/cli/proxy.py‎

Lines changed: 13 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -291,21 +291,16 @@ def dashboard(port: int, no_open: bool) -> None:
291291
help="Max tokens per minute. Env: HEADROOM_TPM. Default: 100000.",
292292
)
293293
@click.option(
294-
"--no-ccr-inject-tool",
294+
"--no-ccr",
295295
is_flag=True,
296-
envvar="HEADROOM_NO_CCR_INJECT_TOOL",
296+
envvar="HEADROOM_NO_CCR",
297297
help=(
298-
"Don't inject the CCR headroom_retrieve tool. Run compression-only — "
299-
"for streaming / non-MCP clients that can't resolve the retrieve tool "
300-
"and would otherwise error on it. Env: HEADROOM_NO_CCR_INJECT_TOOL."
298+
"Disable CCR entirely: no retrieval markers in compressed content AND no "
299+
"headroom_retrieve tool injected. Lossy compression with no recovery path "
300+
"(maximum savings; also right for streaming / non-MCP clients that can't "
301+
"resolve an injected tool). Env: HEADROOM_NO_CCR."
301302
),
302303
)
303-
@click.option(
304-
"--no-ccr-marker",
305-
is_flag=True,
306-
envvar="HEADROOM_NO_CCR_MARKER",
307-
help=("Don't add CCR retrieval markers to compressed content. Env: HEADROOM_NO_CCR_MARKER."),
308-
)
309304
@click.option(
310305
"--lossless",
311306
is_flag=True,
@@ -862,8 +857,7 @@ def proxy(
862857
protect_tool_results: str | None,
863858
rpm: int | None,
864859
tpm: int | None,
865-
no_ccr_inject_tool: bool,
866-
no_ccr_marker: bool,
860+
no_ccr: bool,
867861
lossless: bool,
868862
no_ccr_proactive_expansion: bool,
869863
proxy_extension: tuple[str, ...],
@@ -1102,11 +1096,12 @@ def proxy(
11021096
protect_recent=_get_env_int_optional("HEADROOM_PROTECT_RECENT"),
11031097
protect_analysis_context=_get_env_bool_optional("HEADROOM_PROTECT_ANALYSIS_CONTEXT"),
11041098
accuracy_guard=os.environ.get("HEADROOM_ACCURACY_GUARD") or None,
1105-
# CCR opt-outs for compression-only deployments (streaming / non-MCP
1106-
# clients that can't resolve the injected retrieve tool). Defaults keep
1107-
# CCR fully on; each flag flips one dataclass default to False.
1108-
ccr_inject_tool=not no_ccr_inject_tool,
1109-
ccr_inject_marker=not no_ccr_marker,
1099+
# CCR opt-out: --no-ccr disables both halves at once (markers in content
1100+
# AND the injected retrieve tool). Markers without a tool — or a tool
1101+
# without markers — are useless, so it is a single switch. Default keeps
1102+
# CCR fully on.
1103+
ccr_inject_tool=not no_ccr,
1104+
ccr_inject_marker=not no_ccr,
11101105
lossless=lossless,
11111106
ccr_proactive_expansion=not no_ccr_proactive_expansion,
11121107
# Flatten repeat-flag tuple AND any comma-separated values inside it.

‎headroom/proxy/models.py‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -133,8 +133,8 @@ class ProxyConfig:
133133
ccr_inject_tool: bool = True
134134
ccr_inject_system_instructions: bool = False
135135
# Proxy-level mirror of ContentRouterConfig.ccr_inject_marker, so retrieval
136-
# markers can be toggled from the CLI (--no-ccr-marker). Threaded into the
137-
# router in server.py; default preserves current behavior.
136+
# markers can be toggled from the CLI (--no-ccr, which also drops the retrieve
137+
# tool). Threaded into the router in server.py; default preserves current behavior.
138138
ccr_inject_marker: bool = True
139139

140140
# CCR Response Handling

‎headroom/transforms/content_router.py‎

Lines changed: 347 additions & 40 deletions
Large diffs are not rendered by default.
Lines changed: 231 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,231 @@
1+
"""Cross-turn (whole-conversation) verbatim de-duplication.
2+
3+
Bash coding agents re-display the same file bytes many times across turns
4+
(``cat foo.py`` -> ``sed -n 75,100p foo.py`` -> ``git diff`` -> ``cat foo.py``
5+
again). Every per-block compressor is blind to this: the redundancy is *across*
6+
blocks. This transform replaces a contiguous span in a later tool output that
7+
already appeared verbatim in an earlier tool output with a compact in-context
8+
pointer to the original.
9+
10+
Two hard invariants, both required for production use:
11+
12+
1. CACHE-SAFETY via *prefix-monotonicity*. Blocks are processed in order and a
13+
block is only ever matched against content from *strictly earlier* blocks.
14+
Therefore the rewritten output of blocks ``0..k`` is byte-identical whether
15+
or not block ``k+1`` exists — appending a turn never mutates an earlier turn,
16+
so the upstream prompt-cache prefix stays byte-stable. References are
17+
ABSOLUTE (an earlier block's ordinal), never relative, so a frozen pointer's
18+
text never changes. :func:`is_prefix_monotonic` asserts this.
19+
20+
2. ACCURACY via *no information leaves the window*. Only spans that are present
21+
VERBATIM in an earlier block's already-emitted output are back-referenced
22+
(the "verbatim corpus"), and the earliest occurrence is never rewritten
23+
(keep-earliest), so the original the pointer names is always physically in
24+
context. Only large, non-trivial contiguous spans are folded.
25+
26+
Pure stdlib, deterministic, never raises (returns input unchanged on any error).
27+
"""
28+
29+
from __future__ import annotations
30+
31+
from dataclasses import dataclass
32+
33+
__all__ = ["DedupBlock", "dedup_blocks", "is_prefix_monotonic"]
34+
35+
# A run must be at least this many lines AND this many chars to be worth a
36+
# pointer. Small dups are left alone (fragmenting context is not worth it) —
37+
# and a larger floor keeps the pointer comfortably shorter than the span it
38+
# replaces, so a fold is always a net byte win.
39+
DEFAULT_MIN_LINES = 7
40+
DEFAULT_MIN_CHARS = 120
41+
# Cap anchor candidates examined per line so a hot line (e.g. `` return``)
42+
# can't blow up matching. Deterministic: candidates are kept in first-seen order.
43+
MAX_ANCHOR_CANDIDATES = 16
44+
45+
46+
@dataclass
47+
class DedupBlock:
48+
"""One tool-output block. ``turn`` is a STABLE absolute ordinal used in the
49+
pointer text (must not change as the conversation grows). ``protected`` marks
50+
blocks that must not be rewritten (e.g. carry a cache_control breakpoint) —
51+
they are still indexed as reference targets."""
52+
53+
text: str
54+
turn: int
55+
protected: bool = False
56+
57+
58+
def _is_trivial(line: str) -> bool:
59+
"""A line too common/short to safely anchor a match on its own."""
60+
s = line.strip()
61+
if len(s) < 4:
62+
return True
63+
return s in {
64+
"return",
65+
"pass",
66+
"else:",
67+
"try:",
68+
"except:",
69+
"finally:",
70+
"break",
71+
"continue",
72+
"});",
73+
"})",
74+
"],",
75+
"),",
76+
'"""',
77+
"'''",
78+
"...",
79+
}
80+
81+
82+
def _pointer(span: list[str], ref_turn: int, ref_line: int) -> str:
83+
"""A one-line, obviously-a-reference marker naming the in-context original.
84+
85+
Includes a first-line anchor so the model can locate the block it already
86+
saw. Marker-free of any ``hash=`` retrieval token: recovery is in-context
87+
(the original is physically present earlier in the same request)."""
88+
anchor = next((ln.strip() for ln in span if ln.strip()), "")
89+
if len(anchor) > 80:
90+
anchor = anchor[:77] + "..."
91+
end_line = ref_line + len(span) - 1
92+
return (
93+
f"[headroom: {len(span)} lines identical to output shown earlier "
94+
f"(turn {ref_turn}, lines {ref_line}-{end_line}) — starts: {anchor!r}]"
95+
)
96+
97+
98+
def _index_lines(
99+
lines: list[str | None],
100+
block_pos: int,
101+
anchor_index: dict[str, list[tuple[int, int]]],
102+
) -> None:
103+
"""Record each non-trivial line's (block_pos, line_idx) as a future anchor.
104+
105+
Keeps first-seen order and caps the candidate list per line. Only VERBATIM
106+
(surviving) lines should be passed here — never the lines of a span that was
107+
replaced by a pointer."""
108+
for li, ln in enumerate(lines):
109+
if ln is None or _is_trivial(ln):
110+
continue
111+
bucket = anchor_index.setdefault(ln, [])
112+
if len(bucket) < MAX_ANCHOR_CANDIDATES:
113+
bucket.append((block_pos, li))
114+
115+
116+
def _longest_match(
117+
cur: list[str],
118+
start: int,
119+
anchor_index: dict[str, list[tuple[int, int]]],
120+
corpus: list[list[str | None]],
121+
) -> tuple[int, int, int] | None:
122+
"""Longest contiguous run in ``cur`` starting at ``start`` that appears
123+
verbatim inside a single earlier block. Returns (length, block_pos,
124+
ref_line_idx) or None. ``corpus[block_pos]`` holds that block's VERBATIM
125+
lines (``None`` where a span was already folded, which breaks contiguity)."""
126+
anchor = cur[start]
127+
candidates = anchor_index.get(anchor)
128+
if not candidates:
129+
return None
130+
best_len = 0
131+
best_bp = best_li = -1
132+
for bp, li in candidates:
133+
block_lines = corpus[bp]
134+
k = 0
135+
while (
136+
start + k < len(cur)
137+
and li + k < len(block_lines)
138+
and block_lines[li + k] is not None
139+
and cur[start + k] == block_lines[li + k]
140+
):
141+
k += 1
142+
# Deterministic tie-break: longer wins; on ties keep the earliest
143+
# (smallest block_pos, then line) already held in best_*.
144+
if k > best_len:
145+
best_len, best_bp, best_li = k, bp, li
146+
if best_len == 0:
147+
return None
148+
return best_len, best_bp, best_li
149+
150+
151+
def dedup_blocks(
152+
blocks: list[DedupBlock],
153+
*,
154+
min_lines: int = DEFAULT_MIN_LINES,
155+
min_chars: int = DEFAULT_MIN_CHARS,
156+
) -> tuple[list[DedupBlock], dict]:
157+
"""Rewrite later verbatim spans to in-context pointers. Prefix-monotonic
158+
(cache-safe) and information-preserving (accuracy-safe). Returns
159+
(new_blocks, stats). Never raises."""
160+
stats = {"spans_folded": 0, "lines_removed": 0, "chars_removed": 0, "blocks": len(blocks)}
161+
try:
162+
# corpus[i] = verbatim lines of block i's OUTPUT (None where folded).
163+
corpus: list[list[str | None]] = []
164+
anchor_index: dict[str, list[tuple[int, int]]] = {}
165+
out_blocks: list[DedupBlock] = []
166+
167+
for blk in blocks:
168+
lines = blk.text.split("\n")
169+
170+
if blk.protected:
171+
# Never rewrite; still a valid verbatim reference target.
172+
verbatim: list[str | None] = list(lines)
173+
_index_lines(verbatim, len(corpus), anchor_index)
174+
corpus.append(verbatim)
175+
out_blocks.append(blk)
176+
continue
177+
178+
out: list[str] = []
179+
verbatim = []
180+
i = 0
181+
n = len(lines)
182+
while i < n:
183+
m = _longest_match(lines, i, anchor_index, corpus)
184+
if m is not None and m[0] >= min_lines:
185+
span = lines[i : i + m[0]]
186+
span_text = "\n".join(span)
187+
if len(span_text) >= min_chars:
188+
ref_turn = blocks[m[1]].turn
189+
ptr = _pointer(span, ref_turn, m[2])
190+
out.append(ptr)
191+
# Folded span is NOT verbatim in this block's output:
192+
# mark None so it can't seed a later contiguous match,
193+
# and don't index it (keep-earliest).
194+
verbatim.extend([None] * m[0])
195+
stats["spans_folded"] += 1
196+
stats["lines_removed"] += m[0]
197+
stats["chars_removed"] += len(span_text) - len(ptr)
198+
i += m[0]
199+
continue
200+
out.append(lines[i])
201+
verbatim.append(lines[i])
202+
i += 1
203+
204+
# Index only the surviving verbatim lines of THIS block (first-seen).
205+
# None entries (folded spans) are kept in place so positions stay
206+
# aligned with ``corpus``; _index_lines skips them.
207+
_index_lines(verbatim, len(corpus), anchor_index)
208+
corpus.append(verbatim)
209+
out_blocks.append(DedupBlock(text="\n".join(out), turn=blk.turn, protected=False))
210+
211+
return out_blocks, stats
212+
except Exception: # never break the proxy
213+
return blocks, {"spans_folded": 0, "lines_removed": 0, "chars_removed": 0, "error": True}
214+
215+
216+
def is_prefix_monotonic(
217+
blocks: list[DedupBlock],
218+
*,
219+
min_lines: int = DEFAULT_MIN_LINES,
220+
min_chars: int = DEFAULT_MIN_CHARS,
221+
) -> bool:
222+
"""CACHE-SAFETY invariant: for every k, dedup(blocks[:k]) equals dedup(full)
223+
truncated to its first k blocks. i.e. appending a later turn never changes an
224+
earlier turn's rewritten bytes, so the prompt-cache prefix stays stable."""
225+
full, _ = dedup_blocks(blocks, min_lines=min_lines, min_chars=min_chars)
226+
full_text = [b.text for b in full]
227+
for k in range(1, len(blocks) + 1):
228+
partial, _ = dedup_blocks(blocks[:k], min_lines=min_lines, min_chars=min_chars)
229+
if [b.text for b in partial] != full_text[:k]:
230+
return False
231+
return True

‎tests/test_cli_proxy_env.py‎

Lines changed: 14 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -843,25 +843,25 @@ def mock_run_server(config, **kwargs):
843843
assert cfg.ccr_inject_marker is True
844844
assert cfg.ccr_proactive_expansion is True
845845

846-
def test_no_ccr_inject_tool_flag(self, runner):
847-
"""--no-ccr-inject-tool disables retrieve-tool injection only."""
846+
def test_no_ccr_flag(self, runner):
847+
"""--no-ccr disables BOTH the retrieve-tool injection and the markers."""
848848
captured_config = {}
849849

850850
def mock_run_server(config, **kwargs):
851851
captured_config["config"] = config
852852

853853
with patch("headroom.proxy.server.run_server", mock_run_server):
854-
result = runner.invoke(main, ["proxy", "--no-ccr-inject-tool"], catch_exceptions=False)
854+
result = runner.invoke(main, ["proxy", "--no-ccr"], catch_exceptions=False)
855855

856856
assert result.exit_code == 0, result.output
857857
cfg = captured_config["config"]
858858
assert cfg.ccr_inject_tool is False
859-
# Untouched flags remain on.
860-
assert cfg.ccr_inject_marker is True
859+
assert cfg.ccr_inject_marker is False
860+
# Unrelated CCR knob stays on.
861861
assert cfg.ccr_proactive_expansion is True
862862

863863
def test_compression_only_all_flags(self, runner):
864-
"""All three flags together yield a compression-only config."""
864+
"""--no-ccr + --no-ccr-proactive-expansion yields a compression-only config."""
865865
captured_config = {}
866866

867867
def mock_run_server(config, **kwargs):
@@ -872,8 +872,7 @@ def mock_run_server(config, **kwargs):
872872
main,
873873
[
874874
"proxy",
875-
"--no-ccr-inject-tool",
876-
"--no-ccr-marker",
875+
"--no-ccr",
877876
"--no-ccr-proactive-expansion",
878877
],
879878
catch_exceptions=False,
@@ -885,8 +884,8 @@ def mock_run_server(config, **kwargs):
885884
assert cfg.ccr_inject_marker is False
886885
assert cfg.ccr_proactive_expansion is False
887886

888-
def test_no_ccr_marker_from_env(self, runner):
889-
"""HEADROOM_NO_CCR_MARKER env var disables marker injection."""
887+
def test_no_ccr_from_env(self, runner):
888+
"""HEADROOM_NO_CCR env var disables both markers and tool injection."""
890889
captured_config = {}
891890

892891
def mock_run_server(config, **kwargs):
@@ -896,16 +895,18 @@ def mock_run_server(config, **kwargs):
896895
result = runner.invoke(
897896
main,
898897
["proxy"],
899-
env={"HEADROOM_NO_CCR_MARKER": "1"},
898+
env={"HEADROOM_NO_CCR": "1"},
900899
catch_exceptions=False,
901900
)
902901

903902
assert result.exit_code == 0, result.output
904-
assert captured_config["config"].ccr_inject_marker is False
903+
cfg = captured_config["config"]
904+
assert cfg.ccr_inject_marker is False
905+
assert cfg.ccr_inject_tool is False
905906

906907

907908
class TestNoCcrMarkerCompressors:
908-
"""Verify --no-ccr-marker actually suppresses <<ccr:...>> markers
909+
"""Verify --no-ccr actually suppresses <<ccr:...>> markers
909910
from every compressor, not just SmartCrusher (#1022)."""
910911

911912
def test_content_router_propagates_ccr_inject_marker_false_to_compressors(self):

0 commit comments

Comments
 (0)