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) <noreply@anthropic.com>
This commit is contained in:
Garm 2026-04-17 17:12:38 +02:00
parent 4651f96c92
commit d4cd49af95
14 changed files with 306 additions and 6 deletions

View file

@ -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,

View file

@ -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,

View file

@ -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,

View file

@ -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", {})

View file

@ -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_<agent>`` 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_<agent>``.
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"

View file

@ -186,7 +186,11 @@ export function withHeadroom<T extends AnthropicLike>(
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)

View file

@ -44,7 +44,7 @@ export function withHeadroom<T extends GeminiModelLike>(
// 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<T extends GeminiModelLike>(
const result = await compress(
Array.isArray(contents) ? contents : [contents],
{ ...options, model: modelName },
{ stack: "adapter_ts_gemini", ...options, model: modelName },
);
const newParams = Array.isArray(params)

View file

@ -42,7 +42,11 @@ export function withHeadroom<T extends OpenAILike>(
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,

View file

@ -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<CompressResult & { messages: VercelMessage[] }> {
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 {

View file

@ -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 {

View file

@ -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 {

View file

@ -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,

View file

@ -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;

View file

@ -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"