diff --git a/CHANGELOG.md b/CHANGELOG.md index 793c9a97a..43498c86f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,89 +8,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## Unreleased -### Fixed -- Concurrent large requests no longer 502 on a transient HTTP/2 stream reset. - A single upstream `StreamReset` poisons the shared h2 connection and raises - `RemoteProtocolError` / `LocalProtocolError` on every in-flight request; those - transport errors weren't in the proxy's retry paths, so they collapsed - straight to a 502 with no reconnect. The Anthropic non-streaming and streaming - retry paths now treat any `httpx.TransportError` (including h2 protocol - errors) as retryable before the first client byte, so the bad connection - is dropped and the request re-sent on a fresh one - ([#1639](https://github.com/headroomlabs-ai/headroom/issues/1639)). -- `headroom wrap claude` no longer leaves a dead `ANTHROPIC_BASE_URL` in a - project's `.claude/settings.local.json` after an unclean exit (`SIGKILL`, - OOM, reboot, or terminal/tmux close via `SIGHUP`, which was not caught). - `_write_claude_wrap_base_url`/`_restore_claude_wrap_base_url` only removed - or restored the entry from the wrap process's own `finally` block, so a - crash skipped it and every later bare `claude` invocation in that project - inherited the stale proxy URL and hung indefinitely retrying a dead port. - A wrap session now stamps a sidecar marker (pid, port, prior value); the - next `wrap`, `unwrap`, or `headroom doctor` run detects a marker whose pid - is dead or reused and restores the recorded prior value automatically. - `claude()` also now catches `SIGHUP` alongside the existing `SIGTERM` - handler ([#1768](https://github.com/headroomlabs-ai/headroom/issues/1768)). -- Non-finite values (`NaN`, `Infinity`) in `proxy_savings.json` or in upstream - cost/token metadata no longer crash the proxy or corrupt the savings - dashboard. `SavingsTracker`'s numeric coercion caught only `TypeError` and - `ValueError`, so `int(float('inf'))` raised an uncaught `OverflowError` while - loading persisted state (`SavingsTracker.__init__` failed and the proxy would - not start), and `float('nan')`/`float('inf')` passed straight through, then - serialized to `NaN`/`Infinity` literals that the dashboard's `JSON.parse` - rejects. `json.loads` accepts those literals, so one bad write poisoned every - later start. Both coercion helpers now also catch `OverflowError` and reject - non-finite floats, failing open to safe defaults. -- `headroom learn` now honors `CLAUDE_CONFIG_DIR`. It resolved the Claude - config directory as `~/.claude` and wrote global memory to - `~/.claude/CLAUDE.md`, so users who relocate their Claude config via that - env var had `learn` scan the wrong directory and detect no projects. The - scanner and memory writer now read/write the configured directory - ([#1630](https://github.com/headroomlabs-ai/headroom/issues/1630)). -- `--backend bedrock` now fails fast with an actionable error when temporary - AWS credentials (`AWS_SESSION_TOKEN`) are used but botocore is not installed - (e.g. the slim default Docker image). litellm's session-token auth path - imports botocore, so the missing dependency previously surfaced only at - request time as a misleading `authentication_error: No module named - 'botocore'`. The proxy now tells the user to install the `bedrock` extra up - front ([#1551](https://github.com/headroomlabs-ai/headroom/issues/1551)). -- Content detection no longer crashes the proxy on text containing an - orphaned `+++ ` target line with no preceding `--- ` source line (common in - `set -x` xtrace output and partial diffs). The bundled `unidiff` 0.4.0 parser - panics on that input instead of returning an error; the Rust diff detector now - contains the panic and treats the fragment as plain text, so the request is - compressed and forwarded normally instead of returning HTTP 500 - ([#1547](https://github.com/headroomlabs-ai/headroom/issues/1547)). -* **proxy:** persist lifetime cache-read savings (tokens + USD) in `proxy_savings.json` (schema v4, additive) so cache-mode savings survive proxy restarts and upgrades. Previously prefix-cache read savings lived only in process memory and every restart reset the dashboard's cache figure to zero; the "Cache Reads (lifetime)" tile now reads the persisted value and the Prefix Cache Impact card renders after a restart with zero traffic, marking session-scoped tiles "no activity since restart". -- Proactive expansion blocks injected into user turns are now wrapped in - `` XML tags, giving downstream consumers - (LLMs, loggers, attribution parsers) a machine-readable provenance - boundary and preventing misattribution in multi-agent threads. -- **cli:** the startup banner no longer advertises - `HEADROOM_COMPRESSION_STABLE_AFTER_TURN` and - `HEADROOM_STALE_READ_COMPRESS_AFTER_TURNS` as tuning knobs. Both were read - only to render the `Performance Tuning` banner section and were never wired - into the compression path, so setting them changed the banner but had no - effect on behavior. The banner now surfaces only the embedding sidecar, - which is a real, consumed setting. -- **memory/embedder:** cap CPU thread oversubscription in the local - torch/sentence-transformers embedder. Concurrent encodes previously each - fanned out to ~`os.cpu_count()` BLAS/OpenMP threads, so under load the memory - path starved the asyncio event loop and spiked `/livez` latency to several - seconds. CPU encodes now run on a dedicated, size-limited executor whose - workers each pin their thread pool, bounding total embedding threads to - `HEADROOM_EMBED_CONCURRENCY` × `HEADROOM_EMBED_NUM_THREADS` (defaults - `min(4, cpu)` × 1). The ONNX embedder already capped its threads; this brings - the torch path to parity - ([#198](https://github.com/headroomlabs-ai/headroom/issues/198)). -- Buffered passthrough routes (e.g. `GET /v1/models`) no longer return an - opaque HTTP 502 when an OpenAI-compatible upstream closes a pooled - keep-alive connection mid-response (`httpx.RemoteProtocolError` / - "incomplete chunked read"). Headroom now retries the request once on a - fresh connection — mirroring a direct `curl` — and only returns a clear - `upstream_protocol_error` 502 if the upstream is genuinely sending an - incomplete response - ([#1112](https://github.com/chopratejas/headroom/issues/1112)). - ### Changed * **telemetry:** anonymous usage telemetry is now **opt-in** (off by default) instead of opt-out. Nothing is collected or sent unless you set `HEADROOM_TELEMETRY=on` or pass `--telemetry` to `headroom proxy` / `headroom install apply`. `is_telemetry_enabled()` is fail-closed — only explicit on-values (`on`/`true`/`1`/`yes`/`enable`/`enabled`) enable it; unset, empty, or unrecognized values stay disabled. The existing `--no-telemetry` flag and `HEADROOM_TELEMETRY=off` remain accepted for back-compat, and install manifests now write the `HEADROOM_TELEMETRY` value explicitly so generated deployments are unambiguous. @@ -98,11 +15,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Features -* **proxy:** add provider-only HTTP proxy routing via `--http-proxy` and - `HEADROOM_HTTP_PROXY`. Upstream LLM provider calls can now use an HTTP proxy - without setting process-wide `HTTP_PROXY`/`HTTPS_PROXY` variables that are - inherited by tool executions; proxied provider clients use HTTP/1.1 so HTTPS - provider APIs can tunnel through CONNECT. +* **proxy:** add provider-only HTTP proxy routing via `--http-proxy` and `HEADROOM_HTTP_PROXY`. Upstream LLM provider calls can now use an HTTP proxy without setting process-wide `HTTP_PROXY`/`HTTPS_PROXY` variables that are inherited by tool executions; proxied provider clients use HTTP/1.1 so HTTPS provider APIs can tunnel through CONNECT. * **proxy:** add output shaping for OpenAI Responses traffic on `/v1/responses` HTTP requests and Codex WebSocket `response.create` frames, with stable output-savings holdout keys and counted WS token strata for the experiment. * **wrap:** `headroom wrap claude --1m` preserves the 1M context window. Behind a custom `ANTHROPIC_BASE_URL` (the proxy) Claude Code drops the `context-1m` beta header and caps the window at 200k for entitled subscription users; the opt-in flag sets `ANTHROPIC_MODEL=[1m]` on the launched process so the 1M window activates through Headroom. A model already selected via `ANTHROPIC_MODEL` is preserved (only the `[1m]` suffix is appended) ([#1158](https://github.com/chopratejas/headroom/issues/1158)). * **learn:** weight loops in `headroom learn`. A new loop detector (`headroom/learn/loops.py`) recognizes repeated tool-call patterns — including RTK re-fetch loops, where RTK's output truncation makes the agent re-run larger-limit variants of a *successful* command — collapses output-limit variants to one signature, measures the wasted tokens, surfaces loops as a highest-priority digest section, and weights loop guardrails above one-off rules by their measured waste. Previously loops had no special weight and a no-failure re-fetch loop was skipped entirely. Adds an RTK-loop eval (`benchmarks/rtk_loop_learn_eval.py`) that reproduces a loop, runs it through Learn, and asserts the generated guardrail ranks first and prevents re-triggering. @@ -171,7 +84,18 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 * **install (macOS):** fix `headroom install restart` / `install start` for launchd `persistent-service` deployments. `stop` `bootout`s the job but `start` only ran `launchctl kickstart`, which cannot recover the un-bootstrapped state `stop`/`restart` leave behind (launchctl error 113), so the proxy was left stopped. `start` now tries `kickstart` (fast path for an already-bootstrapped job) and, on failure, `bootstrap`s the plist fresh — retrying for ~15s to ride out the transient `bootstrap` EIO (error 5) window while launchd releases the label after a `bootout`. `stop` tolerates only the already-absent case (`bootout` ESRCH / error 3) and still raises on any other `bootout` failure ([#1289](https://github.com/headroomlabs-ai/headroom/issues/1289)). * **wrap:** isolate wrapped proxy subprocess stdout/stderr into `proxy-stdio.log`, so `proxy.log` remains the canonical rotating runtime log and Windows rollover failures from `RotatingFileHandler` are no longer blocked by wrapper stdio handles ([#1184](https://github.com/chopratejas/headroom/issues/1184)). * **langchain:** fix `HeadroomChatModel.ainvoke()` crashing with `AttributeError: 'AsyncStream' object has no attribute 'model_dump'` when the wrapped model has `streaming=True`. `_agenerate()` now uses a per-call non-streaming copy of the wrapped model instead of mutating shared state across an `await` ([#1285](https://github.com/headroomlabs-ai/headroom/issues/1285)). - +* **proxy:** a transient rtk/lean-ctx stat-read failure (timeout, non-zero exit, bad JSON) no longer corrupts the dashboard's CLI-filtering session metrics. Failed reads now return "no data" instead of a synthetic zero payload, and the session baseline is only ever pinned from successful installed-tool reads — previously one hiccup re-pinned the baseline to zero and the next successful read inflated session savings by the tool's entire lifetime, at every proxy boot and `POST /stats/reset`. +* **proxy:** Concurrent large requests no longer 502 on a transient HTTP/2 stream reset. A single upstream `StreamReset` poisons the shared h2 connection and raises `RemoteProtocolError` / `LocalProtocolError` on every in-flight request; those transport errors weren't in the proxy's retry paths, so they collapsed straight to a 502 with no reconnect. The Anthropic non-streaming and streaming retry paths now treat any `httpx.TransportError` (including h2 protocol errors) as retryable before the first client byte, so the bad connection is dropped and the request re-sent on a fresh one ([#1639](https://github.com/headroomlabs-ai/headroom/issues/1639)). +* **install:** `headroom wrap claude` no longer leaves a dead `ANTHROPIC_BASE_URL` in a project's `.claude/settings.local.json` after an unclean exit (`SIGKILL`, OOM, reboot, or terminal/tmux close via `SIGHUP`, which was not caught). `_write_claude_wrap_base_url`/`_restore_claude_wrap_base_url` only removed or restored the entry from the wrap process's own `finally` block, so a crash skipped it and every later bare `claude` invocation in that project inherited the stale proxy URL and hung indefinitely retrying a dead port. A wrap session now stamps a sidecar marker (pid, port, prior value); the next `wrap`, `unwrap`, or `headroom doctor` run detects a marker whose pid is dead or reused and restores the recorded prior value automatically. `claude()` also now catches `SIGHUP` alongside the existing `SIGTERM` handler ([#1768](https://github.com/headroomlabs-ai/headroom/issues/1768)). +* **proxy:** Non-finite values (`NaN`, `Infinity`) in `proxy_savings.json` or in upstream cost/token metadata no longer crash the proxy or corrupt the savings dashboard. `SavingsTracker`'s numeric coercion caught only `TypeError` and `ValueError`, so `int(float('inf'))` raised an uncaught `OverflowError` while loading persisted state (`SavingsTracker.__init__` failed and the proxy would not start), and `float('nan')`/`float('inf')` passed straight through, then serialized to `NaN`/`Infinity` literals that the dashboard's `JSON.parse` rejects. `json.loads` accepts those literals, so one bad write poisoned every later start. Both coercion helpers now also catch `OverflowError` and reject non-finite floats, failing open to safe defaults. +* **learn:** `headroom learn` now honors `CLAUDE_CONFIG_DIR`. It resolved the Claude config directory as `~/.claude` and wrote global memory to `~/.claude/CLAUDE.md`, so users who relocate their Claude config via that env var had `learn` scan the wrong directory and detect no projects. The scanner and memory writer now read/write the configured directory ([#1630](https://github.com/headroomlabs-ai/headroom/issues/1630)). +* **cli:** `--backend bedrock` now fails fast with an actionable error when temporary AWS credentials (`AWS_SESSION_TOKEN`) are used but botocore is not installed (e.g. the slim default Docker image). litellm's session-token auth path imports botocore, so the missing dependency previously surfaced only at request time as a misleading `authentication_error: No module named 'botocore'`. The proxy now tells the user to install the `bedrock` extra up front ([#1551](https://github.com/headroomlabs-ai/headroom/issues/1551)). +* **compression:** Content detection no longer crashes the proxy on text containing an orphaned `+++ ` target line with no preceding `--- ` source line (common in `set -x` xtrace output and partial diffs). The bundled `unidiff` 0.4.0 parser panics on that input instead of returning an error; the Rust diff detector now contains the panic and treats the fragment as plain text, so the request is compressed and forwarded normally instead of returning HTTP 500 ([#1547](https://github.com/headroomlabs-ai/headroom/issues/1547)). +* **proxy:** persist lifetime cache-read savings (tokens + USD) in `proxy_savings.json` (schema v4, additive) so cache-mode savings survive proxy restarts and upgrades. Previously prefix-cache read savings lived only in process memory and every restart reset the dashboard's cache figure to zero; the "Cache Reads (lifetime)" tile now reads the persisted value and the Prefix Cache Impact card renders after a restart with zero traffic, marking session-scoped tiles "no activity since restart". +* **compression:** Proactive expansion blocks injected into user turns are now wrapped in`` XML tags, giving downstream consumers (LLMs, loggers, attribution parsers) a machine-readable provenance boundary and preventing misattribution in multi-agent threads. +* **cli:** the startup banner no longer advertises `HEADROOM_COMPRESSION_STABLE_AFTER_TURN` and `HEADROOM_STALE_READ_COMPRESS_AFTER_TURNS` as tuning knobs. Both were read only to render the `Performance Tuning` banner section and were never wired into the compression path, so setting them changed the banner but had no effect on behavior. The banner now surfaces only the embedding sidecar, which is a real, consumed setting. +* **memory/embedder:** cap CPU thread oversubscription in the local torch/sentence-transformers embedder. Concurrent encodes previously each fanned out to ~`os.cpu_count()` BLAS/OpenMP threads, so under load the memory path starved the asyncio event loop and spiked `/livez` latency to several seconds. CPU encodes now run on a dedicated, size-limited executor whose workers each pin their thread pool, bounding total embedding threads to `HEADROOM_EMBED_CONCURRENCY` × `HEADROOM_EMBED_NUM_THREADS` (defaults `min(4, cpu)` × 1). The ONNX embedder already capped its threads; this brings the torch path to parity ([#198](https://github.com/headroomlabs-ai/headroom/issues/198)). +* **proxy:** Buffered passthrough routes (e.g. `GET /v1/models`) no longer return an opaque HTTP 502 when an OpenAI-compatible upstream closes a pooled keep-alive connection mid-response (`httpx.RemoteProtocolError` / "incomplete chunked read"). Headroom now retries the request once on a fresh connection — mirroring a direct `curl` — and only returns a clear `upstream_protocol_error` 502 if the upstream is genuinely sending an incomplete response ([#1112](https://github.com/chopratejas/headroom/issues/1112)). ## [0.29.0](https://github.com/headroomlabs-ai/headroom/compare/v0.28.0...v0.29.0) (2026-07-03) diff --git a/headroom/proxy/helpers.py b/headroom/proxy/helpers.py index 2918c9b02..ade8f109f 100644 --- a/headroom/proxy/helpers.py +++ b/headroom/proxy/helpers.py @@ -1287,35 +1287,25 @@ def _read_rtk_lifetime_stats() -> dict[str, Any] | None: summary=summary if isinstance(summary, dict) else {}, ) else: - # PR-G2 remediation (H2): structured log the synthetic-zero path - # so downstream consumers (subscription tracker, dashboards) can - # distinguish a healthy "RTK ran and saved nothing" from a broken - # "RTK failed and we faked zero". + # A failed read is "no data", never a zero counter — a synthetic + # zero here re-pins the session baseline and inflates session + # savings by the tool's whole lifetime on recovery. stderr_excerpt = (result.stderr or "")[:200] logger.warning( "event=rtk_stats_subprocess_failed reason=non_zero_exit rc=%s stderr=%r", result.returncode, stderr_excerpt, ) - return _context_tool_zero_payload( - tool=_CONTEXT_TOOL_RTK, - installed=True, - scope=scope, - ) + return None except Exception as exc: - # PR-G2 remediation (H2): log the exception path too. Reason is the - # exception class name (without payload — RTK exceptions can carry - # filesystem paths). + # Reason is the exception class name (without payload — RTK + # exceptions can carry filesystem paths). logger.warning( "event=rtk_stats_subprocess_failed reason=%s error=%s", type(exc).__name__, exc, ) - return _context_tool_zero_payload( - tool=_CONTEXT_TOOL_RTK, - installed=True, - scope=scope, - ) + return None return payload @@ -1329,8 +1319,6 @@ def _read_lean_ctx_lifetime_stats() -> dict[str, Any] | None: if not lean_ctx_path: return _context_tool_zero_payload(tool=_CONTEXT_TOOL_LEAN_CTX, installed=False) - base_payload = _context_tool_zero_payload(tool=_CONTEXT_TOOL_LEAN_CTX, installed=True) - try: result = run( [str(lean_ctx_path), "gain", "--json"], @@ -1338,21 +1326,32 @@ def _read_lean_ctx_lifetime_stats() -> dict[str, Any] | None: text=True, timeout=5, ) + # Failed reads return None ("no data") — mirrors the rtk reader so + # the baseline logic never sees synthetic zeros from either tool. if result.returncode != 0 or not result.stdout.strip(): - return dict(base_payload) + logger.warning( + "event=lean_ctx_stats_subprocess_failed reason=non_zero_exit rc=%s", + result.returncode, + ) + return None data = json.loads(result.stdout) summary = data.get("summary", data) if isinstance(data, dict) else {} if not isinstance(summary, dict): - return dict(base_payload) + logger.warning("event=lean_ctx_stats_subprocess_failed reason=bad_payload") + return None return _context_tool_summary_payload( tool=_CONTEXT_TOOL_LEAN_CTX, installed=True, summary=summary, ) - except Exception: - return dict(base_payload) + except Exception as exc: + logger.warning( + "event=lean_ctx_stats_subprocess_failed reason=%s", + type(exc).__name__, + ) + return None def _read_context_tool_lifetime_stats(tool: str) -> dict[str, Any] | None: @@ -1367,18 +1366,36 @@ async def initialize_context_tool_session_baseline() -> None: tool = _selected_context_tool() payload = await asyncio.to_thread(_read_context_tool_lifetime_stats, tool) with _context_tool_stats_cache_lock: - _context_tool_session_baseline.update( - { - "initialized": True, - "tool": tool, - "total_commands": int((payload or {}).get("total_commands", 0) or 0), - "input_tokens": int((payload or {}).get("input_tokens", 0) or 0), - "output_tokens": int((payload or {}).get("output_tokens", 0) or 0), - "tokens_saved": int((payload or {}).get("tokens_saved", 0) or 0), - "total_time_ms": int((payload or {}).get("total_time_ms", 0) or 0), - "captured_at": time.time(), - } - ) + if payload is None or not payload.get("installed", False): + # Failed or tool-absent read: defer the pin to the first + # successful read (guarded lazy-init) — pinning zeros here would + # inflate session savings by the tool's whole lifetime once it + # recovers or gets installed. + _context_tool_session_baseline.update( + { + "initialized": False, + "tool": tool, + "total_commands": 0, + "input_tokens": 0, + "output_tokens": 0, + "tokens_saved": 0, + "total_time_ms": 0, + "captured_at": time.time(), + } + ) + else: + _context_tool_session_baseline.update( + { + "initialized": True, + "tool": tool, + "total_commands": int(payload.get("total_commands", 0) or 0), + "input_tokens": int(payload.get("input_tokens", 0) or 0), + "output_tokens": int(payload.get("output_tokens", 0) or 0), + "tokens_saved": int(payload.get("tokens_saved", 0) or 0), + "total_time_ms": int(payload.get("total_time_ms", 0) or 0), + "captured_at": time.time(), + } + ) _context_tool_stats_cache.update( { "expires_at": 0.0, @@ -1416,19 +1433,28 @@ def _get_context_tool_stats() -> dict[str, Any] | None: payload = _read_context_tool_lifetime_stats(tool) with _context_tool_stats_cache_lock: + # Baseline mutations only happen on successful reads from an + # installed tool — a failed read (None) or a tool-absent zero payload + # must never pin or re-pin, or session deltas inflate by the whole + # lifetime when the tool comes back. + tool_installed = payload is not None and bool(payload.get("installed", False)) if ( - not _context_tool_session_baseline["initialized"] - or _context_tool_session_baseline.get("tool") != tool + payload is not None + and tool_installed + and ( + not _context_tool_session_baseline["initialized"] + or _context_tool_session_baseline.get("tool") != tool + ) ): _context_tool_session_baseline.update( { "initialized": True, "tool": tool, - "total_commands": int((payload or {}).get("total_commands", 0) or 0), - "input_tokens": int((payload or {}).get("input_tokens", 0) or 0), - "output_tokens": int((payload or {}).get("output_tokens", 0) or 0), - "tokens_saved": int((payload or {}).get("tokens_saved", 0) or 0), - "total_time_ms": int((payload or {}).get("total_time_ms", 0) or 0), + "total_commands": int(payload.get("total_commands", 0) or 0), + "input_tokens": int(payload.get("input_tokens", 0) or 0), + "output_tokens": int(payload.get("output_tokens", 0) or 0), + "tokens_saved": int(payload.get("tokens_saved", 0) or 0), + "total_time_ms": int(payload.get("total_time_ms", 0) or 0), "captured_at": time.time(), } ) @@ -1444,7 +1470,10 @@ def _get_context_tool_stats() -> dict[str, Any] | None: baseline_output_tokens = int(_context_tool_session_baseline["output_tokens"]) baseline_tokens_saved = int(_context_tool_session_baseline["tokens_saved"]) baseline_total_time_ms = int(_context_tool_session_baseline["total_time_ms"]) - counter_reset_detected = ( + # A tool-absent payload carries zero counters that are not a + # genuine external reset — only successful installed reads may + # re-pin the baseline. + counter_reset_detected = tool_installed and ( lifetime_total_commands < baseline_total_commands or lifetime_input_tokens < baseline_input_tokens or lifetime_output_tokens < baseline_output_tokens diff --git a/headroom/subscription/tracker.py b/headroom/subscription/tracker.py index 4cccf79ea..f4f974514 100644 --- a/headroom/subscription/tracker.py +++ b/headroom/subscription/tracker.py @@ -356,8 +356,9 @@ class SubscriptionTracker(QuotaTracker): Returns ``0`` (never negative) when: - ``HEADROOM_RTK_WIRING=disabled`` — operator opt-out. - - ``_get_rtk_stats()`` returns ``None`` — RTK not selected / not - installed; explicit zero is the right answer. + - ``_get_rtk_stats()`` returns ``None`` — RTK not selected, or the + stat read failed this poll ("no data"); explicit zero contribution + is the right answer and the high-water mark is preserved. - ``_get_rtk_stats()`` raises — transient error, logged loudly. - The session counter regressed (RTK reset / new project) — that path also re-baselines ``_last_rtk_tokens_saved`` to the new @@ -429,7 +430,8 @@ class SubscriptionTracker(QuotaTracker): # the pre-Headroom RTK history. Falls back to the top-level # ``tokens_saved`` (which is also session-scoped in the canonical # payload; see ``_get_context_tool_stats``) and finally to 0 for - # synthetic-zero payloads. + # not-installed zero payloads (failed reads arrive as ``None`` and + # returned above). session_payload = stats.get("session") if isinstance(session_payload, dict) and "tokens_saved" in session_payload: current_total_raw = session_payload.get("tokens_saved", 0) diff --git a/tests/test_rtk_session_savings.py b/tests/test_rtk_session_savings.py index 820a7f922..6b65ee6cc 100644 --- a/tests/test_rtk_session_savings.py +++ b/tests/test_rtk_session_savings.py @@ -68,3 +68,242 @@ def test_session_savings_is_delta_not_lifetime_average(monkeypatch): # Session % is derived from the delta (200 saved / 300 input delta), not the # lifetime-diluted average. assert second["session"]["savings_pct"] == round(200 / 300 * 100, 4) + + +# --- Failure semantics: a failed read is "no data", never a zero counter --- + + +def _fake_run_raises(*args, **kwargs): + import subprocess + + raise subprocess.TimeoutExpired(cmd="rtk", timeout=5) + + +def test_rtk_reader_returns_none_on_timeout(monkeypatch): + import headroom.rtk as rtk_mod + + monkeypatch.setattr(rtk_mod, "get_rtk_path", lambda: "/fake/rtk") + monkeypatch.setattr(helpers, "run", _fake_run_raises) + assert helpers._read_rtk_lifetime_stats() is None + + +def test_rtk_reader_returns_none_on_nonzero_exit(monkeypatch, caplog): + import logging + from types import SimpleNamespace + + import headroom.rtk as rtk_mod + + monkeypatch.setattr(rtk_mod, "get_rtk_path", lambda: "/fake/rtk") + monkeypatch.setattr( + helpers, + "run", + lambda *a, **k: SimpleNamespace(returncode=1, stdout="", stderr="boom"), + ) + with caplog.at_level(logging.WARNING): + assert helpers._read_rtk_lifetime_stats() is None + assert "rtk_stats_subprocess_failed" in caplog.text + + +def test_rtk_reader_returns_none_on_bad_json(monkeypatch): + from types import SimpleNamespace + + import headroom.rtk as rtk_mod + + monkeypatch.setattr(rtk_mod, "get_rtk_path", lambda: "/fake/rtk") + monkeypatch.setattr( + helpers, + "run", + lambda *a, **k: SimpleNamespace(returncode=0, stdout="not-json{", stderr=""), + ) + assert helpers._read_rtk_lifetime_stats() is None + + +def test_rtk_reader_not_installed_keeps_zero_payload(monkeypatch): + import headroom.rtk as rtk_mod + + monkeypatch.setattr(rtk_mod, "get_rtk_path", lambda: None) + payload = helpers._read_rtk_lifetime_stats() + assert payload is not None + assert payload["installed"] is False + assert payload["tokens_saved"] == 0 + + +def test_lean_ctx_reader_returns_none_on_failure_and_logs(monkeypatch, caplog): + import logging + + import headroom.lean_ctx as lean_mod + + monkeypatch.setattr(lean_mod, "get_lean_ctx_path", lambda: "/fake/lean-ctx") + monkeypatch.setattr(helpers, "run", _fake_run_raises) + with caplog.at_level(logging.WARNING): + assert helpers._read_lean_ctx_lifetime_stats() is None + assert "stats_subprocess_failed" in caplog.text + + +# --- Baseline guards: pin only from successful installed-tool reads --- + + +def _payload(saved: int, *, input_tokens: int = 1000, installed: bool = True): + return helpers._context_tool_summary_payload( + tool="rtk", + installed=installed, + scope="global", + summary={"total_input": input_tokens, "total_output": 400, "total_saved": saved}, + ) + + +def _stub_reads(monkeypatch, sequence): + calls = {"n": 0} + + def fake(tool): + idx = min(calls["n"], len(sequence) - 1) + calls["n"] += 1 + item = sequence[idx] + return item() if callable(item) else item + + monkeypatch.setattr(helpers, "_read_context_tool_lifetime_stats", fake) + return calls + + +def test_transient_failure_does_not_repin_baseline_or_inflate_session(monkeypatch): + """The headline regression, end to end through the real reader. + + A transient rtk subprocess failure between two identical successful reads + must not re-pin the session baseline; today the reader converts the + failure into a synthetic zero payload and recovery reports the tool's + entire lifetime as session savings. + """ + import json as json_mod + from types import SimpleNamespace + + import headroom.rtk as rtk_mod + + _reset(monkeypatch) + monkeypatch.setattr(rtk_mod, "get_rtk_path", lambda: "/fake/rtk") + + good = json_mod.dumps( + {"summary": {"total_input": 1000, "total_output": 400, "total_saved": 600}} + ) + behaviors = [ + lambda: SimpleNamespace(returncode=0, stdout=good, stderr=""), + _fake_run_raises, + lambda: SimpleNamespace(returncode=0, stdout=good, stderr=""), + ] + calls = {"n": 0} + + def fake_run(*args, **kwargs): + behavior = behaviors[min(calls["n"], len(behaviors) - 1)] + calls["n"] += 1 + return behavior() + + monkeypatch.setattr(helpers, "run", fake_run) + + first = helpers._get_context_tool_stats() + assert first["session"]["tokens_saved"] == 0 + assert first["lifetime"]["tokens_saved"] == 600 + + _bust_cache() + helpers._get_context_tool_stats() + # Baseline survives the failed poll untouched. + assert helpers._context_tool_session_baseline["tokens_saved"] == 600 + + _bust_cache() + recovered = helpers._get_context_tool_stats() + # Recovery must NOT report the full lifetime as session savings. + assert recovered["session"]["tokens_saved"] == 0 + assert recovered["counter_reset_detected"] is False + + +def test_boot_fail_then_poll_fail_never_pins_zero_baseline(monkeypatch): + import asyncio + + _reset(monkeypatch) + _stub_reads(monkeypatch, [None, None, _payload(600)]) + + asyncio.run(helpers.initialize_context_tool_session_baseline()) + assert helpers._context_tool_session_baseline["initialized"] is False + + _bust_cache() + assert helpers._get_context_tool_stats() is None + # Lazy-init must not have pinned zeros from the failed poll. + assert helpers._context_tool_session_baseline["initialized"] is False + + _bust_cache() + recovered = helpers._get_context_tool_stats() + assert recovered["session"]["tokens_saved"] == 0 + assert recovered["lifetime"]["tokens_saved"] == 600 + + +def test_stats_reset_with_failing_read_defers_to_next_success(monkeypatch): + import asyncio + + _reset(monkeypatch) + _stub_reads(monkeypatch, [_payload(600), None, _payload(650)]) + + first = helpers._get_context_tool_stats() + assert first["session"]["tokens_saved"] == 0 + + # /stats/reset while rtk is down: old baseline dropped, pin deferred. + asyncio.run(helpers.initialize_context_tool_session_baseline()) + assert helpers._context_tool_session_baseline["initialized"] is False + + _bust_cache() + after = helpers._get_context_tool_stats() + # First successful read after the deferred reset pins fresh: delta 0. + assert after["session"]["tokens_saved"] == 0 + assert after["lifetime"]["tokens_saved"] == 650 + + +def test_genuine_counter_reset_still_repins(monkeypatch): + _reset(monkeypatch) + _stub_reads(monkeypatch, [_payload(600), _payload(50)]) + + helpers._get_context_tool_stats() + _bust_cache() + second = helpers._get_context_tool_stats() + assert second["counter_reset_detected"] is True + assert second["session"]["tokens_saved"] == 0 + assert second["lifetime"]["tokens_saved"] == 50 + + +def test_not_installed_payload_does_not_repin_baseline(monkeypatch): + _reset(monkeypatch) + _stub_reads( + monkeypatch, + [_payload(600), _payload(0, input_tokens=0, installed=False), _payload(600)], + ) + + helpers._get_context_tool_stats() + _bust_cache() + absent = helpers._get_context_tool_stats() + # Tool vanished at resolution time: honest zeros display, baseline intact. + assert absent["installed"] is False + assert helpers._context_tool_session_baseline["tokens_saved"] == 600 + + _bust_cache() + back = helpers._get_context_tool_stats() + assert back["session"]["tokens_saved"] == 0 + assert back["counter_reset_detected"] is False + + +def test_tool_switch_with_failing_first_read_does_not_zero_pin(monkeypatch): + _reset(monkeypatch) + _stub_reads(monkeypatch, [_payload(600), None]) + + helpers._get_context_tool_stats() + monkeypatch.setenv("HEADROOM_CONTEXT_TOOL", "lean-ctx") + _bust_cache() + assert helpers._get_context_tool_stats() is None + # Switching tools with a failing first read must not pin a zero baseline + # for the new tool. + assert helpers._context_tool_session_baseline.get("tool") != "lean-ctx" + + +def test_failed_poll_caches_none_for_ttl(monkeypatch): + _reset(monkeypatch) + calls = _stub_reads(monkeypatch, [None]) + + assert helpers._get_context_tool_stats() is None + assert helpers._get_context_tool_stats() is None + # Second call inside the TTL is served from cache — no re-read storm. + assert calls["n"] == 1 diff --git a/tests/test_subscription_tracker_rtk_wired.py b/tests/test_subscription_tracker_rtk_wired.py index d60730986..dcf46e045 100644 --- a/tests/test_subscription_tracker_rtk_wired.py +++ b/tests/test_subscription_tracker_rtk_wired.py @@ -578,10 +578,11 @@ def test_rtk_subprocess_failure_logs_structured_warning( mock_warning = MagicMock() monkeypatch.setattr(_helpers.logger, "warning", mock_warning) + # Failed reads return None ("no data") rather than a synthetic zero + # payload — the zero re-pinned the session baseline and inflated session + # savings by the tool's whole lifetime on recovery. payload = _helpers._read_rtk_lifetime_stats() - assert payload is not None - assert payload["scope"] == "global" - assert payload["tokens_saved"] == 0 + assert payload is None # Concatenate all warning call args so the failure message shows what # the helper actually emitted (debug aid for CI flakes). @@ -636,3 +637,32 @@ def test_multi_worker_only_one_polls(monkeypatch: pytest.MonkeyPatch, tmp_path: # Cleanup so subsequent tests don't see a stale lock. tracker_a._release_rtk_poll_lock() tracker_b._release_rtk_poll_lock() + + +def test_rtk_stats_mid_window_failure_preserves_high_water_mark( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A failed poll (None) mid-window must not reset the high-water mark. + + Failed stat reads now arrive as None ("no data"); the recovery poll's + delta is computed against the preserved mark, so no phantom contribution + lands and nothing is lost. + """ + + tracker = _build_tracker(monkeypatch) + monkeypatch.delenv(tracker_module._RTK_WIRING_ENV, raising=False) + _stub_rtk_stats(monkeypatch, [_session_payload(100), None, _session_payload(150)]) + + tracker.update_contribution() + assert tracker._state.contribution.tokens_saved_rtk == 100 + assert tracker._last_rtk_tokens_saved == 100 + + tracker.update_contribution() + # Outage poll: zero contribution, mark preserved. + assert tracker._state.contribution.tokens_saved_rtk == 100 + assert tracker._last_rtk_tokens_saved == 100 + + tracker.update_contribution() + # Recovery: only the true delta lands. + assert tracker._state.contribution.tokens_saved_rtk == 150 + assert tracker._last_rtk_tokens_saved == 150