From d4cd49af9500eaa366c56107bdec2e30593785db Mon Sep 17 00:00:00 2001 From: Garm Date: Fri, 17 Apr 2026 17:12:38 +0200 Subject: [PATCH] feat(telemetry): add headroom_stack and install_mode identity fields Adds two orthogonal identity fields to the anonymous telemetry beacon so we can segment usage by integration surface and deployment shape: - headroom_stack: how Headroom is invoked (proxy, wrap_claude, wrap_codex, adapter_ts_openai, adapter_ts_anthropic, etc.). Resolved from HEADROOM_STACK env, HEADROOM_AGENT_TYPE fallback, or aggregated request-header counts. - install_mode: how the proxy is deployed (wrapped / persistent / on_demand). Detected from HEADROOM_AGENT_TYPE plus DeploymentManifest lookup. TS SDK adapters now tag every request with X-Headroom-Stack; a FastAPI middleware buckets the counts and surfaces them via /stats so the beacon can report requests_by_stack for mixed-integration sessions. Co-Authored-By: Claude Opus 4.7 (1M context) --- headroom/cli/wrap.py | 1 + headroom/proxy/prometheus_metrics.py | 17 ++++ headroom/proxy/server.py | 14 +++ headroom/telemetry/beacon.py | 13 +++ headroom/telemetry/context.py | 108 +++++++++++++++++++++++ sdk/typescript/src/adapters/anthropic.ts | 6 +- sdk/typescript/src/adapters/gemini.ts | 4 +- sdk/typescript/src/adapters/openai.ts | 6 +- sdk/typescript/src/adapters/vercel-ai.ts | 11 ++- sdk/typescript/src/client.ts | 8 ++ sdk/typescript/src/types.ts | 4 + sql/create_proxy_telemetry_v2.sql | 3 + sql/upgrade_telemetry_stack_context.sql | 15 ++++ tests/test_telemetry_context.py | 102 +++++++++++++++++++++ 14 files changed, 306 insertions(+), 6 deletions(-) create mode 100644 headroom/telemetry/context.py create mode 100644 sql/upgrade_telemetry_stack_context.sql create mode 100644 tests/test_telemetry_context.py diff --git a/headroom/cli/wrap.py b/headroom/cli/wrap.py index 310ed53ff..640ac6a0f 100644 --- a/headroom/cli/wrap.py +++ b/headroom/cli/wrap.py @@ -134,6 +134,7 @@ def _start_proxy( # Tell the proxy which agent is being wrapped (for traffic learning output) if agent_type != "unknown": proxy_env["HEADROOM_AGENT_TYPE"] = agent_type + proxy_env.setdefault("HEADROOM_STACK", f"wrap_{agent_type}") proc = subprocess.Popen( cmd, diff --git a/headroom/proxy/prometheus_metrics.py b/headroom/proxy/prometheus_metrics.py index 43ea17d09..6752f1f6b 100644 --- a/headroom/proxy/prometheus_metrics.py +++ b/headroom/proxy/prometheus_metrics.py @@ -70,6 +70,8 @@ class PrometheusMetrics: self.requests_total = 0 self.requests_by_provider: dict[str, int] = defaultdict(int) self.requests_by_model: dict[str, int] = defaultdict(int) + # Populated via X-Headroom-Stack header (TS SDK adapters, etc.) + self.requests_by_stack: dict[str, int] = defaultdict(int) self.requests_cached = 0 self.requests_rate_limited = 0 self.requests_failed = 0 @@ -192,6 +194,21 @@ class PrometheusMetrics: return total_input_tokens, total_input_cost_usd + def record_stack(self, stack: str | None) -> None: + """Increment the per-stack request counter. + + ``stack`` is the ``X-Headroom-Stack`` header value (e.g. + ``adapter_ts_openai``). Called once per inbound request from the + proxy's stack middleware; a no-op when the header is absent. + """ + + if not stack: + return + slug = stack.strip().lower() + if not slug or len(slug) > 64: + return + self.requests_by_stack[slug] += 1 + async def record_request( self, provider: str, diff --git a/headroom/proxy/server.py b/headroom/proxy/server.py index 7ee1f19a4..1190ba1bf 100644 --- a/headroom/proxy/server.py +++ b/headroom/proxy/server.py @@ -1188,6 +1188,19 @@ def create_app(config: ProxyConfig | None = None) -> FastAPI: allow_headers=["*"], ) + # X-Headroom-Stack: SDK adapters (TS openai/anthropic/etc.) tag their + # requests so telemetry can segment by integration surface. + @app.middleware("http") + async def _record_headroom_stack(request, call_next): + if request.url.path.startswith("/v1/"): + stack = request.headers.get("x-headroom-stack") + if stack: + try: + proxy.metrics.record_stack(stack) + except Exception: + logger.debug("record_stack failed", exc_info=True) + return await call_next(request) + # Health & Metrics @app.get("/livez") async def livez(): @@ -1356,6 +1369,7 @@ def create_app(config: ProxyConfig | None = None) -> FastAPI: "failed": m.requests_failed, "by_provider": dict(m.requests_by_provider), "by_model": dict(m.requests_by_model), + "by_stack": dict(m.requests_by_stack), }, "tokens": { "input": m.tokens_input_total, diff --git a/headroom/telemetry/beacon.py b/headroom/telemetry/beacon.py index e3a2483c1..1eb13a30d 100644 --- a/headroom/telemetry/beacon.py +++ b/headroom/telemetry/beacon.py @@ -19,6 +19,8 @@ import sys import time import uuid +from headroom.telemetry.context import detect_install_mode, detect_stack + logger = logging.getLogger(__name__) # Supabase endpoint for anonymous aggregate telemetry. @@ -89,6 +91,8 @@ class TelemetryBeacon: self._session_id = uuid.uuid4().hex # Stable across restarts — anonymous machine fingerprint (SHA256 of hostname) self._instance_id = hashlib.sha256(platform.node().encode()).hexdigest()[:16] + # Deployment shape is determined once at startup (wrapped / persistent / on_demand) + self._install_mode = detect_install_mode(port) async def start(self) -> None: """Start the periodic beacon. Call from proxy startup.""" @@ -177,8 +181,17 @@ class TelemetryBeacon: "sdk": self._sdk, "backend": self._backend, "session_minutes": session_minutes, + "install_mode": self._install_mode, + "headroom_stack": detect_stack(stats), } + try: + by_stack = (stats.get("requests") or {}).get("by_stack") or {} + if by_stack: + payload["requests_by_stack"] = dict(by_stack) + except Exception: + logger.debug("Beacon: failed to extract requests_by_stack", exc_info=True) + # --- Effectiveness metrics --- try: tokens = stats.get("tokens", {}) diff --git a/headroom/telemetry/context.py b/headroom/telemetry/context.py new file mode 100644 index 000000000..2b77c9533 --- /dev/null +++ b/headroom/telemetry/context.py @@ -0,0 +1,108 @@ +"""Deployment context detection for telemetry. + +Derives two orthogonal identity fields the beacon reports: + +* ``install_mode`` — how the proxy process is deployed + (``persistent`` / ``on_demand`` / ``wrapped`` / ``unknown``). +* ``headroom_stack`` — how Headroom is being invoked + (``proxy``, ``wrap_claude``, ``adapter_ts_openai``, ...). + +Both helpers are best-effort and never raise: telemetry is fire-and-forget and +must not break the proxy. +""" + +from __future__ import annotations + +import logging +import os +from typing import Any + +logger = logging.getLogger(__name__) + + +_KNOWN_WRAP_AGENTS = frozenset( + {"claude", "copilot", "codex", "aider", "cursor", "openclaw"} +) + + +def _slug_from_agent_type(agent_type: str) -> str: + """Return ``wrap_`` for known agents, otherwise ``unknown``.""" + + agent_type = agent_type.strip().lower() + if agent_type and agent_type in _KNOWN_WRAP_AGENTS: + return f"wrap_{agent_type}" + return "unknown" + + +def detect_install_mode(port: int) -> str: + """Classify how the proxy is deployed. + + Resolution order: + + 1. ``HEADROOM_AGENT_TYPE`` env var set → ``wrapped`` (spawned by ``headroom wrap``). + 2. A ``DeploymentManifest`` on disk whose port matches ``port`` → ``persistent``. + 3. Otherwise → ``on_demand``. + + Any failure falls back to ``unknown`` so a broken install subsystem + doesn't silence telemetry. + """ + + try: + if os.environ.get("HEADROOM_AGENT_TYPE"): + return "wrapped" + + try: + from headroom.install.state import list_manifests + + for manifest in list_manifests(): + if getattr(manifest, "port", None) == port: + return "persistent" + except Exception: + logger.debug( + "Beacon: manifest lookup failed during install_mode detection", + exc_info=True, + ) + + return "on_demand" + except Exception: + logger.debug("Beacon: detect_install_mode crashed", exc_info=True) + return "unknown" + + +def detect_stack(stats: dict[str, Any] | None = None) -> str: + """Classify how Headroom is being invoked. + + Resolution order: + + 1. ``HEADROOM_STACK`` env var set → use that slug verbatim. + 2. ``HEADROOM_AGENT_TYPE`` env var set → ``wrap_``. + 3. ``stats['requests']['by_stack']`` dict populated → + pick the stack with >80% of requests, else ``mixed``. + 4. Otherwise → ``proxy``. + + Any failure falls back to ``unknown``. + """ + + try: + explicit = os.environ.get("HEADROOM_STACK") + if explicit: + return explicit.strip().lower() + + agent_type = os.environ.get("HEADROOM_AGENT_TYPE") + if agent_type: + return _slug_from_agent_type(agent_type) + + if stats: + by_stack = (stats.get("requests") or {}).get("by_stack") or {} + if by_stack: + total = sum(by_stack.values()) + if total > 0: + dominant, count = max(by_stack.items(), key=lambda kv: kv[1]) + if count / total >= 0.8: + return dominant + return "mixed" + + return "proxy" + except Exception: + logger.debug("Beacon: detect_stack crashed", exc_info=True) + return "unknown" diff --git a/sdk/typescript/src/adapters/anthropic.ts b/sdk/typescript/src/adapters/anthropic.ts index 0b1772f57..16bceb40e 100644 --- a/sdk/typescript/src/adapters/anthropic.ts +++ b/sdk/typescript/src/adapters/anthropic.ts @@ -186,7 +186,11 @@ export function withHeadroom( options.model ?? params.model ?? "claude-sonnet-4-5-20250929"; const openaiMessages = anthropicToOpenAI(messages); - const result = await compress(openaiMessages, { ...options, model }); + const result = await compress(openaiMessages, { + stack: "adapter_ts_anthropic", + ...options, + model, + }); const anthropicMessages = result.compressed ? openAIToAnthropic(result.messages) diff --git a/sdk/typescript/src/adapters/gemini.ts b/sdk/typescript/src/adapters/gemini.ts index 7522e70db..40b652e20 100644 --- a/sdk/typescript/src/adapters/gemini.ts +++ b/sdk/typescript/src/adapters/gemini.ts @@ -44,7 +44,7 @@ export function withHeadroom( // compress() auto-detects Gemini format const result = await compress( Array.isArray(contents) ? contents : [contents], - { ...options, model: modelName }, + { stack: "adapter_ts_gemini", ...options, model: modelName }, ); const newParams = Array.isArray(params) @@ -61,7 +61,7 @@ export function withHeadroom( const result = await compress( Array.isArray(contents) ? contents : [contents], - { ...options, model: modelName }, + { stack: "adapter_ts_gemini", ...options, model: modelName }, ); const newParams = Array.isArray(params) diff --git a/sdk/typescript/src/adapters/openai.ts b/sdk/typescript/src/adapters/openai.ts index f3bc2b289..50b6bb2c6 100644 --- a/sdk/typescript/src/adapters/openai.ts +++ b/sdk/typescript/src/adapters/openai.ts @@ -42,7 +42,11 @@ export function withHeadroom( const messages: OpenAIMessage[] = params.messages; const model = options.model ?? params.model ?? "gpt-4o"; - const result = await compress(messages, { ...options, model }); + const result = await compress(messages, { + stack: "adapter_ts_openai", + ...options, + model, + }); return originalCreate({ ...params, diff --git a/sdk/typescript/src/adapters/vercel-ai.ts b/sdk/typescript/src/adapters/vercel-ai.ts index 4a8806157..21ab04b9f 100644 --- a/sdk/typescript/src/adapters/vercel-ai.ts +++ b/sdk/typescript/src/adapters/vercel-ai.ts @@ -54,7 +54,11 @@ export function headroomMiddleware(options: CompressOptions = {}) { const openaiMessages = vercelToOpenAI(prompt); // Compress via Headroom - const result = await compress(openaiMessages, { ...options, model }); + const result = await compress(openaiMessages, { + stack: "adapter_ts_vercel_ai", + ...options, + model, + }); if (!result.compressed) return params; @@ -75,7 +79,10 @@ export async function compressVercelMessages( options: CompressOptions = {}, ): Promise { const openaiMessages = vercelToOpenAI(messages); - const result = await compress(openaiMessages, options); + const result = await compress(openaiMessages, { + stack: "adapter_ts_vercel_ai", + ...options, + }); const vercelMessages = openAIToVercel(result.messages); return { diff --git a/sdk/typescript/src/client.ts b/sdk/typescript/src/client.ts index 3402386fd..c3f5dc528 100644 --- a/sdk/typescript/src/client.ts +++ b/sdk/typescript/src/client.ts @@ -201,6 +201,7 @@ export class HeadroomClient implements HeadroomClientInterface { private fallback: boolean; private retries: number; private config: HeadroomConfig | undefined; + private stack: string | undefined; /** @internal */ providerApiKey: string | undefined; @@ -221,6 +222,7 @@ export class HeadroomClient implements HeadroomClientInterface { this.retries = options.retries ?? DEFAULT_RETRIES; this.providerApiKey = options.providerApiKey; this.config = options.config; + this.stack = options.stack; this.chat = { completions: new ChatCompletions(this) }; this.messages = new Messages(this); @@ -513,6 +515,9 @@ export class HeadroomClient implements HeadroomClientInterface { headers["Authorization"] = `Bearer ${this.apiKey}`; } } + if (this.stack && !headers["X-Headroom-Stack"]) { + headers["X-Headroom-Stack"] = this.stack; + } let response: Response; try { @@ -558,6 +563,9 @@ export class HeadroomClient implements HeadroomClientInterface { if (this.apiKey) { headers["Authorization"] = `Bearer ${this.apiKey}`; } + if (this.stack && !headers["X-Headroom-Stack"]) { + headers["X-Headroom-Stack"] = this.stack; + } let response: Response; try { diff --git a/sdk/typescript/src/types.ts b/sdk/typescript/src/types.ts index b9c052671..04738fd6d 100644 --- a/sdk/typescript/src/types.ts +++ b/sdk/typescript/src/types.ts @@ -67,6 +67,8 @@ export interface CompressOptions { tokenBudget?: number; /** Compression hooks for pre/post processing. */ hooks?: CompressionHooks; + /** Integration slug sent as X-Headroom-Stack (e.g. "adapter_ts_openai"). */ + stack?: string; } export interface CompressResult { @@ -89,6 +91,8 @@ export interface HeadroomClientOptions { timeout?: number; fallback?: boolean; retries?: number; + /** Integration slug sent as X-Headroom-Stack on every request. */ + stack?: string; } export interface HeadroomClientInterface { diff --git a/sql/create_proxy_telemetry_v2.sql b/sql/create_proxy_telemetry_v2.sql index 18efdcbbb..f76e2310d 100644 --- a/sql/create_proxy_telemetry_v2.sql +++ b/sql/create_proxy_telemetry_v2.sql @@ -15,6 +15,9 @@ CREATE TABLE IF NOT EXISTS proxy_telemetry_v2 ( sdk text, backend text, session_minutes integer, + headroom_stack text, + install_mode text, + requests_by_stack jsonb, -- Effectiveness metrics tokens_saved bigint, diff --git a/sql/upgrade_telemetry_stack_context.sql b/sql/upgrade_telemetry_stack_context.sql new file mode 100644 index 000000000..396cb70e9 --- /dev/null +++ b/sql/upgrade_telemetry_stack_context.sql @@ -0,0 +1,15 @@ +-- Add deployment-context columns to proxy_telemetry_v2 +-- Run in Supabase SQL Editor. +-- +-- headroom_stack: how Headroom is being invoked — e.g. "proxy", +-- "wrap_claude", "adapter_ts_openai", "mixed", "unknown". +-- install_mode: how the proxy process is deployed — one of +-- "wrapped", "persistent", "on_demand", "unknown". +-- requests_by_stack: JSONB dict {stack_slug: count} for sessions that see +-- multiple integration surfaces (e.g. a persistent proxy +-- serving both wrap_claude and TS adapter callers). + +ALTER TABLE proxy_telemetry_v2 + ADD COLUMN IF NOT EXISTS headroom_stack text, + ADD COLUMN IF NOT EXISTS install_mode text, + ADD COLUMN IF NOT EXISTS requests_by_stack jsonb; diff --git a/tests/test_telemetry_context.py b/tests/test_telemetry_context.py new file mode 100644 index 000000000..df0cf6015 --- /dev/null +++ b/tests/test_telemetry_context.py @@ -0,0 +1,102 @@ +"""Tests for headroom.telemetry.context (install_mode + headroom_stack detection).""" + +from __future__ import annotations + +from types import SimpleNamespace + +import pytest + +from headroom.telemetry.context import detect_install_mode, detect_stack + + +@pytest.fixture(autouse=True) +def _clean_env(monkeypatch): + """Every test starts without our env vars set.""" + + monkeypatch.delenv("HEADROOM_STACK", raising=False) + monkeypatch.delenv("HEADROOM_AGENT_TYPE", raising=False) + yield + + +class TestDetectInstallMode: + def test_wrapped_when_agent_type_set(self, monkeypatch): + monkeypatch.setenv("HEADROOM_AGENT_TYPE", "claude") + assert detect_install_mode(8787) == "wrapped" + + def test_on_demand_when_no_env_and_no_manifest(self, monkeypatch): + monkeypatch.setattr( + "headroom.install.state.list_manifests", lambda: [] + ) + assert detect_install_mode(8787) == "on_demand" + + def test_persistent_when_manifest_matches_port(self, monkeypatch): + manifest = SimpleNamespace(port=8787, profile="default") + monkeypatch.setattr( + "headroom.install.state.list_manifests", lambda: [manifest] + ) + assert detect_install_mode(8787) == "persistent" + + def test_on_demand_when_manifest_port_mismatches(self, monkeypatch): + manifest = SimpleNamespace(port=9000, profile="other") + monkeypatch.setattr( + "headroom.install.state.list_manifests", lambda: [manifest] + ) + assert detect_install_mode(8787) == "on_demand" + + def test_wrapped_takes_precedence_over_manifest(self, monkeypatch): + monkeypatch.setenv("HEADROOM_AGENT_TYPE", "codex") + manifest = SimpleNamespace(port=8787, profile="default") + monkeypatch.setattr( + "headroom.install.state.list_manifests", lambda: [manifest] + ) + assert detect_install_mode(8787) == "wrapped" + + def test_manifest_crash_falls_back_to_on_demand(self, monkeypatch): + def _boom(): + raise RuntimeError("disk gone") + + monkeypatch.setattr("headroom.install.state.list_manifests", _boom) + # install_mode should not raise; graceful fallback + assert detect_install_mode(8787) == "on_demand" + + +class TestDetectStack: + def test_explicit_env_wins(self, monkeypatch): + monkeypatch.setenv("HEADROOM_STACK", "custom_slug") + assert detect_stack() == "custom_slug" + + def test_explicit_env_overrides_agent_type(self, monkeypatch): + monkeypatch.setenv("HEADROOM_STACK", "proxy") + monkeypatch.setenv("HEADROOM_AGENT_TYPE", "claude") + assert detect_stack() == "proxy" + + def test_wrap_slug_from_agent_type(self, monkeypatch): + monkeypatch.setenv("HEADROOM_AGENT_TYPE", "claude") + assert detect_stack() == "wrap_claude" + + def test_unknown_agent_type_rejected(self, monkeypatch): + monkeypatch.setenv("HEADROOM_AGENT_TYPE", "somebespoke") + assert detect_stack() == "unknown" + + def test_default_is_proxy(self): + assert detect_stack() == "proxy" + + def test_default_is_proxy_with_empty_stats(self): + assert detect_stack({"requests": {"by_stack": {}}}) == "proxy" + + def test_dominant_stack_from_stats(self): + stats = {"requests": {"by_stack": {"adapter_ts_openai": 90, "adapter_ts_anthropic": 10}}} + assert detect_stack(stats) == "adapter_ts_openai" + + def test_mixed_when_no_dominant_stack(self): + stats = {"requests": {"by_stack": {"adapter_ts_openai": 40, "adapter_ts_anthropic": 60}}} + assert detect_stack(stats) == "mixed" + + def test_single_stack_is_dominant(self): + stats = {"requests": {"by_stack": {"adapter_ts_openai": 3}}} + assert detect_stack(stats) == "adapter_ts_openai" + + def test_env_beats_stats(self, monkeypatch): + monkeypatch.setenv("HEADROOM_STACK", "wrap_claude") + stats = {"requests": {"by_stack": {"adapter_ts_openai": 100}}} + assert detect_stack(stats) == "wrap_claude"