headroom/deploy/beacon/worker.js
Tejas Chopra 9cfb00838a
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>
2026-08-03 05:43:25 -07:00

174 lines
6 KiB
JavaScript

/**
* 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 });
},
};