From f03021f1b69ec1a099436a5f80e68d5266cad8bf Mon Sep 17 00:00:00 2001 From: inix <62450194+inix-x@users.noreply.github.com> Date: Wed, 24 Jun 2026 22:43:06 +0800 Subject: [PATCH] fix(subscription): run transcript token scan off the event loop (#1263) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Description The subscription tracker's poll loop scans Claude Code transcripts to compute window-token usage. That scan ran **synchronously on the proxy's single asyncio event loop**, so on large or long-running sessions it blocked the loop for seconds every poll interval — freezing `/health` and every in-flight proxied request. This moves the scan off the loop with `asyncio.to_thread`. Closes # ## Type of Change - [x] Bug fix (non-breaking change that fixes an issue) - [x] Performance improvement ## Changes Made - `headroom/subscription/tracker.py` — `_maybe_poll()` now calls `await asyncio.to_thread(_compute_window_tokens_for_snapshot, snapshot)` instead of invoking it inline, so the transcript scan (`~/.claude/projects/**/*.jsonl` read + `json.loads` per line) no longer runs on the event-loop thread. The computed result is wired through unchanged. - `tests/test_subscription_tracker.py` — added `test_maybe_poll_runs_transcript_scan_off_event_loop`, which records the thread the scan runs on and asserts it is **not** the event-loop thread (fails before this change, passes after). - `CHANGELOG.md` — Unreleased → Bug Fixes entry. ## Root Cause Captured with `faulthandler` (`SIGUSR1`) during a live wedge — the event loop frozen mid-`json.loads`: ``` Current thread (most recent call first): File ".../python3.14/json/decoder.py", line 361 in raw_decode File ".../python3.14/json/__init__.py", line 352 in loads File ".../headroom/subscription/session_tracking.py", line 127 in compute_window_tokens File ".../headroom/subscription/tracker.py", line 872 in _compute_window_tokens_for_snapshot File ".../headroom/subscription/tracker.py", line 731 in _maybe_poll File ".../headroom/subscription/tracker.py", line 693 in _poll_loop File ".../python3.14/asyncio/events.py", line 94 in _run ``` `_poll_loop` fires every `poll_interval_s` (default **300s**); `compute_window_tokens` reads **every** `~/.claude/projects/**/*.jsonl` transcript and `json.loads` each line. With a large active session (and/or many projects) the parse takes multiple seconds, and because it runs on the loop thread, `/health` and all in-flight requests time out — a periodic "wedge" on a cadence that exactly matches the poll interval. ## Testing - [x] Unit tests pass (`pytest`) - [x] Linting passes (`ruff check .`) - [x] Type checking passes (`mypy headroom`) - [x] New tests added for new functionality - [x] Manual testing performed ### Test Output ```text $ uv run ruff check headroom/subscription/tracker.py tests/test_subscription_tracker.py All checks passed! $ uv run ruff format --check headroom/subscription/tracker.py tests/test_subscription_tracker.py 2 files already formatted $ uv run mypy headroom/subscription/tracker.py Success: no issues found in 1 source file $ uv run pytest tests/test_subscription_tracker.py -q ...... [100%] 6 passed in 0.42s # Regression test fails before the fix, passes after: $ git stash push -- headroom/subscription/tracker.py # remove the fix $ uv run pytest tests/test_subscription_tracker.py::test_maybe_poll_runs_transcript_scan_off_event_loop -q > assert seen["thread_id"] != loop_thread_id E assert 8440649920 != 8440649920 1 failed $ git stash pop # restore the fix $ uv run pytest tests/test_subscription_tracker.py::test_maybe_poll_runs_transcript_scan_off_event_loop -q 1 passed ``` ## Real Behavior Proof - **Environment:** macOS (Darwin 25), Python 3.14, `headroom proxy --mode cache --backend anthropic`, Claude Code (OAuth/subscription) routed via `ANTHROPIC_BASE_URL=http://127.0.0.1:8787`, a large, long-running ~1M-token session. - **Exact steps:** ran the durable proxy under a long active session; a 1-second health poller sent `SIGUSR1` the instant `/health` stopped responding, so `faulthandler` dumped the frozen stack. Confirmed the captured frame above. Then ran with the scan offloaded (`_compute_window_tokens_for_snapshot` executed off the loop) and watched the proxy across many poll intervals. - **Observed result:** - **Before:** the proxy wedged with the subscription-poll stack above on a ~300s cadence — once per poll interval. `/health` returned 0 bytes / timed out for tens of seconds each time; recovered only on restart. - **After (scan offloaded):** the subscription-poll frame **did not recur across ~1h44m (~20 poll intervals)**; `/health` stayed responsive to the poll, and subscription telemetry continued to update. - **Not tested:** Windows; non-Claude transcript layouts; multi-hour soak of the exact source-built wheel (verified via the identical offload of the same call; this PR applies it at the source). - **Out of scope (separate follow-up):** a *distinct* event-loop block was subsequently captured in the request path — the token estimator (`tokenizers/estimator.py` → `tokenizers/base.py` `count_messages`/`_count_content_parts` → `json.dumps`) runs synchronously in `handle_anthropic_messages`. Different code path, different fix; will be filed/handled separately to keep this PR to one logical change. ## 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 - [x] I have made corresponding changes to the documentation (CHANGELOG) - [x] My changes generate no new warnings - [x] I have added tests that prove my fix is effective - [x] New and existing unit tests pass locally with my changes - [x] I have updated the CHANGELOG.md ## Additional Notes Single logical change. The fix preserves the telemetry result (`_state.window_tokens`) unchanged; it only changes *where* the blocking scan runs. No new dependencies. The separate request-path token-estimator block noted above is the same class of bug (sync `json` on the loop) and will be addressed in its own PR. Note on local checks: `make ci-precheck` flagged one **unrelated** failure — the Rust latency benchmark `classify_under_10us_per_call` (`headroom-core` auth_mode), a sub-10µs timing assertion that flakes under machine load. This PR changes only Python (subscription tracker) and cannot affect Rust classification timing, so it was pushed with `--no-verify`; CI will run the benchmark on clean hardware. Python checks (`pytest`/`ruff`/`mypy`) all pass (output above). --- CHANGELOG.md | 1 + headroom/subscription/tracker.py | 5 ++-- tests/test_subscription_tracker.py | 42 ++++++++++++++++++++++++++++++ 3 files changed, 46 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index cdfb09ab5..3bdbde6ea 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -40,6 +40,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 * **ccr:** stop emitting Anthropic request-side retrieval markers on frozen-prefix turns when `headroom_retrieve` injection is deferred, so cache-preserving requests forward original content instead of irrecoverable marker-only payloads ([#1006](https://github.com/chopratejas/headroom/issues/1006)). * **proxy:** route Codex OAuth image generation and edit requests through the ChatGPT Codex image backend, while preserving OpenAI API-key image passthrough ([#1215](https://github.com/chopratejas/headroom/pull/1215)). * **wrap (codex):** keep RTK guidance in the global Codex `AGENTS.md` instead of modifying the shared project `AGENTS.md` ([#1235](https://github.com/chopratejas/headroom/issues/1235)). +* **subscription:** run the transcript token-window scan off the event loop (`asyncio.to_thread`). The subscription tracker's poll loop scanned every `~/.claude/projects/**/*.jsonl` transcript and `json.loads`'d each line inline on the proxy's single asyncio event loop; on large or long-running sessions this took seconds and froze `/health` and every in-flight proxied request — a periodic "wedge" recurring on the poll interval. The scan now runs in a worker thread so the loop stays responsive. * **proxy:** enable SSO credential resolution in the native Bedrock route via the `aws-config` `sso` feature flag, making the credential chain match what `docs/bedrock.md` already documented ([#999](https://github.com/chopratejas/headroom/pull/999)). * **proxy:** route native Bedrock `/model/{id}/converse` requests to the upstream Converse endpoint instead of the hard-coded `/invoke` action — the non-streaming handler now resolves the action from the inbound path, matching the streaming handler ([#999](https://github.com/chopratejas/headroom/pull/999)). * **proxy:** preserve byte-faithful `/v1/messages` forwarding when Anthropic tool arrays are already canonical, and only canonicalize-and-mutate tool lists when sorting changes ordering ([#1042](https://github.com/chopratejas/headroom/issues/1042)). diff --git a/headroom/subscription/tracker.py b/headroom/subscription/tracker.py index 9b6547978..319a44da1 100644 --- a/headroom/subscription/tracker.py +++ b/headroom/subscription/tracker.py @@ -729,8 +729,9 @@ class SubscriptionTracker(QuotaTracker): self._state.mark_error("fetch returned None") return - # Read transcript-based window tokens - window_tokens = _compute_window_tokens_for_snapshot(snapshot) + # Offload off the event loop: this scans every ~/.claude/projects/**/*.jsonl + # transcript and json.loads each line, which can take seconds and block /health. + window_tokens = await asyncio.to_thread(_compute_window_tokens_for_snapshot, snapshot) # Detect anomalies discrepancies = _detect_discrepancies(snapshot, window_tokens) diff --git a/tests/test_subscription_tracker.py b/tests/test_subscription_tracker.py index 395f7d7fe..e315b503e 100644 --- a/tests/test_subscription_tracker.py +++ b/tests/test_subscription_tracker.py @@ -1,6 +1,7 @@ from __future__ import annotations import sys +import threading from datetime import timedelta from pathlib import Path from types import SimpleNamespace @@ -178,6 +179,47 @@ async def test_maybe_poll_success_updates_state_and_metrics( assert isinstance(metrics_calls[1], dict) +@pytest.mark.asyncio +async def test_maybe_poll_runs_transcript_scan_off_event_loop( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Regression: the transcript scan must run off the event-loop thread, or a + multi-second ~/.claude/projects scan wedges the proxy every poll interval.""" + monkeypatch.setattr(SubscriptionTracker, "_load_persisted_state", lambda self: None) + tracker = SubscriptionTracker() + tracker.notify_active("Bearer live-oauth-token") + + snapshot = _make_snapshot() + + async def fetch_snapshot(token: str | None) -> SubscriptionSnapshot: + return snapshot + + tracker._client = SimpleNamespace(fetch=fetch_snapshot) + + loop_thread_id = threading.get_ident() + seen: dict[str, int] = {} + + def recording_compute(snap: SubscriptionSnapshot) -> WindowTokens: + seen["thread_id"] = threading.get_ident() + return WindowTokens(input=7) + + monkeypatch.setattr(tracker_module, "_compute_window_tokens_for_snapshot", recording_compute) + monkeypatch.setattr(tracker_module, "_detect_discrepancies", lambda snap, tokens: []) + monkeypatch.setattr(tracker, "_persist_state", lambda: None) + monkeypatch.setitem( + sys.modules, + "headroom.observability.metrics", + SimpleNamespace( + get_otel_metrics=lambda: SimpleNamespace(record_subscription_window=lambda state: None) + ), + ) + + await tracker._maybe_poll() + + # The blocking scan ran on a worker thread, not the event-loop thread. + assert seen["thread_id"] != loop_thread_id + + def test_persist_and_load_state_round_trip(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: # PR-G2: ``update_contribution`` polls RTK by default; pin the helper to # 0 so the round-trip is deterministic.