mirror of
https://github.com/headroomlabs-ai/headroom.git
synced 2026-08-27 14:17:10 -04:00
## Description
A data-driven push for better compression savings without accuracy loss,
in four parts: expose and tune the Rust compressor knobs, harden the CCR
retrieval store, add traffic-audit tooling that sizes opportunities from
real transcripts, and introduce **read maturation** — a new,
live-validated mechanism that compresses Read outputs *before* they ever
enter the provider prefix cache.
## Type of Change
- [x] Bug fix (non-breaking change that fixes an issue)
- [x] New feature (non-breaking change that adds functionality)
- [ ] Breaking change (fix or feature that would cause existing
functionality to change)
- [ ] Documentation update
- [x] Performance improvement
- [ ] Code refactoring (no functional changes)
## Changes Made
### 1. Rust compressor extraction
- Expose `lossless_min_savings_ratio` end-to-end and lower the default
0.30 → 0.15 (lockstep across Rust, PyO3, and both Python config classes)
so the lossless Table/CSV compaction path wins more often.
- Expose the `CompactConfig` heuristics (core-field fraction,
heterogeneity ratio, flatten cap, bucket bounds) through PyO3 + Python.
- `SearchCompressor` grouped-by-file output (`rg --heading` style — path
once per file instead of per match). Library default off; the proxy
enables it in token mode.
- Complete `factor_out_constants`: constant fields now emit once in a
`_constant_fields` sentinel with slim rows (defensive per-item value
match; default off).
- `ContentRouter` accepts a SmartCrusher config override and the
search-grouping knob.
### 2. CCR store hardening
- Session-scale TTL: 300s → 1800s (CCRConfig, CompressionEntry,
CompressionStore, Rust `DEFAULT_TTL` — lockstep).
- **SQLite is the default CCR backend** (`~/.headroom/ccr_store.db`,
WAL): survives proxy restarts and is shared across workers.
`HEADROOM_CCR_BACKEND=memory` opts out.
- Multi-worker safety: `busy_timeout`, and corruption detection narrowed
so transient `SQLITE_BUSY` errors can never trigger database deletion.
- Data-at-rest hygiene: `chmod 600` on db + sidecars, expired rows swept
at open.
- Retrieval-miss messages are actionable (re-read the file / re-run the
command).
### 3. Traffic audit tooling (measure before tuning)
- `headroom audit-reads`: sizes Read opportunities from local Claude
Code transcripts (read share, stale %, line-number overhead, context
residency, cache-death windows).
- `--simulate-maturation`: Mechanism B risk sizing (re-read rates,
never-touched-again share, quiesce coverage, at-risk edits).
- `--codex`: shell-read classifier for Codex transcripts (rtk-wrapper
aware, workdir resolution).
- Findings that shaped this PR (81 sessions): Reads are 67% of tool
bytes; median Read lingers 118 turns (~13x lifetime cost); a prototyped
repeat-Read dedup measured 0.1% and was **removed** rather than shipped
as dead code.
### 4. Read maturation (Mechanism B) — experimental, default OFF
- Activity-based: a fresh large Read is held **out** of the provider
cache (trailing breakpoint relocated before it), stays verbatim while
its file is active, and matures into a CCR-backed marker once the file
is quiet for `quiesce_turns` (default 5; `max_hold_turns` bounds busy
files).
- Only the final compressed form ever enters the cache — **no cached
byte is ever mutated**; matured markers replay byte-identically.
- Wired into the Anthropic handler behind `--read-maturation` /
`HEADROOM_READ_MATURATION=1`; session state rides on the prefix tracker;
advisory (can never fail a request).
- Live-validated against the Anthropic API: held content excluded from
cache_creation; after maturation the prior cached prefix still served —
the no-bust invariant holds end-to-end.
### 5. Rebase / CI fixups (this update)
- Rebased onto latest `main` (was 28 commits behind): picks up `ci: pass
CODECOV_TOKEN to coverage uploads (#968)`, which is what was turning the
4 test shards red — the tests themselves passed (1528) but the post-test
codecov upload exited non-zero on a protected branch.
- Resolved the duplicate `lossless_min_savings_ratio` that two
independent main/branch additions left in `SmartCrusherConfig` and the
Rust-config kwarg (import-time `SyntaxError` + mypy `no-redef`).
- Aligned CCR tests with the new defaults (SQLite backend, 1800s TTL)
across `test_ccr`, `test_adapter_hooks`, `test_compression_store`,
`test_proxy_ccr`, and the lossy row-drop bridge test.
## Testing
<!-- Check what you actually ran, then paste the real command output
below. -->
- [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
### Test Output
```text
$ python -m pytest tests/test_proxy_ccr.py tests/test_ccr.py tests/test_compression_store.py tests/test_adapter_hooks.py tests/test_ccr_row_drop_store_bridge.py -q
170 passed, 4 warnings in 42.49s
$ python -m pytest tests/test_audit_reads.py tests/test_audit_codex.py tests/test_read_maturation.py tests/test_transforms_content_router.py tests/test_smart_crusher_toin_attachment.py -q
83 passed
$ mypy headroom/
Success: no issues found in 365 source files
$ python -m compileall headroom/ -q
COMPILE-OK
# CI (run 27488990477, pre-rebase head): all 4 shards ran to completion —
# "1528 passed, 120 skipped, 4922 deselected"
# The red shards were the codecov upload step, not test failures; fixed by
# the #968 rebase above.
```
## Real Behavior Proof
- Environment: macOS (darwin), Python 3.12 venv; branch
`feat/compression-extraction` rebased onto `origin/main` (head
7cb0f43b); GitHub Actions CI run 27488990477 for the test shards
- Exact command / steps: rebased onto latest main (clean, 13 commits
replayed, 0 conflicts); ran the pytest suites and mypy above locally;
inspected CI shard logs to confirm the failure was the codecov upload,
not the test phase
- Observed result: 253 targeted tests pass locally; mypy clean on 365
files; CI test phase reports `1528 passed, 120 skipped`; the only red
step (codecov `upload-coverage` → "Token required because branch is
protected") is resolved by the rebased-in #968 CODECOV_TOKEN fix
- Not tested: the read-maturation live-API no-bust validation
(`tests/test_live/`) was not re-run in this rebase pass (requires
provider keys); it was validated when the feature first landed, and no
maturation code changed in the rebase — only CCR-default test assertions
and the duplicate-field resolution
## 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
- [ ] I have updated the CHANGELOG.md if applicable
## Additional Notes
CHANGELOG is generated by release-please from the conventional commits,
so the CHANGELOG box is intentionally left unchecked. "Manual testing
performed" is unchecked deliberately — see `Real Behavior Proof` → `Not
tested` for the exact boundary (the live-API maturation validation was
not re-run in this rebase pass).
### Follow-ups (tracked, not in this PR)
- Mechanism B provider extensions: OpenAI-family wiring (no breakpoint
hold — bounded near-tail bust) and the Codex runtime read-detector (the
audit classifier is the prototype).
- Pilot enablement playbook: run `audit-reads --simulate-maturation` on
target traffic → pick `quiesce_turns` → enable via env → watch cache hit
rate + `read_maturation:N` transform tags.
742 lines
28 KiB
Python
742 lines
28 KiB
Python
"""Tests for CCR endpoints in the proxy server.
|
|
|
|
These tests verify the /v1/retrieve endpoints work correctly.
|
|
"""
|
|
|
|
import json
|
|
from unittest.mock import patch
|
|
|
|
import pytest
|
|
|
|
# Skip if fastapi not available
|
|
pytest.importorskip("fastapi")
|
|
|
|
from fastapi.testclient import TestClient
|
|
|
|
from headroom.cache.compression_store import get_compression_store, reset_compression_store
|
|
from headroom.proxy.server import ProxyConfig, create_app
|
|
|
|
|
|
@pytest.fixture
|
|
def client():
|
|
"""Create test client with fresh compression store."""
|
|
reset_compression_store()
|
|
config = ProxyConfig(
|
|
optimize=False, # Disable optimization for simpler tests
|
|
cache_enabled=False,
|
|
rate_limit_enabled=False,
|
|
cost_tracking_enabled=False,
|
|
)
|
|
app = create_app(config)
|
|
with TestClient(app) as client:
|
|
yield client
|
|
reset_compression_store()
|
|
|
|
|
|
@pytest.fixture
|
|
def client_with_data(client):
|
|
"""Test client with pre-populated compression store."""
|
|
store = get_compression_store()
|
|
|
|
# Store some test data
|
|
items = [{"id": i, "content": f"Item {i} about Python programming"} for i in range(100)]
|
|
store.store(
|
|
original=json.dumps(items),
|
|
compressed=json.dumps(items[:10]),
|
|
original_tokens=1000,
|
|
compressed_tokens=100,
|
|
original_item_count=100,
|
|
compressed_item_count=10,
|
|
tool_name="test_tool",
|
|
)
|
|
|
|
return client
|
|
|
|
|
|
class TestCCRRetrieveEndpoint:
|
|
"""Test the /v1/retrieve POST endpoint."""
|
|
|
|
def test_retrieve_requires_hash(self, client):
|
|
"""Request without hash should return 400."""
|
|
response = client.post("/v1/retrieve", json={})
|
|
assert response.status_code == 400
|
|
assert "hash required" in response.json()["detail"]
|
|
|
|
def test_retrieve_nonexistent_hash(self, client):
|
|
"""Request with nonexistent hash should return 404."""
|
|
response = client.post("/v1/retrieve", json={"hash": "nonexistent123"})
|
|
assert response.status_code == 404
|
|
assert "Entry not found" in response.json()["detail"]
|
|
assert "CCR TTL: 1800 seconds" in response.json()["detail"]
|
|
|
|
def test_retrieve_expired_hash_reports_expiration_detail(self, client):
|
|
"""Expired entries report expiration separately from missing hashes."""
|
|
store = get_compression_store(default_ttl=1)
|
|
with patch("headroom.cache.compression_store.time.time", return_value=1000.0):
|
|
hash_key = store.store(original="payload", compressed="payload")
|
|
|
|
with patch("headroom.cache.compression_store.time.time", return_value=1002.0):
|
|
response = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
|
|
assert response.status_code == 404
|
|
detail = response.json()["detail"]
|
|
assert "Entry expired" in detail
|
|
assert "CCR TTL: 1 seconds" in detail
|
|
assert "age: 2 seconds" in detail
|
|
|
|
def test_retrieve_full_content(self, client):
|
|
"""Full retrieval returns original content."""
|
|
store = get_compression_store()
|
|
items = [{"id": i} for i in range(50)]
|
|
hash_key = store.store(
|
|
original=json.dumps(items),
|
|
compressed="[]",
|
|
original_item_count=50,
|
|
compressed_item_count=0,
|
|
)
|
|
|
|
response = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["hash"] == hash_key
|
|
assert data["original_item_count"] == 50
|
|
assert "original_content" in data
|
|
|
|
# Verify content is correct
|
|
retrieved_items = json.loads(data["original_content"])
|
|
assert len(retrieved_items) == 50
|
|
assert retrieved_items[0]["id"] == 0
|
|
|
|
def test_retrieve_with_search(self, client):
|
|
"""Search retrieval filters by query."""
|
|
store = get_compression_store()
|
|
items = [
|
|
{"id": 1, "text": "Python programming language"},
|
|
{"id": 2, "text": "JavaScript web development"},
|
|
{"id": 3, "text": "Python data science"},
|
|
{"id": 4, "text": "Java enterprise"},
|
|
]
|
|
hash_key = store.store(
|
|
original=json.dumps(items),
|
|
compressed="[]",
|
|
original_item_count=4,
|
|
compressed_item_count=0,
|
|
)
|
|
|
|
response = client.post(
|
|
"/v1/retrieve", json={"hash": hash_key, "query": "Python programming"}
|
|
)
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["hash"] == hash_key
|
|
assert data["query"] == "Python programming"
|
|
assert "results" in data
|
|
assert data["count"] >= 1
|
|
|
|
def test_retrieve_with_search_plain_text_original(self, client):
|
|
"""Query retrieval searches plain-text originals stored by Kompress."""
|
|
store = get_compression_store()
|
|
original = (
|
|
"Codex WS compression stores plain text originals. "
|
|
"The target symbol is _compress_openai_responses_payload."
|
|
)
|
|
hash_key = store.store(original=original, compressed="compressed")
|
|
|
|
response = client.post(
|
|
"/v1/retrieve",
|
|
json={"hash": hash_key, "query": "_compress_openai_responses_payload"},
|
|
)
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["hash"] == hash_key
|
|
assert data["count"] == 1
|
|
assert data["results"][0]["type"] == "text"
|
|
assert "_compress_openai_responses_payload" in data["results"][0]["text"]
|
|
|
|
def test_retrieve_with_search_nonexistent_hash_returns_404(self, client):
|
|
"""Query mode should not mask a missing hash as an empty search."""
|
|
response = client.post(
|
|
"/v1/retrieve",
|
|
json={"hash": "nonexistent123", "query": "anything"},
|
|
)
|
|
assert response.status_code == 404
|
|
|
|
def test_retrieve_increments_count(self, client):
|
|
"""Each retrieval increments the retrieval count."""
|
|
store = get_compression_store()
|
|
hash_key = store.store(original="[]", compressed="[]")
|
|
|
|
# First retrieval
|
|
response1 = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response1.status_code == 200
|
|
count1 = response1.json()["retrieval_count"]
|
|
|
|
# Second retrieval
|
|
response2 = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response2.status_code == 200
|
|
count2 = response2.json()["retrieval_count"]
|
|
|
|
assert count2 > count1
|
|
|
|
|
|
class TestCCRRetrieveGetEndpoint:
|
|
"""Test the /v1/retrieve/{hash_key} GET endpoint."""
|
|
|
|
def test_get_retrieve_full(self, client):
|
|
"""GET retrieval returns full content."""
|
|
store = get_compression_store()
|
|
items = [{"id": i} for i in range(20)]
|
|
hash_key = store.store(
|
|
original=json.dumps(items),
|
|
compressed="[]",
|
|
original_item_count=20,
|
|
compressed_item_count=0,
|
|
tool_name="get_test_tool",
|
|
)
|
|
|
|
response = client.get(f"/v1/retrieve/{hash_key}")
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["hash"] == hash_key
|
|
assert data["original_item_count"] == 20
|
|
assert data["tool_name"] == "get_test_tool"
|
|
|
|
def test_get_retrieve_with_query(self, client):
|
|
"""GET retrieval with query parameter invokes search."""
|
|
store = get_compression_store()
|
|
# Create items with distinctive content
|
|
items = [
|
|
{"id": 1, "msg": "Python programming language tutorial for beginners"},
|
|
{"id": 2, "msg": "JavaScript web development framework guide"},
|
|
{"id": 3, "msg": "Python data science machine learning pandas"},
|
|
{"id": 4, "msg": "Java enterprise application development"},
|
|
]
|
|
hash_key = store.store(
|
|
original=json.dumps(items),
|
|
compressed="[]",
|
|
)
|
|
|
|
response = client.get(f"/v1/retrieve/{hash_key}?query=Python programming")
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["query"] == "Python programming"
|
|
# Response includes search results structure
|
|
assert "results" in data
|
|
assert "count" in data
|
|
# Results should be a list (may be empty if BM25 threshold not met)
|
|
assert isinstance(data["results"], list)
|
|
|
|
def test_get_retrieve_with_query_plain_text_original(self, client):
|
|
"""GET query retrieval searches plain-text originals."""
|
|
store = get_compression_store()
|
|
hash_key = store.store(
|
|
original="plain text contains _compress_openai_responses_payload",
|
|
compressed="plain text",
|
|
)
|
|
|
|
response = client.get(f"/v1/retrieve/{hash_key}?query=_compress_openai_responses_payload")
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["count"] == 1
|
|
assert data["results"][0]["type"] == "text"
|
|
|
|
def test_get_retrieve_nonexistent(self, client):
|
|
"""GET with nonexistent hash returns 404."""
|
|
response = client.get("/v1/retrieve/nonexistent123")
|
|
assert response.status_code == 404
|
|
|
|
|
|
class TestCCRStatsEndpoint:
|
|
"""Test the /v1/retrieve/stats endpoint."""
|
|
|
|
def test_stats_empty_store(self, client):
|
|
"""Stats with empty store returns zeros."""
|
|
response = client.get("/v1/retrieve/stats")
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert "store" in data
|
|
assert data["store"]["entry_count"] == 0
|
|
assert data["store"]["default_ttl_seconds"] == 1800
|
|
assert "recent_retrievals" in data
|
|
|
|
def test_stats_exposes_env_configured_ttl(self, client, monkeypatch):
|
|
"""Stats expose the effective CCR TTL configured through env."""
|
|
reset_compression_store()
|
|
monkeypatch.setenv("HEADROOM_CCR_TTL_SECONDS", "7200")
|
|
|
|
response = client.get("/v1/retrieve/stats")
|
|
|
|
assert response.status_code == 200
|
|
assert response.json()["store"]["default_ttl_seconds"] == 7200
|
|
|
|
def test_stats_with_entries(self, client):
|
|
"""Stats reflect store contents."""
|
|
store = get_compression_store()
|
|
|
|
# Add some entries
|
|
store.store(original="[1]", compressed="[]", original_tokens=100)
|
|
store.store(original="[2]", compressed="[]", original_tokens=200)
|
|
|
|
response = client.get("/v1/retrieve/stats")
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["store"]["entry_count"] == 2
|
|
assert data["store"]["total_original_tokens"] == 300
|
|
|
|
def test_stats_tracks_retrievals(self, client):
|
|
"""Stats include recent retrieval events."""
|
|
import json as json_module
|
|
|
|
store = get_compression_store()
|
|
|
|
# Use non-empty content so search actually logs
|
|
content = json_module.dumps(
|
|
[
|
|
{"id": "1", "name": "test item", "value": 100},
|
|
{"id": "2", "name": "another item", "value": 200},
|
|
]
|
|
)
|
|
hash_key = store.store(
|
|
original=content,
|
|
compressed=content,
|
|
tool_name="stats_test_tool",
|
|
)
|
|
|
|
# Make some retrievals
|
|
client.post("/v1/retrieve", json={"hash": hash_key}) # Full retrieval
|
|
client.post("/v1/retrieve", json={"hash": hash_key, "query": "test"}) # Search retrieval
|
|
|
|
response = client.get("/v1/retrieve/stats")
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["store"]["total_retrievals"] >= 2
|
|
assert len(data["recent_retrievals"]) >= 2
|
|
|
|
# Verify we have both retrieval types (no double-logging of full)
|
|
retrieval_types = [r["retrieval_type"] for r in data["recent_retrievals"]]
|
|
assert "full" in retrieval_types
|
|
assert "search" in retrieval_types
|
|
|
|
|
|
class TestCCRIntegration:
|
|
"""Integration tests for CCR with proxy."""
|
|
|
|
def test_health_endpoint(self, client):
|
|
"""Health endpoint works."""
|
|
response = client.get("/health")
|
|
assert response.status_code == 200
|
|
assert response.json()["status"] == "healthy"
|
|
|
|
def test_stats_endpoint(self, client):
|
|
"""Stats endpoint includes CCR-relevant info."""
|
|
response = client.get("/stats")
|
|
assert response.status_code == 200
|
|
# Proxy stats endpoint is separate from CCR stats
|
|
data = response.json()
|
|
assert "requests" in data
|
|
assert "tokens" in data
|
|
|
|
|
|
class TestCCREdgeCases:
|
|
"""Edge cases for CCR endpoints."""
|
|
|
|
def test_retrieve_empty_content(self, client):
|
|
"""Retrieve works with empty content."""
|
|
store = get_compression_store()
|
|
hash_key = store.store(original="[]", compressed="[]")
|
|
|
|
response = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response.status_code == 200
|
|
assert response.json()["original_content"] == "[]"
|
|
|
|
def test_retrieve_large_content(self, client):
|
|
"""Retrieve works with large content."""
|
|
store = get_compression_store()
|
|
items = [{"id": i, "data": "x" * 100} for i in range(1000)]
|
|
hash_key = store.store(
|
|
original=json.dumps(items),
|
|
compressed=json.dumps(items[:10]),
|
|
original_item_count=1000,
|
|
)
|
|
|
|
response = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["original_item_count"] == 1000
|
|
|
|
def test_search_no_matches(self, client):
|
|
"""Search with no matches returns empty results."""
|
|
store = get_compression_store()
|
|
items = [{"id": 1, "text": "hello world"}]
|
|
hash_key = store.store(original=json.dumps(items), compressed="[]")
|
|
|
|
response = client.post("/v1/retrieve", json={"hash": hash_key, "query": "xyznonexistent"})
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
assert data["count"] == 0
|
|
assert data["results"] == []
|
|
|
|
def test_unicode_content(self, client):
|
|
"""Unicode content is handled correctly."""
|
|
store = get_compression_store()
|
|
items = [
|
|
{"id": 1, "text": "日本語テキスト"},
|
|
{"id": 2, "text": "Émoji 🎉 test"},
|
|
]
|
|
hash_key = store.store(original=json.dumps(items, ensure_ascii=False), compressed="[]")
|
|
|
|
response = client.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response.status_code == 200
|
|
|
|
data = response.json()
|
|
retrieved = json.loads(data["original_content"])
|
|
assert retrieved[0]["text"] == "日本語テキスト"
|
|
assert "🎉" in retrieved[1]["text"]
|
|
|
|
|
|
class TestEndToEndTOINIntegration:
|
|
"""End-to-end tests verifying the production path from proxy → TOIN.
|
|
|
|
These tests verify that:
|
|
1. SmartCrusher compresses tool outputs when called through the proxy pipeline
|
|
2. TOIN records compression events
|
|
3. Retrieval events update TOIN field semantics
|
|
4. The full feedback loop works
|
|
|
|
This catches bugs where components are wired correctly but don't communicate
|
|
(e.g., compression_store not passing retrieved_items to TOIN).
|
|
"""
|
|
|
|
@pytest.fixture
|
|
def fresh_toin(self):
|
|
"""Create a fresh TOIN instance."""
|
|
import tempfile
|
|
from pathlib import Path
|
|
|
|
from headroom.telemetry.toin import (
|
|
TOINConfig,
|
|
get_toin,
|
|
reset_toin,
|
|
)
|
|
|
|
reset_toin()
|
|
with tempfile.TemporaryDirectory() as tmpdir:
|
|
storage_path = str(Path(tmpdir) / "toin.json")
|
|
toin = get_toin(
|
|
TOINConfig(
|
|
storage_path=storage_path,
|
|
auto_save_interval=0,
|
|
)
|
|
)
|
|
yield toin
|
|
reset_toin()
|
|
|
|
@pytest.fixture
|
|
def client_with_optimization(self, fresh_toin):
|
|
"""Create test client with optimization enabled."""
|
|
reset_compression_store()
|
|
config = ProxyConfig(
|
|
optimize=True, # Enable optimization
|
|
cache_enabled=False,
|
|
rate_limit_enabled=False,
|
|
cost_tracking_enabled=False,
|
|
)
|
|
app = create_app(config)
|
|
with TestClient(app) as client:
|
|
yield client
|
|
reset_compression_store()
|
|
|
|
def test_pipeline_compresses_tool_output_and_records_toin(
|
|
self, fresh_toin, client_with_optimization
|
|
):
|
|
"""CRITICAL: Verify SmartCrusher compression records events in TOIN.
|
|
|
|
This tests the production code path:
|
|
1. Tool output comes in through proxy
|
|
2. SmartCrusher compresses it
|
|
3. TOIN records the compression event
|
|
"""
|
|
from headroom.config import CCRConfig, SmartCrusherConfig
|
|
from headroom.providers import AnthropicProvider
|
|
from headroom.telemetry import ToolSignature
|
|
from headroom.transforms import SmartCrusher, TransformPipeline
|
|
|
|
# Create tool output with 100 items that will trigger compression
|
|
# Key: score field with varying values signals sortable data
|
|
# Having repetitive category values helps trigger compression
|
|
items = [
|
|
{
|
|
"id": i,
|
|
"score": 1000 - i, # Decreasing scores signal sorting
|
|
"category": f"cat_{i % 3}", # Only 3 unique categories
|
|
"status": "active" if i % 2 == 0 else "inactive", # Binary status
|
|
}
|
|
for i in range(100)
|
|
]
|
|
tool_output = json.dumps(items)
|
|
|
|
# Create messages with tool_result containing our data
|
|
messages = [
|
|
{"role": "user", "content": "Search for items"},
|
|
{
|
|
"role": "assistant",
|
|
"content": [
|
|
{
|
|
"type": "tool_use",
|
|
"id": "tool_123",
|
|
"name": "search_api",
|
|
"input": {"query": "test"},
|
|
}
|
|
],
|
|
},
|
|
{
|
|
"role": "user",
|
|
"content": [
|
|
{
|
|
"type": "tool_result",
|
|
"tool_use_id": "tool_123",
|
|
"content": tool_output,
|
|
}
|
|
],
|
|
},
|
|
]
|
|
|
|
# Create pipeline with SmartCrusher (same as proxy does).
|
|
# Use with_compaction=False so we exercise the lossy + CCR
|
|
# caching path that this test asserts. The PR4 lossless
|
|
# default substitutes a CSV+schema string and skips CCR
|
|
# caching (nothing dropped → no cache entry).
|
|
pipeline = TransformPipeline(
|
|
transforms=[
|
|
SmartCrusher(
|
|
SmartCrusherConfig(
|
|
enabled=True,
|
|
min_tokens_to_crush=100,
|
|
max_items_after_crush=15,
|
|
),
|
|
ccr_config=CCRConfig(
|
|
enabled=True,
|
|
inject_retrieval_marker=True,
|
|
min_items_to_cache=10,
|
|
),
|
|
with_compaction=False,
|
|
),
|
|
],
|
|
provider=AnthropicProvider(),
|
|
)
|
|
|
|
# Apply pipeline (this is what the proxy does)
|
|
result = pipeline.apply(
|
|
messages=messages,
|
|
model="claude-sonnet-4-20250514",
|
|
model_limit=200000,
|
|
)
|
|
|
|
# Verify SmartCrusher was invoked (transform name starts with smart_crush)
|
|
smart_crush_applied = any(
|
|
t.startswith("smart_crush") or t.startswith("smart:") for t in result.transforms_applied
|
|
)
|
|
assert smart_crush_applied, (
|
|
f"SmartCrusher should be in transforms: {result.transforms_applied}"
|
|
)
|
|
|
|
# Check if compression was actually performed (not skipped)
|
|
# Skip messages look like "smart:skip:reason(100->100)"
|
|
compression_was_skipped = any(
|
|
"skip" in t.lower() for t in result.transforms_applied if "smart:" in t.lower()
|
|
)
|
|
|
|
# If compression happened, verify TOIN and store
|
|
if not compression_was_skipped:
|
|
# Verify compression store has the entry
|
|
store = get_compression_store()
|
|
stats = store.get_stats()
|
|
assert stats["entry_count"] >= 1, "Should have cached entry"
|
|
|
|
# Verify TOIN recorded the compression
|
|
signature = ToolSignature.from_items(items)
|
|
pattern = fresh_toin._patterns.get(signature.structure_hash)
|
|
assert pattern is not None, (
|
|
"TOIN should have recorded compression event. "
|
|
"If this fails, SmartCrusher is not calling TOIN.record_compression."
|
|
)
|
|
assert pattern.total_compressions >= 1, "Should have at least 1 compression"
|
|
else:
|
|
# Compression was skipped - this is expected for some data patterns
|
|
# The important thing is that SmartCrusher was invoked and made a decision
|
|
# The other tests verify the full loop when compression does happen
|
|
pass
|
|
|
|
def test_retrieval_through_proxy_updates_toin_field_semantics(
|
|
self, fresh_toin, client_with_optimization
|
|
):
|
|
"""CRITICAL: Verify retrieval through proxy updates TOIN field semantics.
|
|
|
|
This tests the full feedback loop:
|
|
1. Store compressed content (simulating prior compression)
|
|
2. Retrieve through proxy endpoint
|
|
3. Verify TOIN learned field semantics from retrieved items
|
|
"""
|
|
from headroom.telemetry import ToolSignature
|
|
|
|
# Create items with distinctive field types
|
|
items = [
|
|
{
|
|
"id": i,
|
|
"error_code": 500 if i % 10 == 0 else 200,
|
|
"timestamp": f"2024-01-{i:02d}T00:00:00Z",
|
|
"message": f"Log entry {i}",
|
|
}
|
|
for i in range(50)
|
|
]
|
|
original_content = json.dumps(items)
|
|
compressed_content = json.dumps(items[:10])
|
|
|
|
# Get the signature hash
|
|
signature = ToolSignature.from_items(items)
|
|
|
|
# Store in compression store with correct metadata
|
|
store = get_compression_store()
|
|
hash_key = store.store(
|
|
original=original_content,
|
|
compressed=compressed_content,
|
|
original_item_count=50,
|
|
compressed_item_count=10,
|
|
tool_name="logs_api",
|
|
tool_signature_hash=signature.structure_hash,
|
|
compression_strategy="smart_sample",
|
|
)
|
|
|
|
# Pre-record some compressions in TOIN (needed for pattern to exist)
|
|
for _ in range(3):
|
|
fresh_toin.record_compression(
|
|
tool_signature=signature,
|
|
original_count=50,
|
|
compressed_count=10,
|
|
original_tokens=5000,
|
|
compressed_tokens=1000,
|
|
strategy="smart_sample",
|
|
)
|
|
|
|
# Retrieve through proxy endpoint
|
|
response = client_with_optimization.post("/v1/retrieve", json={"hash": hash_key})
|
|
assert response.status_code == 200
|
|
|
|
# Process pending feedback (this is what triggers TOIN learning)
|
|
# Note: get_compression_store is already imported at module level
|
|
store = get_compression_store()
|
|
store.process_pending_feedback()
|
|
|
|
# PR-B5: pattern key is now `(auth_mode, model_family, sig_hash)`.
|
|
# Callers that don't supply auth/model land on the
|
|
# `("unknown", "unknown", sig_hash)` slot.
|
|
from headroom.telemetry.toin import _make_pattern_key
|
|
|
|
pattern = fresh_toin._patterns.get(_make_pattern_key(None, None, signature.structure_hash))
|
|
assert pattern is not None, "Pattern should exist after compression and retrieval"
|
|
|
|
# CRITICAL ASSERTION: This catches the bug where compression_store
|
|
# wasn't passing retrieved_items to TOIN
|
|
assert len(pattern.field_semantics) > 0, (
|
|
"TOIN should have learned field semantics from retrieved items. "
|
|
"If this fails, the production code path "
|
|
"(CompressionStore.process_pending_feedback -> TOIN.record_retrieval) "
|
|
"is not passing retrieved_items."
|
|
)
|
|
|
|
# Verify specific field types were learned
|
|
field_names = list(pattern.field_semantics.keys())
|
|
assert len(field_names) > 0, "Should have learned at least one field"
|
|
|
|
def test_full_proxy_ccr_feedback_loop(self, fresh_toin, client_with_optimization):
|
|
"""CRITICAL: Test the complete CCR feedback loop through proxy.
|
|
|
|
This is the most important integration test - it verifies:
|
|
1. Compression happens and TOIN records it
|
|
2. Retrieval happens and TOIN learns from it
|
|
3. Future recommendations reflect the learning
|
|
"""
|
|
from headroom.telemetry import ToolSignature
|
|
|
|
# Create items for the full feedback loop test
|
|
items = [
|
|
{
|
|
"id": i,
|
|
"score": 1000 - i,
|
|
"category": f"cat_{i % 5}",
|
|
"status": "active" if i % 2 == 0 else "inactive",
|
|
}
|
|
for i in range(100)
|
|
]
|
|
signature = ToolSignature.from_items(items)
|
|
|
|
# Store content directly (simulating what SmartCrusher does)
|
|
# This ensures we have entries regardless of whether compression was triggered
|
|
store = get_compression_store()
|
|
hash_key = store.store(
|
|
original=json.dumps(items),
|
|
compressed=json.dumps(items[:15]),
|
|
original_item_count=100,
|
|
compressed_item_count=15,
|
|
tool_name="search_api",
|
|
tool_signature_hash=signature.structure_hash,
|
|
compression_strategy="smart_sample",
|
|
)
|
|
|
|
# Record compressions in TOIN (simulating what SmartCrusher does)
|
|
for _ in range(3):
|
|
fresh_toin.record_compression(
|
|
tool_signature=signature,
|
|
original_count=100,
|
|
compressed_count=15,
|
|
original_tokens=5000,
|
|
compressed_tokens=1000,
|
|
strategy="smart_sample",
|
|
)
|
|
|
|
# Step 2: Retrieve through proxy endpoint
|
|
response = client_with_optimization.post(
|
|
"/v1/retrieve",
|
|
json={"hash": hash_key, "query": "category:cat_1"},
|
|
)
|
|
assert response.status_code == 200
|
|
|
|
# Process feedback (this triggers TOIN learning)
|
|
store.process_pending_feedback()
|
|
|
|
# Step 3: Verify TOIN learned
|
|
# PR-B5: pattern key is now `(auth_mode, model_family, sig_hash)`.
|
|
# Callers that don't supply auth/model land on the
|
|
# `("unknown", "unknown", sig_hash)` slot.
|
|
from headroom.telemetry.toin import _make_pattern_key
|
|
|
|
pattern = fresh_toin._patterns.get(_make_pattern_key(None, None, signature.structure_hash))
|
|
assert pattern is not None, "Pattern should exist"
|
|
assert pattern.total_compressions >= 1, "Should have compression count"
|
|
assert pattern.total_retrievals >= 1, "Should have retrieval count"
|
|
|
|
# Step 4: Verify field semantics were learned
|
|
assert len(pattern.field_semantics) > 0, (
|
|
"TOIN should learn field semantics through the full proxy CCR loop. "
|
|
"This is the ultimate integration test - if this fails, "
|
|
"the production feedback loop is broken."
|
|
)
|
|
|
|
# Step 5: PR-B5 retired the request-time recommendation API in favor of
|
|
# observation-only learning + startup-published recommendations.toml.
|
|
# `get_recommendation()` now returns None and emits a deprecation
|
|
# warning; the dispatcher consumes published advice via the Rust
|
|
# `RecommendationStore`. Assert the deprecation contract here so a
|
|
# future revival of the API doesn't slip past silently.
|
|
assert fresh_toin.get_recommendation(signature, "find category") is None
|