Skip to content

Commit 7c476fd

Browse files
authored
feat: add anonymous opt-out telemetry (#684)
# Anonymous opt-out telemetry ## Summary Ships the researched daily anonymous telemetry (plan: `Plans/2026-09-18-anonymous-telemetry.md`, spec: `Runs/anonymous-telemetry-research/research.md`). - **SQL collector** — new `pgflow_telemetry` schema: `sent_reports` audit/dedup table, bucket helpers, `build_payload()` covering all 26 metrics, `preview()` (never sends), `report()` with gates (inactive day → already reported → local), one `pg_net` attempt with 5 s timeout, no retries, and `enable()`/`disable()` where the `cron.job` row is the switch. The Atlas migration schedules the daily job exactly once; later migrations never re-schedule it. - **Version stamping** — `pgflow.workers.pgflow_version` column; edge workers stamp their package version at registration (Deno/JSR-compatible JSON import, vendored for e2e). - **Cloudflare ingest worker** — `apps/telemetry-worker`: strict allowlist validation (metric names, bucket strings, semver), one Analytics Engine data point per contribution, no cookies, no body logging, 2 KB / 64-contribution caps. - **Privacy** — payloads contain only allowlisted metric names, bucket strings, semver versions, and count buckets. Never slugs, names, IDs, hosts, URLs, inputs, outputs, errors, IPs, or free text. Every sent payload stays auditable in `pgflow_telemetry.sent_reports`; `preview()` shows any day's exact bytes; local dev stacks never send. - **Docs & release** — new `/reference/telemetry/` page, sidebar + reference card + redirect for the renamed news article (combined 0.17.1), changeset (`@pgflow/core`, `@pgflow/edge-worker`: patch). **Deferred (manual follow-up, not in this PR):** `wrangler deploy` of the ingest worker and the live endpoint smoke test. The SQL endpoint constant is `https://pgflow-telemetry.workers.dev`; if the deployed subdomain differs, it needs a follow-up change per the plan's Task 8 Step 5. ## Checks - 4 new pgTAP files (red/green cycles; migration replay with dependencies): `telemetry/{workers_version,schema,payload,report}.test.sql` - Full pgTAP suite: 314 files, 1606 tests — pass (one pre-existing timing-flaky perf test passes standalone) - `nx affected -t build test` — pass, except `demo:test`: 4 Groq-API tests fail only when a `GROQ_API_KEY` is present in the environment (model `llama-3.1-8b-instant` decommissioned); demo sources untouched by this PR and CI has no such key - telemetry-worker: 10 vitest tests + typecheck — pass - edge-worker: build + 265 tests, JSR publish dry-run, e2e 13/13 — pass - Migration verification: `verify-migrations`, `gen-types`, `verify-gen-types` — pass - Live local round-trip: `preview()` shows real contributions; `report()` returns `skipped: local`, zero audit rows; worker registration stamps `pgflow_version = 0.17.0` - Website build with link validator and redirect check — pass - Lint on all touched packages — pass
2 parents 48fdb7b + 56bef34 commit 7c476fd

42 files changed

Lines changed: 3220 additions & 28 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.changeset/anonymous-telemetry.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
'@pgflow/core': patch
3+
'@pgflow/edge-worker': patch
4+
---
5+
6+
Add anonymous, opt-out telemetry. A daily `pg_cron` job aggregates yesterday's `pgflow.*` activity into coarse, identifier-free buckets (versions in use, run and worker counts, flow shapes, feature adoption, durations) and sends one small payload through `pg_net` to a Cloudflare Worker backed by Workers Analytics Engine. Workers stamp their package version at registration. Nothing identifying is ever sent; every payload is stored locally in `pgflow_telemetry.sent_reports` for audit; `pgflow_telemetry.preview()` shows any day's payload without sending; `pgflow_telemetry.disable()` opts out permanently (the cron job row is the switch) and local development databases never report.

‎.changeset/config.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,5 +15,5 @@
1515
"access": "public",
1616
"baseBranch": "main",
1717
"updateInternalDependencies": "patch",
18-
"ignore": ["@pgflow/demo", "@pgflow/website", "@pgflow/example-flows"]
18+
"ignore": ["@pgflow/demo", "@pgflow/website", "@pgflow/example-flows", "@pgflow/telemetry-worker"]
1919
}

