mirror of
https://github.com/headroomlabs-ai/headroom.git
synced 2026-08-27 14:17:10 -04:00
fix(telemetry): anonymous compression stats — no prompts, no data (#2728)
## In one line
Headroom starts reporting **how well compression is working** — counters
and percentages only. **No prompts. No code. No file paths. Nothing
about what you're building.**
## Why
Right now nobody knows whether compression actually helps real users.
You can see your own numbers in `/stats`, but that's it — there's no way
to tell whether a given workload compresses well, or why it sometimes
doesn't. This closes that loop so we can make compression better for
everyone.
## Exactly what gets sent
One message per session, and every 5 minutes while you're active:
```json
{
"session": { "id": "random", "turns": 47, "duration_s": 4210, "seq": 3 },
"tokens": { "original": 890000, "attempted": 410000, "saved": 320000,
"tool_saved": 48000, "cache_read": 210000 },
"rates": { "saved_pct": 35.96, "eligible_pct": 46.07, "yield_pct": 78.05,
"cache_read_pct": 23.60, "overhead_pct": 1.96 },
"compression": { "transforms": {"crush": 47}, "passthrough_turns": 0 },
"skips": {},
"sources": { "proxy": 47 },
"providers": ["anthropic"],
"models": ["claude-sonnet-4-5-20250929"],
"failures": 2
}
```
Plus a random install ID, the Headroom version, and OS/architecture
(`darwin`, `arm64`).
That's the whole thing. A full example lives at
`deploy/beacon/sample-event.json`.
## What is never sent
- Your prompts or the model's responses
- Your code
- File paths, project names, repo names
- Tool names or MCP server names
- Hostname, username, or IP address
- Custom or fine-tuned model names (an id like `ft:gpt-4o:acme-corp:…`
contains a company name, so only models in a public registry are
reported)
**This is structural, not a pinky-swear.** Every value in the payload is
a number, a fixed word, or a random ID — there is no free-text field
anywhere for content to hide in. The receiver
(`deploy/beacon/worker.js`, in this repo so you can read it) drops
anything not on an explicit allowlist before storing.
## Turning it off
Any one of these:
```bash
HEADROOM_BEACON=off # or
DO_NOT_TRACK=1 # or
# offline mode
```
It's on by default, and Headroom says so at startup:
```
Telemetry: anonymous compression stats — never prompts, code, or file paths.
Helps us improve compression | Off: HEADROOM_BEACON=off
```
`HEADROOM_TELEMETRY` is a **separate** switch that still only affects
local stats. If you had turned that on, this change does not start
uploading anything — you answered a different question, and upgrading
should not change the answer.
## Why the percentages, not just "tokens saved"
"We saved 36%" hides the interesting part. In the example above only
**46% of tokens were eligible** for compression at all — the rest is
frozen cache prefix and system prompts we deliberately do not touch. Of
what we *could* touch, we removed **78%**.
Those are two separate problems. Raising eligibility is proxy work;
raising yield is compressor work. A single number cannot tell us which
to fix.
## Coverage
`emit_request_outcome` is a single chokepoint —
`handler.metrics.record_request` is called from exactly one place,
inside the funnel — so all 30 `RequestOutcome` construction sites are
covered: Anthropic, OpenAI, Gemini, Bedrock, batch, streaming, and the
long-lived Codex Responses-WS path.
The `headroom_compress` MCP path bypassed that funnel and is now wired
in separately. It has a different shape (no provider, no upstream
latency, and everything handed to the tool is eligible by construction),
so `sources` counts turns by origin — MCP turns always read
`eligible_pct: 100` and must not drag the proxy's real eligibility
ceiling upward.
**Subagents.** All subagent traffic through the proxy merges into one
session, which is correct for savings and retention but means `turns`
conflates fan-out with depth. Fan-out is still derivable —
`compression.latency_ms_total / session.duration_s` gives the
concurrency ratio (~1x serial, ~4x for four parallel agents), so no
extra field is needed. Verified no lost updates under 6-way concurrency
(1,200 turns).
**Known gap:** `--workers N` gives each process its own aggregator, so
one user session becomes up to N. Token totals and fleet rates stay
correct; session counts inflate. This matches the existing documented
limitation that TOIN state, CostTracker, and the prefix tracker are all
per-process.
## Notes for reviewers
- **Cumulative snapshots, not deltas.** Every report restates running
totals under one session ID, so the highest `seq` per `(install,
session)` is the complete session. Dedupe is a window function, and a
lost report costs nothing.
- **Never breaks the proxy.** Every path swallows its own exceptions;
uploads go out on a daemon thread so nothing blocks the request loop.
- **Explicit User-Agent is load-bearing.** urllib's default is blocked
by Cloudflare (error 1010). Combined with fire-and-forget error
handling, that would have failed every upload while looking perfectly
healthy.
- **The exit flush was broken and is fixed.** `atexit` handed the POST
to a daemon thread, and daemon threads are killed before they finish
during interpreter shutdown — so nothing was sent. That silently dropped
*every session shorter than the 5-minute heartbeat*, plus all
short-lived subagent MCP processes. The exit path now posts
synchronously with a 2s timeout.
- Receiver and query tooling are in `deploy/beacon/`.
## Testing
- `python -m headroom.telemetry.session` self-check: dedupe, cumulative
totals, dropped-report recovery, payload contains no model id or
prompt-derived string, allowlist coverage
- 175 telemetry/outcome tests pass; 6 new ones cover the opt-out notice
- Verified end to end against a live deployment: client → receiver →
storage → query
## Still to do before release
The default endpoint currently points at a temporary `workers.dev` URL.
It needs to move to a Headroom-owned hostname before this ships in a
tagged release — noted inline at `DEFAULT_ENDPOINT`.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
007446c73a
commit
9cfb00838a
11 changed files with 1855 additions and 13 deletions
2
deploy/beacon/.gitignore
vendored
Normal file
2
deploy/beacon/.gitignore
vendored
Normal file
|
|
@ -0,0 +1,2 @@
|
|||
# wrangler local state, caches, and account info — never commit
|
||||
.wrangler/
|
||||
94
deploy/beacon/query.sh
Executable file
94
deploy/beacon/query.sh
Executable file
|
|
@ -0,0 +1,94 @@
|
|||
#!/usr/bin/env bash
|
||||
# Query the telemetry corpus in R2 with DuckDB.
|
||||
#
|
||||
# ./query.sh # fleet summary
|
||||
# ./query.sh sessions # one row per session (deduped)
|
||||
# ./query.sh "SELECT ..." # your own SQL against the corpus
|
||||
#
|
||||
# Setup, once:
|
||||
# brew install duckdb
|
||||
# Cloudflare > R2 > API > Create Account API Token (Object Read only,
|
||||
# scoped to headroom-telemetry), then put the values in ~/env.txt
|
||||
# (or any file named by HEADROOM_ENV_FILE):
|
||||
#
|
||||
# R2_ACCOUNT_ID=...
|
||||
# R2_ACCESS_KEY_ID=...
|
||||
# R2_SECRET_ACCESS_KEY=...
|
||||
#
|
||||
# R2_ACCOUNT_TOKEN is Cloudflare's REST-API token and is NOT used here — the
|
||||
# S3 protocol wants the access-key pair.
|
||||
set -euo pipefail
|
||||
|
||||
BUCKET="${R2_BUCKET:-headroom-telemetry}"
|
||||
_repo_env="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)/.env"
|
||||
ENV_FILE="${HEADROOM_ENV_FILE:-$HOME/env.txt}"
|
||||
[ -f "$ENV_FILE" ] || ENV_FILE="$_repo_env"
|
||||
|
||||
[ -f "$ENV_FILE" ] || { echo "no env file (~/env.txt or $_repo_env) — see this script's header" >&2; exit 1; }
|
||||
# shellcheck disable=SC1090
|
||||
set -a; source "$ENV_FILE"; set +a
|
||||
|
||||
for v in R2_ACCOUNT_ID R2_ACCESS_KEY_ID R2_SECRET_ACCESS_KEY; do
|
||||
[ -n "${!v:-}" ] || { echo "$v not set in $ENV_FILE" >&2; exit 1; }
|
||||
done
|
||||
command -v duckdb >/dev/null || { echo "duckdb not installed: brew install duckdb" >&2; exit 1; }
|
||||
|
||||
# Credentials go in via a heredoc on stdin, never on the command line, so they
|
||||
# stay out of `ps` and shell history.
|
||||
SECRET="
|
||||
INSTALL httpfs; LOAD httpfs;
|
||||
CREATE OR REPLACE SECRET r2corpus (
|
||||
TYPE r2,
|
||||
KEY_ID '${R2_ACCESS_KEY_ID}',
|
||||
SECRET '${R2_SECRET_ACCESS_KEY}',
|
||||
ACCOUNT_ID '${R2_ACCOUNT_ID}'
|
||||
);
|
||||
"
|
||||
|
||||
# The corpus is heartbeats: a session reports every 5 minutes with CUMULATIVE
|
||||
# totals under one id. So the row with the highest seq per (install, session) is
|
||||
# the whole session — never SUM across heartbeats, you would count each session
|
||||
# once per report.
|
||||
DEDUPE="
|
||||
CREATE OR REPLACE TEMP VIEW sessions AS
|
||||
SELECT * FROM read_ndjson('r2://${BUCKET}/sessions/**/*.json', union_by_name = true)
|
||||
QUALIFY row_number() OVER (
|
||||
PARTITION BY resource['headroom.install_id'], session.id
|
||||
ORDER BY session.seq DESC
|
||||
) = 1;
|
||||
"
|
||||
|
||||
case "${1:-summary}" in
|
||||
summary)
|
||||
# Fleet rates come from summing raw counts. Averaging the per-session
|
||||
# rates.*_pct fields would weight a 10-token session equal to a 1M one.
|
||||
QUERY="
|
||||
SELECT count(*) AS sessions,
|
||||
count(DISTINCT resource['headroom.install_id']) AS installs,
|
||||
sum(session.turns) AS turns,
|
||||
sum(tokens.saved) AS tokens_saved,
|
||||
sum(tokens.tool_saved) AS tool_tokens_saved,
|
||||
round(sum(tokens.attempted) * 100.0
|
||||
/ nullif(sum(tokens.original), 0), 2) AS eligible_pct,
|
||||
round(sum(tokens.saved) * 100.0
|
||||
/ nullif(sum(tokens.attempted), 0), 2) AS yield_pct,
|
||||
round(sum(tokens.saved) * 100.0
|
||||
/ nullif(sum(tokens.original), 0), 2) AS saved_pct,
|
||||
sum(failures) AS failures
|
||||
FROM sessions;"
|
||||
;;
|
||||
sessions)
|
||||
QUERY="
|
||||
SELECT resource['headroom.install_id'][1:8] AS install,
|
||||
session.id, session.seq, session.turns, session.duration_s,
|
||||
tokens.original, tokens.attempted, tokens.saved,
|
||||
rates.saved_pct, rates.eligible_pct, rates.yield_pct,
|
||||
providers, models, skips
|
||||
FROM sessions
|
||||
ORDER BY session.duration_s DESC
|
||||
LIMIT 50;"
|
||||
;;
|
||||
*) QUERY="$1" ;;
|
||||
esac
|
||||
|
||||
printf '%s\n%s\n%s\n' "$SECRET" "$DEDUPE" "$QUERY" | duckdb -box
|
||||
297
deploy/beacon/sample-event.json
Normal file
297
deploy/beacon/sample-event.json
Normal file
|
|
@ -0,0 +1,297 @@
|
|||
{
|
||||
"resourceLogs": [
|
||||
{
|
||||
"resource": {
|
||||
"attributes": [
|
||||
{
|
||||
"key": "service.name",
|
||||
"value": {
|
||||
"stringValue": "headroom"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "service.version",
|
||||
"value": {
|
||||
"stringValue": "0.34.0"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "headroom.install_id",
|
||||
"value": {
|
||||
"stringValue": "00000000000000000000000000000000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "os.type",
|
||||
"value": {
|
||||
"stringValue": "darwin"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "host.arch",
|
||||
"value": {
|
||||
"stringValue": "arm64"
|
||||
}
|
||||
}
|
||||
]
|
||||
},
|
||||
"scopeLogs": [
|
||||
{
|
||||
"scope": {
|
||||
"name": "headroom.telemetry.session"
|
||||
},
|
||||
"logRecords": [
|
||||
{
|
||||
"timeUnixNano": "1785731364402434048",
|
||||
"body": {
|
||||
"kvlistValue": {
|
||||
"values": [
|
||||
{
|
||||
"key": "schema_version",
|
||||
"value": {
|
||||
"intValue": "1"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "session",
|
||||
"value": {
|
||||
"kvlistValue": {
|
||||
"values": [
|
||||
{
|
||||
"key": "id",
|
||||
"value": {
|
||||
"stringValue": "sample0000000001"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "seq",
|
||||
"value": {
|
||||
"intValue": "0"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "duration_s",
|
||||
"value": {
|
||||
"intValue": "4210"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "turns",
|
||||
"value": {
|
||||
"intValue": "47"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "ended",
|
||||
"value": {
|
||||
"stringValue": "active"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "final",
|
||||
"value": {
|
||||
"boolValue": false
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "tokens",
|
||||
"value": {
|
||||
"kvlistValue": {
|
||||
"values": [
|
||||
{
|
||||
"key": "original",
|
||||
"value": {
|
||||
"intValue": "890000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "attempted",
|
||||
"value": {
|
||||
"intValue": "410000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "input",
|
||||
"value": {
|
||||
"intValue": "570000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "output",
|
||||
"value": {
|
||||
"intValue": "41000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "saved",
|
||||
"value": {
|
||||
"intValue": "320000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "tool_saved",
|
||||
"value": {
|
||||
"intValue": "48000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "cache_read",
|
||||
"value": {
|
||||
"intValue": "210000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "cache_write",
|
||||
"value": {
|
||||
"intValue": "30000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "uncached",
|
||||
"value": {
|
||||
"intValue": "650000"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "rates",
|
||||
"value": {
|
||||
"kvlistValue": {
|
||||
"values": [
|
||||
{
|
||||
"key": "saved_pct",
|
||||
"value": {
|
||||
"doubleValue": 35.96
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "eligible_pct",
|
||||
"value": {
|
||||
"doubleValue": 46.07
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "yield_pct",
|
||||
"value": {
|
||||
"doubleValue": 78.05
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "cache_read_pct",
|
||||
"value": {
|
||||
"doubleValue": 23.6
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "overhead_pct",
|
||||
"value": {
|
||||
"doubleValue": 1.96
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "compression",
|
||||
"value": {
|
||||
"kvlistValue": {
|
||||
"values": [
|
||||
{
|
||||
"key": "transforms",
|
||||
"value": {
|
||||
"kvlistValue": {
|
||||
"values": [
|
||||
{
|
||||
"key": "crush",
|
||||
"value": {
|
||||
"intValue": "47"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "overhead_ms_total",
|
||||
"value": {
|
||||
"intValue": "1840"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "latency_ms_total",
|
||||
"value": {
|
||||
"intValue": "94000"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "passthrough_turns",
|
||||
"value": {
|
||||
"intValue": "0"
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "response_cache_hits",
|
||||
"value": {
|
||||
"intValue": "3"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "skips",
|
||||
"value": {
|
||||
"kvlistValue": {
|
||||
"values": []
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "providers",
|
||||
"value": {
|
||||
"arrayValue": {
|
||||
"values": [
|
||||
{
|
||||
"stringValue": "anthropic"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "models",
|
||||
"value": {
|
||||
"arrayValue": {
|
||||
"values": [
|
||||
{
|
||||
"stringValue": "claude-sonnet-4-5-20250929"
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"key": "failures",
|
||||
"value": {
|
||||
"intValue": "2"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
]
|
||||
}
|
||||
174
deploy/beacon/worker.js
Normal file
174
deploy/beacon/worker.js
Normal file
|
|
@ -0,0 +1,174 @@
|
|||
/**
|
||||
* Headroom telemetry beacon receiver.
|
||||
*
|
||||
* This file is open source on purpose. It is the other half of the promise
|
||||
* made in headroom/telemetry/session.py: users can read exactly what the
|
||||
* client sends AND exactly what happens to it on arrival. "Trust us" is not a
|
||||
* privacy policy.
|
||||
*
|
||||
* Deployed at otlp.headroomlabs.ai. Three jobs:
|
||||
*
|
||||
* 1. Allowlist. Drop every field not on ALLOWED_KEYS before anything is
|
||||
* written. This is the only privacy control that works retroactively —
|
||||
* if a future client version ships a bug that leaks a field, we cannot
|
||||
* patch the installs already in the wild, but we can stop storing it
|
||||
* here in one deploy.
|
||||
*
|
||||
* 2. Flatten. OTLP AnyValue nesting is portable but miserable to query
|
||||
* ({"kvlistValue":{"values":[{"key":"tokens",...}]}}). We keep OTLP on
|
||||
* the wire so the backend stays vendor-swappable, and store plain JSON so
|
||||
* DuckDB can read it without unwrapping anything.
|
||||
*
|
||||
* 3. Fan out. R2 for the durable corpus; optionally a metrics vendor for
|
||||
* dashboards. Adding a destination is one more call here — never a
|
||||
* client release.
|
||||
*
|
||||
* What this deliberately does NOT do: log, store, or forward the source IP.
|
||||
* Cloudflare offers it as cf-connecting-ip; it is the one field that would
|
||||
* deanonymise install_id, so it is never read.
|
||||
*/
|
||||
|
||||
// Mirrors the payload built by _Session.payload(). A key absent here is
|
||||
// dropped, not stored. Adding a metric means adding it here first — that
|
||||
// friction is the point.
|
||||
const ALLOWED_KEYS = [
|
||||
'schema_version',
|
||||
'session',
|
||||
'tokens',
|
||||
'rates',
|
||||
'compression',
|
||||
'skips',
|
||||
'sources',
|
||||
'providers',
|
||||
'models',
|
||||
'failures',
|
||||
];
|
||||
|
||||
// Resource attributes we keep. Same rule: allowlist, not denylist.
|
||||
const ALLOWED_RESOURCE = [
|
||||
'service.name',
|
||||
'service.version',
|
||||
'headroom.install_id',
|
||||
'headroom.install_mode',
|
||||
'headroom.stack',
|
||||
'os.type',
|
||||
'host.arch',
|
||||
];
|
||||
|
||||
// A beacon event is ~2KB. Anything far past that is a bug or an attack.
|
||||
const MAX_BODY_BYTES = 64 * 1024;
|
||||
|
||||
/** OTLP AnyValue -> plain JS. The inverse of _any_value() in session.py. */
|
||||
function unwrap(value) {
|
||||
if (value == null) return null;
|
||||
if ('stringValue' in value) return value.stringValue;
|
||||
if ('boolValue' in value) return value.boolValue;
|
||||
if ('intValue' in value) return Number(value.intValue);
|
||||
if ('doubleValue' in value) return value.doubleValue;
|
||||
if ('arrayValue' in value) return (value.arrayValue.values || []).map(unwrap);
|
||||
if ('kvlistValue' in value) {
|
||||
const out = {};
|
||||
for (const kv of value.kvlistValue.values || []) out[kv.key] = unwrap(kv.value);
|
||||
return out;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function pick(obj, allowed) {
|
||||
const out = {};
|
||||
if (!obj || typeof obj !== 'object') return out;
|
||||
for (const key of allowed) {
|
||||
if (key in obj) out[key] = obj[key];
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/** OTLP ExportLogsServiceRequest -> flat, allowlisted records. */
|
||||
function extract(payload) {
|
||||
const records = [];
|
||||
for (const rl of payload.resourceLogs || []) {
|
||||
const resource = {};
|
||||
for (const attr of rl.resource?.attributes || []) {
|
||||
resource[attr.key] = unwrap(attr.value);
|
||||
}
|
||||
const cleanResource = pick(resource, ALLOWED_RESOURCE);
|
||||
|
||||
for (const sl of rl.scopeLogs || []) {
|
||||
for (const rec of sl.logRecords || []) {
|
||||
const body = unwrap(rec.body);
|
||||
if (!body || typeof body !== 'object') continue;
|
||||
records.push({
|
||||
...pick(body, ALLOWED_KEYS),
|
||||
resource: cleanResource,
|
||||
// Server-stamped. A client clock can be wrong or forged; this is the
|
||||
// timestamp partitioning and retention actually rely on.
|
||||
received_at: new Date().toISOString(),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
return records;
|
||||
}
|
||||
|
||||
export default {
|
||||
async fetch(request, env, ctx) {
|
||||
if (request.method !== 'POST') {
|
||||
return new Response('beacon: POST OTLP logs to /v1/logs', { status: 405 });
|
||||
}
|
||||
const url = new URL(request.url);
|
||||
if (url.pathname !== '/v1/logs') {
|
||||
return new Response('not found', { status: 404 });
|
||||
}
|
||||
|
||||
const raw = await request.arrayBuffer();
|
||||
if (raw.byteLength > MAX_BODY_BYTES) {
|
||||
return new Response('payload too large', { status: 413 });
|
||||
}
|
||||
|
||||
let records;
|
||||
try {
|
||||
records = extract(JSON.parse(new TextDecoder().decode(raw)));
|
||||
} catch {
|
||||
// Malformed input is not worth a retry storm from clients.
|
||||
return new Response('bad request', { status: 400 });
|
||||
}
|
||||
if (records.length === 0) return new Response(null, { status: 204 });
|
||||
|
||||
const now = new Date();
|
||||
const day = now.toISOString().slice(0, 10);
|
||||
const hour = now.toISOString().slice(11, 13);
|
||||
// Hive-style partitioning so DuckDB can prune by date without a catalog.
|
||||
// ponytail: one object per request. At beacon volume that is a few hundred
|
||||
// thousand objects a month, which globs fine. Add a daily compaction job
|
||||
// when the file count starts to slow queries, not before.
|
||||
const key = `sessions/dt=${day}/hh=${hour}/${crypto.randomUUID()}.json`;
|
||||
const ndjson = records.map((r) => JSON.stringify(r)).join('\n');
|
||||
|
||||
// Respond immediately; durability work continues after the response.
|
||||
// The client is fire-and-forget and ignores the status anyway — making it
|
||||
// wait on R2 would only add latency to someone else's coding session.
|
||||
ctx.waitUntil(
|
||||
env.CORPUS.put(key, ndjson, {
|
||||
httpMetadata: { contentType: 'application/x-ndjson' },
|
||||
})
|
||||
);
|
||||
|
||||
// Optional second lane: forward verbatim OTLP to a metrics backend for
|
||||
// dashboards. Configured by secret, so it can be added or swapped with a
|
||||
// `wrangler secret put` and no code change.
|
||||
if (env.METRICS_OTLP_URL) {
|
||||
ctx.waitUntil(
|
||||
fetch(env.METRICS_OTLP_URL, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'content-type': 'application/json',
|
||||
authorization: env.METRICS_OTLP_AUTH || '',
|
||||
},
|
||||
body: JSON.stringify({ resourceLogs: [{ scopeLogs: [{ logRecords: records.map((r) => ({ body: { stringValue: JSON.stringify(r) } })) }] }] }),
|
||||
}).catch(() => {})
|
||||
);
|
||||
}
|
||||
|
||||
return new Response(null, { status: 204 });
|
||||
},
|
||||
};
|
||||
44
deploy/beacon/wrangler.toml
Normal file
44
deploy/beacon/wrangler.toml
Normal file
|
|
@ -0,0 +1,44 @@
|
|||
name = "headroom-beacon"
|
||||
main = "worker.js"
|
||||
compatibility_date = "2025-01-01"
|
||||
|
||||
# PHASE 1 — deploy to <name>.<subdomain>.workers.dev with no DNS changes.
|
||||
# Lets the whole path be tested against a real client before headroomlabs.ai
|
||||
# nameservers move anywhere.
|
||||
workers_dev = true
|
||||
|
||||
# PHASE 2 — the permanent address. Uncomment once headroomlabs.ai is on
|
||||
# Cloudflare nameservers, then redeploy. This string is baked into every
|
||||
# released client (DEFAULT_ENDPOINT in headroom/telemetry/session.py), so it can
|
||||
# never change afterwards — everything behind it can.
|
||||
#
|
||||
# Deploying this while the zone is still on Namecheap fails: wrangler cannot
|
||||
# find the zone. That is the intended guardrail, not a bug.
|
||||
#
|
||||
# [[routes]]
|
||||
# pattern = "otlp.headroomlabs.ai/v1/logs"
|
||||
# zone_name = "headroomlabs.ai"
|
||||
# custom_domain = false
|
||||
|
||||
# The corpus. R2 rather than S3 specifically for zero egress: training jobs
|
||||
# re-read the whole dataset, and on S3 that is a recurring bill for data we
|
||||
# already own.
|
||||
[[r2_buckets]]
|
||||
binding = "CORPUS"
|
||||
bucket_name = "headroom-telemetry"
|
||||
|
||||
# Optional metrics lane, added later without touching this file:
|
||||
# npx wrangler secret put METRICS_OTLP_URL
|
||||
# npx wrangler secret put METRICS_OTLP_AUTH
|
||||
# Absent = R2 only, which is the right place to start.
|
||||
|
||||
[observability]
|
||||
enabled = true
|
||||
|
||||
# Rate limiting is configured in the Cloudflare dashboard, not here — this
|
||||
# endpoint is unauthenticated by design (anonymity is the product), so it is
|
||||
# the only thing between the Worker and a bored stranger:
|
||||
# Security > WAF > Rate limiting rules
|
||||
# otlp.headroomlabs.ai/v1/logs -> 60 requests / minute / IP
|
||||
# A real client sends ~2 requests/hour, so that is ~1000x headroom while still
|
||||
# capping a single abusive source hard.
|
||||
|
|
@ -36,6 +36,7 @@ from typing import Any
|
|||
from headroom import paths as _paths
|
||||
from headroom import savings_ledger
|
||||
from headroom.cache.compression_store import format_retrieval_miss_detail
|
||||
from headroom.telemetry import session as telemetry_session
|
||||
|
||||
# fcntl is Unix-only; on Windows we skip file locking (stats are best-effort).
|
||||
# Keep the module typed as Any so Windows mypy runs don't try to resolve Unix-only attrs.
|
||||
|
|
@ -795,6 +796,16 @@ class HeadroomMCPServer:
|
|||
client=self._current_client(),
|
||||
source="mcp",
|
||||
)
|
||||
# Anonymous beacon (opt-out; no-op unless enabled). Without this the MCP
|
||||
# path is invisible to aggregate stats even though it is a first-class
|
||||
# way to use Headroom — and subagents make that worse, since each runs
|
||||
# its own MCP server. Swallows its own errors; the caller's try/except
|
||||
# is a second net, not the first.
|
||||
telemetry_session.record_mcp_compression(
|
||||
original_tokens=before,
|
||||
compressed_tokens=after,
|
||||
model=os.environ.get("HEADROOM_MCP_MODEL"),
|
||||
)
|
||||
|
||||
def _current_client(self) -> str:
|
||||
"""Name of the MCP client driving this session (best-effort)."""
|
||||
|
|
|
|||
|
|
@ -314,8 +314,11 @@ class RequestOutcome:
|
|||
async def emit_request_outcome(handler: Any, outcome: RequestOutcome) -> None:
|
||||
"""Single funnel for per-request bookkeeping. The contract.
|
||||
|
||||
Owns the four downstream effects in canonical order:
|
||||
Owns the downstream effects in canonical order:
|
||||
|
||||
0. ``telemetry.session.record_outcome(...)`` — anonymous session beacon
|
||||
(opt-in, no-op unless ``HEADROOM_TELEMETRY`` is on). Runs ahead of the
|
||||
5xx guard below so session error rates see upstream failures.
|
||||
1. ``handler.metrics.record_request(...)`` — Prometheus / SavingsTracker
|
||||
2. ``handler.cost_tracker.record_tokens(...)`` — cost dashboard
|
||||
(skipped when cost_tracker is None, i.e. ``--no-cost``)
|
||||
|
|
@ -324,8 +327,8 @@ async def emit_request_outcome(handler: Any, outcome: RequestOutcome) -> None:
|
|||
4. structured PERF log line — consumed by ``headroom perf``
|
||||
|
||||
A failure outcome (``status_code >= 500``, e.g. a 529 surfaced after retry
|
||||
exhaustion) short-circuits before these four effects: it records a failed
|
||||
request and returns, so an upstream failure cannot feed the success stats.
|
||||
exhaustion) short-circuits before effects 1-4: it records a failed request
|
||||
and returns, so an upstream failure cannot feed the success stats.
|
||||
|
||||
Takes the handler as a free argument rather than ``self`` so this
|
||||
function is callable from:
|
||||
|
|
@ -343,6 +346,7 @@ async def emit_request_outcome(handler: Any, outcome: RequestOutcome) -> None:
|
|||
from headroom.proxy.cost import _summarize_transforms
|
||||
from headroom.proxy.models import RequestLog
|
||||
from headroom.proxy.project_context import get_current_project
|
||||
from headroom.telemetry.session import record_outcome
|
||||
|
||||
# GitHub Copilot: requests routed to the Copilot API travel on the OpenAI or
|
||||
# Anthropic wire, so the handlers stamp the wire provider. Relabel to
|
||||
|
|
@ -357,6 +361,21 @@ async def emit_request_outcome(handler: Any, outcome: RequestOutcome) -> None:
|
|||
|
||||
outcome = dataclasses.replace(outcome, provider="copilot")
|
||||
|
||||
# 0. Anonymous session beacon. Opt-in and a no-op unless HEADROOM_TELEMETRY
|
||||
# is explicitly on, in which case it folds this outcome into an in-memory
|
||||
# per-session aggregate and POSTs one content-free event when the session
|
||||
# goes idle (off-thread — see telemetry.session.post_session_event).
|
||||
#
|
||||
# Placed before the 5xx short-circuit, for the same reason the Copilot
|
||||
# relabel above is: a session's failure count has to see upstream
|
||||
# failures, and per-session error rate is exactly the signal that shows a
|
||||
# provider going flaky for real users. Everything below this point is
|
||||
# success-only bookkeeping.
|
||||
#
|
||||
# Synchronous but allocation-light, and swallows its own exceptions — the
|
||||
# beacon must never add latency to, or take down, the request path.
|
||||
record_outcome(outcome)
|
||||
|
||||
# Upstream failure (>= 500, e.g. a 529 Overloaded surfaced after retry
|
||||
# exhaustion) must not feed the savings/cost/log success stats; that would
|
||||
# let a failed request inflate the save-rate. Record it as failed and stop,
|
||||
|
|
|
|||
|
|
@ -1,14 +1,23 @@
|
|||
"""Telemetry opt-in state for Headroom.
|
||||
"""Telemetry state for Headroom.
|
||||
|
||||
Headroom collects only **local**, aggregate telemetry (tokens saved, compression
|
||||
ratios, timings) for the in-process collector and the ``/stats`` /
|
||||
``/v1/telemetry`` endpoints. **Nothing is sent to Headroom Labs:** the anonymous
|
||||
telemetry beacon that previously shipped aggregate stats has been removed.
|
||||
Operational metrics can still be exported to *your own* OpenTelemetry collector
|
||||
via ``HEADROOM_OTEL_METRICS_*`` (a destination you control).
|
||||
Two independent switches, because they answer two different questions:
|
||||
|
||||
This module holds the ``HEADROOM_TELEMETRY`` opt-in predicate (off by default)
|
||||
that gates local collection, plus the CLI notice helpers.
|
||||
``HEADROOM_TELEMETRY`` — **off by default, opt-in.** Aggregates stats locally
|
||||
for the in-process collector and the ``/stats`` / ``/v1/telemetry`` endpoints.
|
||||
Never leaves the machine.
|
||||
|
||||
``HEADROOM_BEACON`` — **on by default, opt-out.** Uploads an anonymous,
|
||||
content-free session summary to Headroom Labs. Disable with
|
||||
``HEADROOM_BEACON=off``, ``DO_NOT_TRACK=1``, or offline mode. See
|
||||
:mod:`headroom.telemetry.session` for the exact payload and
|
||||
``deploy/beacon/worker.js`` for what the receiver does with it.
|
||||
|
||||
Operational metrics are separate again, and go to *your own* OpenTelemetry
|
||||
collector via ``HEADROOM_OTEL_METRICS_*`` — a destination you control, that
|
||||
Headroom Labs never sees.
|
||||
|
||||
Keeping these separate matters: an operator who turned on local stats has not
|
||||
thereby agreed to upload anything, and must not start doing so on upgrade.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
|
@ -36,6 +45,53 @@ def is_telemetry_enabled() -> bool:
|
|||
return val in _ON_VALUES
|
||||
|
||||
|
||||
# Whether the upload beacon runs when the operator has expressed no preference.
|
||||
#
|
||||
# True makes Headroom opt-OUT: every install uploads anonymous session summaries
|
||||
# unless it is told not to. Flip to False for opt-in and nothing uploads without
|
||||
# an explicit HEADROOM_BEACON=on.
|
||||
#
|
||||
# This single constant is the whole policy — deliberately, so the decision is
|
||||
# one line to audit and one line to reverse.
|
||||
BEACON_DEFAULT_ON = True
|
||||
|
||||
|
||||
def is_beacon_enabled() -> bool:
|
||||
"""Check if the anonymous upload beacon is enabled.
|
||||
|
||||
Default is :data:`BEACON_DEFAULT_ON`. Three things switch it off regardless,
|
||||
in this order:
|
||||
|
||||
* offline mode — no network is no network;
|
||||
* ``DO_NOT_TRACK`` — the cross-tool convention, honoured because a user who
|
||||
has already stated this preference should not have to state it again;
|
||||
* ``HEADROOM_BEACON=off`` (or false/0/no/disable/disabled).
|
||||
|
||||
Deliberately a *separate* switch from ``HEADROOM_TELEMETRY``. That flag
|
||||
means "aggregate stats locally so ``/stats`` has something to show" — data
|
||||
that never leaves the machine. Overloading it to also authorise upload
|
||||
would mean every operator who had turned on local stats begins transmitting
|
||||
the moment they upgrade, having answered a different question.
|
||||
|
||||
Note the asymmetry with :func:`is_telemetry_enabled`, which is fail-closed:
|
||||
this one is fail-open by construction while ``BEACON_DEFAULT_ON`` is True,
|
||||
so an unrecognised value uploads rather than staying silent.
|
||||
"""
|
||||
from headroom.offline import is_offline
|
||||
|
||||
if is_offline():
|
||||
return False
|
||||
dnt = os.environ.get("DO_NOT_TRACK", "").lower().strip()
|
||||
if dnt and dnt not in _OFF_VALUES:
|
||||
return False
|
||||
raw = os.environ.get("HEADROOM_BEACON", "").lower().strip()
|
||||
if raw in _OFF_VALUES:
|
||||
return False
|
||||
if raw in _ON_VALUES:
|
||||
return True
|
||||
return BEACON_DEFAULT_ON
|
||||
|
||||
|
||||
def is_telemetry_warn_enabled() -> bool:
|
||||
"""Check if telemetry warnings are enabled (feature flag, on by default).
|
||||
|
||||
|
|
@ -56,8 +112,30 @@ def format_telemetry_notice(*, prefix: str = "") -> str:
|
|||
Returns an empty string when telemetry or warnings are disabled so callers
|
||||
can unconditionally include the result in their output.
|
||||
"""
|
||||
if not is_telemetry_enabled() or not is_telemetry_warn_enabled():
|
||||
if not is_telemetry_warn_enabled():
|
||||
return ""
|
||||
|
||||
beacon = is_beacon_enabled()
|
||||
local = is_telemetry_enabled()
|
||||
if not beacon and not local:
|
||||
return ""
|
||||
|
||||
# The beacon line comes first and is unconditional when on. It is opt-out,
|
||||
# so this notice is the only place a user learns it is running at all —
|
||||
# surprise is what turns anonymous telemetry into a trust incident.
|
||||
#
|
||||
# Wording: lead with what the data *is* (compression ratios) and why it is
|
||||
# useful (understanding when compression works), not with the fact of
|
||||
# transmission. "Stats are sent to Headroom Labs" is accurate but reads like
|
||||
# extraction, which misrepresents a payload that is entirely counters. The
|
||||
# reassurance has to stay concrete, though — "never prompts, code, or file
|
||||
# paths" names the things people actually worry about, and every one of
|
||||
# those is enforced by the allowlist, not just promised here.
|
||||
if beacon:
|
||||
return (
|
||||
f"{prefix}Telemetry: anonymous compression stats — never prompts, code, or "
|
||||
"file paths. Helps us improve compression | Off: HEADROOM_BEACON=off"
|
||||
)
|
||||
return (
|
||||
f"{prefix}Telemetry: ENABLED (local aggregate stats only — nothing sent externally) | "
|
||||
"Disable: HEADROOM_TELEMETRY=off or --no-telemetry"
|
||||
|
|
|
|||
1056
headroom/telemetry/session.py
Normal file
1056
headroom/telemetry/session.py
Normal file
File diff suppressed because it is too large
Load diff
|
|
@ -30,6 +30,21 @@ def _scrub_developer_headroom_env(monkeypatch):
|
|||
monkeypatch.delenv("ANTHROPIC_CUSTOM_HEADERS", raising=False)
|
||||
|
||||
|
||||
# The scrub above deletes every HEADROOM_* var — which includes HEADROOM_BEACON,
|
||||
# and the beacon defaults to ON. So scrubbing for hermeticity is precisely what
|
||||
# switches it on, and with HEADROOM_TELEMETRY_ENDPOINT scrubbed too it falls back
|
||||
# to the real production endpoint. Every test that reaches the outcome funnel
|
||||
# then POSTs a session event for real: observed writing into the live corpus
|
||||
# during a local run, and CI would do the same on every push.
|
||||
#
|
||||
# Depends on the scrub fixture so it is guaranteed to run after it rather than
|
||||
# relying on declaration order. A test that wants the beacon on just sets the
|
||||
# var itself — monkeypatch inside the test wins over this.
|
||||
@pytest.fixture(autouse=True)
|
||||
def _disable_telemetry_beacon(monkeypatch, _scrub_developer_headroom_env):
|
||||
monkeypatch.setenv("HEADROOM_BEACON", "off")
|
||||
|
||||
|
||||
# The MCP install ledger defaults to ``~/.headroom/mcp_installs.json``, so any
|
||||
# test that registers a server (directly or through `wrap`) writes into the
|
||||
# developer's REAL ledger — observed adding a live `claude/serena` entry during a
|
||||
|
|
|
|||
|
|
@ -53,6 +53,7 @@ class TestFormatTelemetryNotice:
|
|||
|
||||
def test_returns_notice_when_telemetry_on(self, monkeypatch):
|
||||
monkeypatch.setenv("HEADROOM_TELEMETRY", "on")
|
||||
monkeypatch.setenv("HEADROOM_BEACON", "off")
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
notice = format_telemetry_notice()
|
||||
assert notice != ""
|
||||
|
|
@ -62,6 +63,43 @@ class TestFormatTelemetryNotice:
|
|||
|
||||
def test_empty_when_telemetry_off(self, monkeypatch):
|
||||
monkeypatch.setenv("HEADROOM_TELEMETRY", "off")
|
||||
monkeypatch.setenv("HEADROOM_BEACON", "off")
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
assert format_telemetry_notice() == ""
|
||||
|
||||
def test_beacon_is_announced_by_default(self, monkeypatch):
|
||||
"""The beacon is opt-out, so the notice is the only place a user finds
|
||||
out it is running. Silence here is how anonymous telemetry becomes a
|
||||
trust incident."""
|
||||
for var in ("HEADROOM_TELEMETRY", "HEADROOM_BEACON", "DO_NOT_TRACK"):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
notice = format_telemetry_notice()
|
||||
assert "compression stats" in notice
|
||||
assert "HEADROOM_BEACON=off" in notice
|
||||
|
||||
def test_beacon_notice_names_what_is_not_sent(self, monkeypatch):
|
||||
"""Vague reassurance is worse than none. The notice has to name the
|
||||
three things users actually worry about."""
|
||||
for var in ("HEADROOM_TELEMETRY", "HEADROOM_BEACON", "DO_NOT_TRACK"):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
notice = format_telemetry_notice()
|
||||
assert "never prompts" in notice
|
||||
assert "code" in notice
|
||||
assert "file paths" in notice
|
||||
|
||||
def test_silent_when_beacon_disabled_and_no_local(self, monkeypatch):
|
||||
monkeypatch.setenv("HEADROOM_BEACON", "off")
|
||||
for var in ("HEADROOM_TELEMETRY", "DO_NOT_TRACK"):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
assert format_telemetry_notice() == ""
|
||||
|
||||
def test_do_not_track_silences_the_beacon_notice(self, monkeypatch):
|
||||
monkeypatch.setenv("DO_NOT_TRACK", "1")
|
||||
for var in ("HEADROOM_TELEMETRY", "HEADROOM_BEACON"):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
assert format_telemetry_notice() == ""
|
||||
|
||||
|
|
@ -176,6 +214,7 @@ class TestWrapCLITelemetryNotice:
|
|||
|
||||
def test_print_notice_outputs_when_telemetry_on(self, monkeypatch, capsys):
|
||||
monkeypatch.setenv("HEADROOM_TELEMETRY", "on")
|
||||
monkeypatch.setenv("HEADROOM_BEACON", "off")
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
|
||||
from headroom.cli.wrap import _print_telemetry_notice
|
||||
|
|
@ -185,8 +224,21 @@ class TestWrapCLITelemetryNotice:
|
|||
assert "Telemetry" in captured.out
|
||||
assert "HEADROOM_TELEMETRY=off" in captured.out
|
||||
|
||||
def test_print_notice_announces_beacon_by_default(self, monkeypatch, capsys):
|
||||
for var in ("HEADROOM_TELEMETRY", "HEADROOM_BEACON", "DO_NOT_TRACK"):
|
||||
monkeypatch.delenv(var, raising=False)
|
||||
monkeypatch.delenv("HEADROOM_TELEMETRY_WARN", raising=False)
|
||||
|
||||
from headroom.cli.wrap import _print_telemetry_notice
|
||||
|
||||
_print_telemetry_notice()
|
||||
captured = capsys.readouterr()
|
||||
assert "compression stats" in captured.out
|
||||
assert "HEADROOM_BEACON=off" in captured.out
|
||||
|
||||
def test_print_notice_silent_when_telemetry_off(self, monkeypatch, capsys):
|
||||
monkeypatch.setenv("HEADROOM_TELEMETRY", "off")
|
||||
monkeypatch.setenv("HEADROOM_BEACON", "off")
|
||||
|
||||
from headroom.cli.wrap import _print_telemetry_notice
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue