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.
278 lines
8.8 KiB
Python
278 lines
8.8 KiB
Python
"""Tests for proxy scalability features.
|
|
|
|
These tests verify connection pooling, HTTP/2, and worker configuration.
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
import os
|
|
from unittest.mock import patch
|
|
|
|
import httpx
|
|
import pytest
|
|
|
|
|
|
class TestConnectionPoolConfig:
|
|
"""Test connection pool configuration."""
|
|
|
|
def test_httpx_limits_basic(self):
|
|
"""Test that httpx accepts our connection limits."""
|
|
limits = httpx.Limits(
|
|
max_connections=500,
|
|
max_keepalive_connections=100,
|
|
)
|
|
assert limits.max_connections == 500
|
|
assert limits.max_keepalive_connections == 100
|
|
|
|
def test_httpx_limits_custom(self):
|
|
"""Test custom connection limits."""
|
|
limits = httpx.Limits(
|
|
max_connections=1000,
|
|
max_keepalive_connections=200,
|
|
)
|
|
assert limits.max_connections == 1000
|
|
assert limits.max_keepalive_connections == 200
|
|
|
|
def test_httpx_timeout_config(self):
|
|
"""Test timeout configuration for proxy."""
|
|
timeout = httpx.Timeout(
|
|
connect=10.0,
|
|
read=300.0,
|
|
write=300.0,
|
|
pool=10.0,
|
|
)
|
|
assert timeout.connect == 10.0
|
|
assert timeout.read == 300.0
|
|
assert timeout.write == 300.0
|
|
assert timeout.pool == 10.0
|
|
|
|
def test_async_client_with_limits(self):
|
|
"""Test AsyncClient accepts connection pool limits."""
|
|
|
|
async def _run():
|
|
limits = httpx.Limits(
|
|
max_connections=500,
|
|
max_keepalive_connections=100,
|
|
)
|
|
async with httpx.AsyncClient(
|
|
limits=limits,
|
|
timeout=httpx.Timeout(10.0),
|
|
) as client:
|
|
assert client is not None
|
|
assert limits.max_connections == 500
|
|
assert limits.max_keepalive_connections == 100
|
|
|
|
asyncio.run(_run())
|
|
|
|
|
|
class TestHTTP2Config:
|
|
"""Test HTTP/2 configuration."""
|
|
|
|
def test_http2_requires_h2_package(self):
|
|
"""Test that http2=True requires h2 package."""
|
|
import importlib.util
|
|
|
|
h2_available = importlib.util.find_spec("h2") is not None
|
|
|
|
if h2_available:
|
|
client = httpx.Client(http2=True)
|
|
assert client._base_url is not None
|
|
client.close()
|
|
else:
|
|
with pytest.raises(ImportError):
|
|
httpx.Client(http2=True)
|
|
|
|
def test_async_client_http2(self):
|
|
"""Test AsyncClient with HTTP/2 enabled."""
|
|
import importlib.util
|
|
|
|
if not importlib.util.find_spec("h2"):
|
|
pytest.skip("h2 package not installed")
|
|
|
|
async def _run():
|
|
async with httpx.AsyncClient(
|
|
http2=True,
|
|
limits=httpx.Limits(max_connections=100),
|
|
) as client:
|
|
assert client is not None
|
|
|
|
asyncio.run(_run())
|
|
|
|
|
|
class TestProxyConfigDataclass:
|
|
"""Test ProxyConfig dataclass with new fields."""
|
|
|
|
def test_proxy_config_defaults(self):
|
|
"""Test default values for scalability settings."""
|
|
from dataclasses import dataclass
|
|
|
|
@dataclass
|
|
class ProxyConfigTest:
|
|
"""Minimal proxy config for testing."""
|
|
|
|
host: str = "127.0.0.1"
|
|
port: int = 8787
|
|
request_timeout_seconds: int = 300
|
|
connect_timeout_seconds: int = 10
|
|
max_connections: int = 500
|
|
max_keepalive_connections: int = 100
|
|
http2: bool = True
|
|
|
|
config = ProxyConfigTest()
|
|
assert config.max_connections == 500
|
|
assert config.max_keepalive_connections == 100
|
|
assert config.http2 is True
|
|
|
|
def test_proxy_config_custom_values(self):
|
|
"""Test custom values for scalability settings."""
|
|
from dataclasses import dataclass
|
|
|
|
@dataclass
|
|
class ProxyConfigTest:
|
|
max_connections: int = 500
|
|
max_keepalive_connections: int = 100
|
|
http2: bool = True
|
|
|
|
config = ProxyConfigTest(
|
|
max_connections=1000,
|
|
max_keepalive_connections=200,
|
|
http2=False,
|
|
)
|
|
assert config.max_connections == 1000
|
|
assert config.max_keepalive_connections == 200
|
|
assert config.http2 is False
|
|
|
|
|
|
class TestConcurrencyPatterns:
|
|
"""Test async concurrency patterns used in proxy."""
|
|
|
|
def test_semaphore_for_backpressure(self):
|
|
"""Test semaphore pattern for limiting concurrent requests."""
|
|
|
|
async def _run():
|
|
semaphore = asyncio.Semaphore(3)
|
|
active = []
|
|
completed = []
|
|
|
|
async def task(task_id: int):
|
|
async with semaphore:
|
|
active.append(task_id)
|
|
assert len(active) <= 3
|
|
await asyncio.sleep(0.01)
|
|
active.remove(task_id)
|
|
completed.append(task_id)
|
|
|
|
tasks = [task(i) for i in range(10)]
|
|
await asyncio.gather(*tasks)
|
|
assert len(completed) == 10
|
|
|
|
asyncio.run(_run())
|
|
|
|
def test_connection_reuse_pattern(self):
|
|
"""Test that single client instance is reused (not recreated)."""
|
|
|
|
async def _run():
|
|
clients_created = []
|
|
|
|
class MockProxyWithClient:
|
|
def __init__(self):
|
|
self.http_client = None
|
|
|
|
async def startup(self):
|
|
self.http_client = httpx.AsyncClient(
|
|
limits=httpx.Limits(max_connections=100),
|
|
)
|
|
clients_created.append(self.http_client)
|
|
|
|
async def shutdown(self):
|
|
if self.http_client:
|
|
await self.http_client.aclose()
|
|
|
|
async def make_request(self, url: str):
|
|
return self.http_client
|
|
|
|
proxy = MockProxyWithClient()
|
|
await proxy.startup()
|
|
|
|
client1 = await proxy.make_request("http://example1.com")
|
|
client2 = await proxy.make_request("http://example2.com")
|
|
client3 = await proxy.make_request("http://example3.com")
|
|
|
|
assert client1 is client2 is client3
|
|
assert len(clients_created) == 1
|
|
|
|
await proxy.shutdown()
|
|
|
|
asyncio.run(_run())
|
|
|
|
|
|
class TestTimeoutOverrides:
|
|
"""Test per-request timeout overrides."""
|
|
|
|
def test_request_level_timeout_override(self):
|
|
"""Test that timeout can be overridden per-request."""
|
|
|
|
async def _run():
|
|
async with httpx.AsyncClient(
|
|
timeout=httpx.Timeout(10.0),
|
|
):
|
|
override_timeout = httpx.Timeout(120.0)
|
|
assert override_timeout.read == 120.0
|
|
assert override_timeout.connect == 120.0
|
|
|
|
asyncio.run(_run())
|
|
|
|
|
|
class TestWorkerConfiguration:
|
|
"""Test worker process configuration."""
|
|
|
|
def test_uvicorn_workers_parameter(self):
|
|
"""Test that uvicorn accepts workers parameter."""
|
|
uvicorn = pytest.importorskip("uvicorn")
|
|
|
|
config = uvicorn.Config(
|
|
app="app:app",
|
|
workers=4,
|
|
limit_concurrency=1000,
|
|
)
|
|
assert config.workers == 4
|
|
assert config.limit_concurrency == 1000
|
|
|
|
def test_single_worker_default(self):
|
|
"""Test that default is single worker (None)."""
|
|
uvicorn = pytest.importorskip("uvicorn")
|
|
|
|
config = uvicorn.Config(app="app:app")
|
|
assert config.workers is None or config.workers == 1
|
|
|
|
def test_run_server_uses_import_string_for_multiple_workers(self, monkeypatch):
|
|
from headroom.proxy.models import ProxyConfig
|
|
from headroom.proxy.server import _MULTI_WORKER_CONFIG_ENV, run_server
|
|
|
|
captured = {}
|
|
config = ProxyConfig(host="0.0.0.0", port=8787, max_connections=200)
|
|
|
|
def fake_run(app, **kwargs):
|
|
captured["app"] = app
|
|
captured["kwargs"] = kwargs
|
|
|
|
monkeypatch.delenv(_MULTI_WORKER_CONFIG_ENV, raising=False)
|
|
|
|
try:
|
|
with patch("headroom.proxy.server.uvicorn.run", fake_run):
|
|
run_server(config, workers=4, limit_concurrency=250)
|
|
|
|
assert captured["app"] == "headroom.proxy.server:create_app_from_env"
|
|
assert captured["kwargs"]["workers"] == 4
|
|
assert captured["kwargs"]["limit_concurrency"] == 250
|
|
assert captured["kwargs"]["factory"] is True
|
|
payload = json.loads(os.environ[_MULTI_WORKER_CONFIG_ENV])
|
|
assert payload["host"] == "0.0.0.0"
|
|
assert payload["port"] == 8787
|
|
assert payload["max_connections"] == 200
|
|
finally:
|
|
# run_server sets this via raw os.environ. Pop it directly rather
|
|
# than via monkeypatch.delenv: delenv records the current (JSON)
|
|
# value and re-restores it on teardown, leaking the config into
|
|
# later tests (e.g. _proxy_config_from_env then ignores HEADROOM_*).
|
|
os.environ.pop(_MULTI_WORKER_CONFIG_ENV, None)
|