‎apps/telemetry-worker/package.json‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
{
2+
"name": "@pgflow/telemetry-worker",
3+
"private": true,
4+
"type": "module",
5+
"scripts": {
6+
"typecheck": "tsc --noEmit",
7+
"test": "vitest run",
8+
"deploy": "wrangler deploy"
9+
},
10+
"devDependencies": {
11+
"@cloudflare/workers-types": "^4.20250901.0",
12+
"typescript": "^5.5.0",
13+
"vitest": "^3.0.0",
14+
"wrangler": "^4.0.0"
15+
}
16+
}
Lines changed: 280 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,280 @@
1+
import { describe, expect, it } from 'vitest';
2+
import worker, { type Env } from './index';
3+
// Worst-case payload produced by pgflow_telemetry.preview() on the real
4+
// SQL (65 runs with distinct millisecond durations + 70 worker versions):
5+
// 31 bounded contributions, 2015 bytes. Regenerate from a local database
6+
// whenever the sender schema changes.
7+
import fixture from './sender-fixture.json';
8+
9+
function makeEnv() {
10+
const points: Array<{ indexes: string[]; blobs: string[]; doubles: number[] }> = [];
11+
const env = {
12+
PGFLOW_TELEMETRY: {
13+
writeDataPoint: (p: { indexes: string[]; blobs: string[]; doubles: number[] }) =>
14+
points.push(p),
15+
},
16+
} as unknown as Env;
17+
return { env, points };
18+
}
19+
20+
function post(
21+
body: unknown,
22+
env: Env,
23+
raw = false,
24+
contentType = 'application/json',
25+
): Promise<Response> {
26+
return worker.fetch(
27+
new Request('https://pgflow-telemetry.workers.dev/', {
28+
method: 'POST',
29+
headers: { 'content-type': contentType },
30+
body: raw ? (body as string) : JSON.stringify(body),
31+
}),
32+
env,
33+
);
34+
}
35+
36+
const valid = {
37+
schema: 1,
38+
contributions: [
39+
{ metric: 'runs_started_day', bucket: '2-3' },
40+
{ metric: 'workers_by_version', bucket: '0.17.0', count: '1' },
41+
{ metric: 'uses_map', bucket: 'yes' },
42+
],
43+
};
44+
45+
describe('telemetry ingest', () => {
46+
it('accepts a valid payload, writes one point per contribution, returns 204', async () => {
47+
const { env, points } = makeEnv();
48+
const res = await post(valid, env);
49+
expect(res.status).toBe(204);
50+
expect(points).toHaveLength(3);
51+
});
52+
53+
it('stores the count bucket as a second blob', async () => {
54+
const { env, points } = makeEnv();
55+
await post(valid, env);
56+
expect(points[1]).toEqual({
57+
indexes: ['workers_by_version'],
58+
blobs: ['0.17.0', '1'],
59+
doubles: [1],
60+
});
61+
});
62+
63+
it('writes a single blob for uncounted contributions', async () => {
64+
const { env, points } = makeEnv();
65+
await post(valid, env);
66+
expect(points[0]).toEqual({
67+
indexes: ['runs_started_day'],
68+
blobs: ['2-3'],
69+
doubles: [1],
70+
});
71+
expect(points[2].blobs).toEqual(['yes']);
72+
});
73+
74+
it('rejects GET with 405', async () => {
75+
const { env } = makeEnv();
76+
const res = await worker.fetch(
77+
new Request('https://pgflow-telemetry.workers.dev/', { method: 'GET' }),
78+
env,
79+
);
80+
expect(res.status).toBe(405);
81+
});
82+
83+
it('rejects malformed JSON with 400', async () => {
84+
const { env, points } = makeEnv();
85+
const res = await post('{not json', env, true);
86+
expect(res.status).toBe(400);
87+
expect(points).toHaveLength(0);
88+
});
89+
90+
it('rejects wrong schema version with 400', async () => {
91+
const { env, points } = makeEnv();
92+
const res = await post({ schema: 2, contributions: valid.contributions }, env);
93+
expect(res.status).toBe(400);
94+
expect(points).toHaveLength(0);
95+
});
96+
97+
it('rejects null and array roots with 400', async () => {
98+
const { env, points } = makeEnv();
99+
expect((await post('null', env, true)).status).toBe(400);
100+
expect((await post('[1,2]', env, true)).status).toBe(400);
101+
expect(points).toHaveLength(0);
102+
});
103+
104+
it('rejects unknown root keys with 400', async () => {
105+
const { env, points } = makeEnv();
106+
const res = await post(
107+
{ schema: 1, contributions: valid.contributions, extra: 'x' },
108+
env,
109+
);
110+
expect(res.status).toBe(400);
111+
expect(points).toHaveLength(0);
112+
});
113+
114+
it('rejects unknown contribution keys with 400', async () => {
115+
const { env, points } = makeEnv();
116+
const res = await post(
117+
{
118+
schema: 1,
119+
contributions: [{ metric: 'uses_map', bucket: 'yes', surprise: 1 }],
120+
},
121+
env,
122+
);
123+
expect(res.status).toBe(400);
124+
expect(points).toHaveLength(0);
125+
});
126+
127+
it('rejects prototype-named metrics with 400', async () => {
128+
const { env, points } = makeEnv();
129+
for (const metric of ['constructor', '__proto__', 'toString', 'hasOwnProperty']) {
130+
const res = await post(
131+
{ schema: 1, contributions: [{ metric, bucket: 'yes' }] },
132+
env,
133+
);
134+
expect(res.status, `metric ${metric}`).toBe(400);
135+
}
136+
expect(points).toHaveLength(0);
137+
});
138+
139+
it('rejects non-string metrics with 400', async () => {
140+
const { env, points } = makeEnv();
141+
const res = await post(
142+
{ schema: 1, contributions: [{ metric: 42, bucket: 'yes' }] },
143+
env,
144+
);
145+
expect(res.status).toBe(400);
146+
expect(points).toHaveLength(0);
147+
});
148+
149+
it('rejects unknown metric with 400', async () => {
150+
const { env } = makeEnv();
151+
const res = await post({
152+
schema: 1,
153+
contributions: [{ metric: 'flow_slug', bucket: 'yes' }],
154+
}, env);
155+
expect(res.status).toBe(400);
156+
});
157+
158+
it('rejects arbitrary bucket strings with 400', async () => {
159+
const { env } = makeEnv();
160+
const res = await post({
161+
schema: 1,
162+
contributions: [{ metric: 'runs_started_day', bucket: 'banana' }],
163+
}, env);
164+
expect(res.status).toBe(400);
165+
});
166+
167+
it('rejects non-semantic-version buckets with 400', async () => {
168+
const { env } = makeEnv();
169+
const res = await post({
170+
schema: 1,
171+
contributions: [{ metric: 'workers_by_version', bucket: 'my-flow-name' }],
172+
}, env);
173+
expect(res.status).toBe(400);
174+
});
175+
176+
it('rejects non-bucket count values with 400', async () => {
177+
const { env, points } = makeEnv();
178+
const res = await post({
179+
schema: 1,
180+
contributions: [{ metric: 'workers_by_version', bucket: '0.17.0', count: 7 }],
181+
}, env);
182+
expect(res.status).toBe(400);
183+
expect(points).toHaveLength(0);
184+
});
185+
186+
it('rejects duplicate metric/bucket pairs with 400 before any write', async () => {
187+
const { env, points } = makeEnv();
188+
const res = await post({
189+
schema: 1,
190+
contributions: [
191+
{ metric: 'uses_map', bucket: 'yes' },
192+
{ metric: 'uses_map', bucket: 'yes' },
193+
{ metric: 'runs_started_day', bucket: '2-3' },
194+
],
195+
}, env);
196+
expect(res.status).toBe(400);
197+
expect(points).toHaveLength(0);
198+
});
199+
200+
it('rejects a repeated pair even when the counts differ', async () => {
201+
const { env, points } = makeEnv();
202+
const res = await post({
203+
schema: 1,
204+
contributions: [
205+
{ metric: 'workers_by_version', bucket: '0.17.0', count: '1' },
206+
{ metric: 'workers_by_version', bucket: '0.17.0', count: '2-3' },
207+
],
208+
}, env);
209+
expect(res.status).toBe(400);
210+
expect(points).toHaveLength(0);
211+
});
212+
213+
it('accepts application/json with charset parameters', async () => {
214+
const { env } = makeEnv();
215+
const res = await post(valid, env, false, 'application/json; charset=utf-8');
216+
expect(res.status).toBe(204);
217+
});
218+
219+
it('rejects media types that merely mention application/json', async () => {
220+
const { env, points } = makeEnv();
221+
for (const contentType of [
222+
'text/plain; application/json',
223+
'application/jsonp',
224+
'text/application/json',
225+
]) {
226+
const res = await post(valid, env, false, contentType);
227+
expect(res.status, contentType).toBe(415);
228+
}
229+
expect(points).toHaveLength(0);
230+
});
231+
232+
it('rejects oversized bodies with 413', async () => {
233+
const { env } = makeEnv();
234+
const big = 'x'.repeat(2049);
235+
const res = await post(big, env, true);
236+
expect(res.status).toBe(413);
237+
});
238+
239+
it('measures body bytes, not UTF-16 code units', async () => {
240+
const { env, points } = makeEnv();
241+
// 1100 two-byte characters: 2200+ UTF-8 bytes, ~1130 UTF-16 units.
242+
const multibyte = '{"schema":1,"contributions":[{"metric":"' +
243+
'é'.repeat(1100) + '","bucket":"yes"}]}';
244+
expect(multibyte.length).toBeLessThan(2048);
245+
const encoder = new TextEncoder();
246+
expect(encoder.encode(multibyte).byteLength).toBeGreaterThan(2048);
247+
const res = await post(multibyte, env, true);
248+
expect(res.status).toBe(413);
249+
expect(points).toHaveLength(0);
250+
});
251+
252+
it('rejects more than 64 contributions', async () => {
253+
const { env, points } = makeEnv();
254+
const res = await post({
255+
schema: 1,
256+
contributions: Array.from({ length: 65 }, () => ({
257+
metric: 'uses_map',
258+
bucket: 'yes',
259+
})),
260+
}, env);
261+
// 65 valid contributions always exceed the 2 KB body cap, so the size
262+
// guard (413) fires before the count guard. MAX_CONTRIBUTIONS stays as
263+
// defense-in-depth below the byte cap.
264+
expect(res.status).toBe(413);
265+
expect(points).toHaveLength(0);
266+
});
267+
268+
it('accepts a real sender-produced worst-case payload', async () => {
269+
const { env, points } = makeEnv();
270+
const res = await post(fixture, env);
271+
expect(res.status).toBe(204);
272+
expect(points).toHaveLength(fixture.contributions.length);
273+
});
274+
275+
it('sets no cookies on success', async () => {
276+
const { env } = makeEnv();
277+
const res = await post(valid, env);
278+
expect(res.headers.getSetCookie()).toHaveLength(0);
279+
});
280+
});

0 commit comments

Comments
 (0)