diff --git a/services/diagnostics-api/README.md b/services/diagnostics-api/README.md index 3f48f08..ac8ab90 100644 --- a/services/diagnostics-api/README.md +++ b/services/diagnostics-api/README.md @@ -1,13 +1,13 @@ # VniDrop diagnostics API -Cloudflare Worker for ingesting batched telemetry, crash reports, and user-submitted -bug reports. D1 stores searchable metadata; R2 stores larger stack traces and logs. +Cloudflare Worker for ingesting user-submitted bug reports. D1 stores searchable +metadata; R2 stores the larger attached logs. The service is designed for modest traffic and low operating cost: -- one D1 row is written per telemetry batch, not per event; -- crash stacks and bug logs are stored in R2 instead of D1; -- request and batch limits reject oversized work before storage writes; +- one D1 row is written per bug report; +- bug logs are stored in R2 instead of D1; +- request limits reject oversized work before storage writes; - an hourly scheduled cleanup and an R2 lifecycle rule enforce retention; - no Queue, Durable Object, or KV resources are required. @@ -35,15 +35,13 @@ X-VniDrop-Install-Id: |--------|------|------| | `GET` | `/live` | process liveness; does not touch storage | | `GET` | `/health` | authenticated readiness; checks required configuration and the D1 schema | -| `POST` | `/v1/events` | `{ batchId, installId, appVersion?, platform?, events: [...] }` | -| `POST` | `/v1/crashes` | app crash payload | | `POST` | `/v1/bugs` | app bug-report payload | -Batch and report IDs are client-generated UUIDs. A client must reuse the same ID -when retrying so D1 can acknowledge the request without storing it twice. +Report IDs are client-generated UUIDs. A client must reuse the same ID when +retrying so D1 can acknowledge the request without storing it twice. -Accepted reports return `202`. Defaults are a 262,144-byte request limit and at -most 50 events per batch. Cloudflare rate-limit bindings allow 30 requests per +Accepted reports return `202`. The default is a 262,144-byte request limit. +Cloudflare rate-limit bindings allow 30 requests per installation and 120 requests per source, per ingest route, per minute. Source limits run before shared-key verification so rejected traffic is bounded too. These counters are eventually consistent and local to a Cloudflare location, so @@ -153,11 +151,11 @@ migrations to the isolated local database assigned to each test file. `RETENTION_DAYS` defaults to 90. The `17 * * * *` cron trigger runs cleanup at 17 minutes past every hour. Cleanup works in bounded batches: it deletes each expired report's referenced R2 object before deleting that exact D1 row. The R2 -lifecycle rule is an independent backstop for stack and log objects, including -objects left behind by a partial ingest failure. Each scheduled run can remove -8,000 event batches and 7,200 rows from each report table while staying below -D1's per-invocation query ceiling. Later hourly runs continue any backlog. -Reaching the cap emits a structured warning with the remaining expired-row counts; +lifecycle rule is an independent backstop for log objects, including objects left +behind by a partial ingest failure. Each scheduled run can remove 7,200 bug rows +while staying below D1's per-invocation query ceiling. Later hourly runs continue +any backlog. +Reaching the cap emits a structured warning with the remaining expired-row count; alert on that warning because retention is necessarily best-effort during sustained distributed abuse. @@ -179,23 +177,20 @@ vnidrop.diagnostics.ingestKey= Both the endpoint and key are required. When both are empty the app uses its offline-safe no-op transport; configuring only one fails the Gradle build. -`vnidrop.diagnostics.included=false` disables -automatic telemetry and crash upload, but a configured endpoint can still accept -an explicit user-submitted bug report. Treat the app-side key as an abuse-control -token with the limitations described above. +`vnidrop.diagnostics.included=false` routes bug reports to that no-op transport +(never sent); a configured endpoint accepts an explicit user-submitted bug report. +Treat the app-side key as an abuse-control token with the limitations described +above. ## Reading reports ```bash -npx wrangler d1 execute vnidrop-diagnostics --remote \ - --command "SELECT id, exception_type, platform, occurred_at FROM crashes ORDER BY occurred_at DESC LIMIT 20" - npx wrangler d1 execute vnidrop-diagnostics --remote \ --command "SELECT id, what_happened, status, occurred_at FROM bugs WHERE status = 'open' ORDER BY occurred_at DESC LIMIT 20" ``` -R2 object keys use `crashes///stack.txt` and -`bugs///logs.txt`. The unique attempt segment prevents a retry -from overwriting an already accepted object before D1 detects the duplicate. +R2 object keys use `bugs///logs.txt`. The unique attempt segment +prevents a retry from overwriting an already accepted object before D1 detects +the duplicate. There is no public administration endpoint; inspect reports through authenticated Cloudflare tools or a future Access-protected dashboard. diff --git a/services/diagnostics-api/migrations/0003_drop_telemetry_tables.sql b/services/diagnostics-api/migrations/0003_drop_telemetry_tables.sql new file mode 100644 index 0000000..946eff8 --- /dev/null +++ b/services/diagnostics-api/migrations/0003_drop_telemetry_tables.sql @@ -0,0 +1,10 @@ +-- Telemetry and crash auto-reporting were removed from the app; only user-initiated +-- bug reports remain. Drop the now-unused ingestion tables and their indexes. +DROP INDEX IF EXISTS idx_event_batches_received; +DROP INDEX IF EXISTS idx_event_batches_install; +DROP TABLE IF EXISTS event_batches; + +DROP INDEX IF EXISTS idx_crashes_received; +DROP INDEX IF EXISTS idx_crashes_fingerprint; +DROP INDEX IF EXISTS idx_crashes_install; +DROP TABLE IF EXISTS crashes; diff --git a/services/diagnostics-api/src/index.ts b/services/diagnostics-api/src/index.ts index f83b5fe..d22c11f 100644 --- a/services/diagnostics-api/src/index.ts +++ b/services/diagnostics-api/src/index.ts @@ -1,20 +1,15 @@ import { normalizeBug, - normalizeCrash, - normalizeEvents, readJsonObject, } from "./input"; import { type DiagnosticsEnv, runRetention, storeBug, - storeCrash, - storeEvents, } from "./storage"; const DEFAULT_MAX_BODY_BYTES = 262_144; const HARD_MAX_BODY_BYTES = 1_048_576; -const DEFAULT_MAX_EVENTS = 50; export default { async fetch(request: Request, env: DiagnosticsEnv, _ctx: ExecutionContext): Promise { @@ -72,41 +67,6 @@ export default { if (!parsed.ok) return json({ error: parsed.error }, parsed.status, requestId); switch (url.pathname) { - case "/v1/events": { - const maxEvents = boundedPositiveInt(env.MAX_EVENTS_PER_BATCH, DEFAULT_MAX_EVENTS, 1, 100); - const normalized = normalizeEvents(parsed.value, maxEvents); - if (!normalized.ok) { - return json({ error: normalized.error }, normalized.status, requestId); - } - const result = await storeEvents(normalized.value, env); - return json( - { - ok: true, - id: result.id, - stored: result.stored, - duplicate: result.duplicate, - }, - 202, - requestId, - ); - } - case "/v1/crashes": { - const normalized = normalizeCrash(parsed.value); - if (!normalized.ok) { - return json({ error: normalized.error }, normalized.status, requestId); - } - const result = await storeCrash(normalized.value, env); - return json( - { - ok: true, - id: result.id, - fingerprint: result.fingerprint, - duplicate: result.duplicate, - }, - 202, - requestId, - ); - } case "/v1/bugs": { const normalized = normalizeBug(parsed.value); if (!normalized.ok) { @@ -159,12 +119,6 @@ async function readiness(env: DiagnosticsEnv, requestId: string): Promise { diff --git a/services/diagnostics-api/src/input.ts b/services/diagnostics-api/src/input.ts index d5c8983..e4ea0b4 100644 --- a/services/diagnostics-api/src/input.ts +++ b/services/diagnostics-api/src/input.ts @@ -10,13 +10,6 @@ export type InputResult = { ok: true; value: T } | InputFailure; export type NormalizedProperties = Record; -export interface NormalizedEvent { - name: string; - timestampMillis: number; - properties: NormalizedProperties; - schemaVersion: 1; -} - export interface NormalizedBreadcrumb { name: string; timestampMillis: number; @@ -31,28 +24,6 @@ export interface NormalizedDevice { batteryLevel: string; } -export interface NormalizedEventsPayload { - batchId: string; - installId: string; - appVersion: string; - platform: string; - events: NormalizedEvent[]; -} - -export interface NormalizedCrashPayload { - id: string; - installId: string; - appVersion: string; - platform: string; - exceptionType: string; - exceptionMessage: string; - stackTrace: string; - occurredAt: number; - diagnosticsEnabledAtCapture: boolean; - breadcrumbs: NormalizedBreadcrumb[]; - schemaVersion: 1; -} - export interface NormalizedBugPayload { id: string; installId: string; @@ -164,119 +135,6 @@ export async function readJsonObject( return success(parsed); } -export function normalizeEvents( - body: JsonObject, - maxEvents = 50, -): InputResult { - if (!isPlainObject(body)) return failure(400, "invalid_body"); - if (!Number.isSafeInteger(maxEvents) || maxEvents <= 0) { - throw new RangeError("maxEvents must be a positive safe integer"); - } - - const batchId = idField(body, ["batchId", "batch_id"], "invalid_batch_id"); - if (!batchId.ok) return batchId; - const installId = installIdField(body); - if (!installId.ok) return installId; - const appVersion = stringField(body, ["appVersion", "app_version"], 40, "invalid_app_version"); - if (!appVersion.ok) return appVersion; - const platform = stringField(body, ["platform"], 40, "invalid_platform"); - if (!platform.ok) return platform; - - const batchSchema = schemaVersion(body); - if (!batchSchema.ok) return batchSchema; - const rawEvents = pick(body, ["events"]); - if (!Array.isArray(rawEvents)) return failure(400, "invalid_events"); - if (rawEvents.length === 0) return failure(400, "empty_batch"); - if (rawEvents.length > maxEvents) return failure(400, "batch_too_large"); - - const events: NormalizedEvent[] = []; - for (const rawEvent of rawEvents) { - const event = normalizeEvent(rawEvent); - if (!event.ok) return event; - events.push(event.value); - } - - return success({ - batchId: batchId.value, - installId: installId.value, - appVersion: appVersion.value, - platform: platform.value, - events, - }); -} - -export function normalizeCrash(body: JsonObject): InputResult { - if (!isPlainObject(body)) return failure(400, "invalid_body"); - - const id = idField(body, ["id"], "invalid_id"); - if (!id.ok) return id; - const installId = installIdField(body); - if (!installId.ok) return installId; - const appVersion = stringField(body, ["appVersion", "app_version"], 40, "invalid_app_version"); - if (!appVersion.ok) return appVersion; - const platform = stringField(body, ["platform"], 40, "invalid_platform"); - if (!platform.ok) return platform; - const exceptionType = stringField( - body, - ["exceptionType", "exception_type"], - 120, - "invalid_exception_type", - true, - true, - ); - if (!exceptionType.ok) return exceptionType; - const exceptionMessage = stringField( - body, - ["exceptionMessage", "exception_message"], - 2_000, - "invalid_exception_message", - true, - ); - if (!exceptionMessage.ok) return exceptionMessage; - const stackTrace = stringField( - body, - ["stackTrace", "stack_trace"], - 32_000, - "invalid_stack_trace", - true, - ); - if (!stackTrace.ok) return stackTrace; - const occurredAt = timestampField( - body, - ["timestampMillis", "timestamp_millis", "occurredAt", "occurred_at"], - ); - if (!occurredAt.ok) return occurredAt; - const diagnosticsEnabled = booleanField( - body, - [ - "diagnosticsEnabledAtCapture", - "diagnostics_enabled_at_capture", - "diagnostics_enabled", - ], - "invalid_diagnostics_enabled", - true, - ); - if (!diagnosticsEnabled.ok) return diagnosticsEnabled; - const version = schemaVersion(body); - if (!version.ok) return version; - const breadcrumbs = normalizeBreadcrumbs(pick(body, ["breadcrumbs"])); - if (!breadcrumbs.ok) return breadcrumbs; - - return success({ - id: id.value, - installId: installId.value, - appVersion: appVersion.value, - platform: platform.value, - exceptionType: exceptionType.value, - exceptionMessage: exceptionMessage.value, - stackTrace: stackTrace.value, - occurredAt: occurredAt.value, - diagnosticsEnabledAtCapture: diagnosticsEnabled.value, - breadcrumbs: breadcrumbs.value, - schemaVersion: version.value, - }); -} - export function normalizeBug(body: JsonObject): InputResult { if (!isPlainObject(body)) return failure(400, "invalid_body"); @@ -341,24 +199,6 @@ export function normalizeBug(body: JsonObject): InputResult { - if (!isPlainObject(raw)) return failure(400, "invalid_event"); - const name = stringField(raw, ["name"], 64, "invalid_event", true, true); - if (!name.ok) return name; - const timestamp = timestampField(raw, ["timestampMillis", "timestamp_millis", "ts"]); - if (!timestamp.ok) return failure(400, "invalid_event"); - const properties = normalizeProperties(pick(raw, ["properties", "props"]), "invalid_event"); - if (!properties.ok) return properties; - const version = schemaVersion(raw); - if (!version.ok) return version; - return success({ - name: name.value, - timestampMillis: timestamp.value, - properties: properties.value, - schemaVersion: version.value, - }); -} - function normalizeBreadcrumbs(raw: unknown | typeof MISSING): InputResult { if (raw === MISSING) return success([]); if (!Array.isArray(raw)) return failure(400, "invalid_breadcrumbs"); diff --git a/services/diagnostics-api/src/storage.ts b/services/diagnostics-api/src/storage.ts index 30fd2fe..d5df546 100644 --- a/services/diagnostics-api/src/storage.ts +++ b/services/diagnostics-api/src/storage.ts @@ -1,12 +1,7 @@ -import type { - NormalizedBugPayload, - NormalizedCrashPayload, - NormalizedEventsPayload, -} from "./input"; +import type { NormalizedBugPayload } from "./input"; export type DiagnosticsEnv = Cloudflare.Env & { INGEST_KEY?: string; - AE?: AnalyticsEngineDataset; }; export interface StoreResult { @@ -15,130 +10,6 @@ export interface StoreResult { stored: number; } -export async function storeEvents( - payload: NormalizedEventsPayload, - env: DiagnosticsEnv, -): Promise { - const result = await env.DB.prepare( - `INSERT INTO event_batches (id, received_at, install_id, app_version, platform, event_count, payload_json) - VALUES (?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(id) DO NOTHING`, - ) - .bind( - payload.batchId, - Date.now(), - payload.installId, - payload.appVersion, - payload.platform, - payload.events.length, - JSON.stringify(payload.events), - ) - .run(); - const duplicate = result.meta.changes === 0; - if (!duplicate && env.AE) { - try { - for (const event of payload.events) { - env.AE.writeDataPoint({ - blobs: [ - event.name, - payload.platform, - payload.appVersion, - payload.installId, - JSON.stringify(event.properties), - payload.batchId, - ], - doubles: [event.timestampMillis, event.schemaVersion], - indexes: [payload.installId], - }); - } - } catch (error) { - // D1 remains the durable source of truth if the optional analytics index is unavailable. - console.error( - JSON.stringify({ - message: "failed to index diagnostics event batch", - batchId: payload.batchId, - error: error instanceof Error ? error.message : String(error), - }), - ); - } - } - return { - id: payload.batchId, - duplicate, - stored: duplicate ? 0 : payload.events.length, - }; -} - -export async function storeCrash( - payload: NormalizedCrashPayload, - env: DiagnosticsEnv, -): Promise { - const database = env.DB.withSession("first-primary"); - const existing = await database - .prepare("SELECT fingerprint FROM crashes WHERE id = ?") - .bind(payload.id) - .first<{ fingerprint: string }>(); - if (existing) { - return { id: payload.id, duplicate: true, stored: 0, fingerprint: existing.fingerprint }; - } - - const fingerprint = await crashFingerprint(payload.exceptionType, payload.stackTrace); - const stackKey = payload.stackTrace - ? `crashes/${payload.id}/${crypto.randomUUID()}/stack.txt` - : null; - - if (stackKey) { - await env.BLOBS.put(stackKey, payload.stackTrace, { - httpMetadata: { contentType: "text/plain; charset=utf-8" }, - customMetadata: { installId: payload.installId, fingerprint }, - }); - } - - try { - const result = await database - .prepare( - `INSERT INTO crashes ( - id, received_at, occurred_at, install_id, app_version, platform, - exception_type, exception_message, fingerprint, diagnostics_enabled, - stack_r2_key, breadcrumbs_json, schema_version - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) - ON CONFLICT(id) DO NOTHING`, - ) - .bind( - payload.id, - Date.now(), - payload.occurredAt, - payload.installId, - payload.appVersion, - payload.platform, - payload.exceptionType, - payload.exceptionMessage, - fingerprint, - payload.diagnosticsEnabledAtCapture ? 1 : 0, - stackKey, - JSON.stringify(payload.breadcrumbs), - payload.schemaVersion, - ) - .run(); - const duplicate = result.meta.changes === 0; - if (duplicate) { - const stored = await database - .prepare("SELECT fingerprint FROM crashes WHERE id = ?") - .bind(payload.id) - .first<{ fingerprint: string }>(); - if (!stored) throw new Error("duplicate crash row was not readable"); - if (stackKey) await deleteAttemptBlob(env, stackKey); - return { id: payload.id, duplicate: true, stored: 0, fingerprint: stored.fingerprint }; - } - return { id: payload.id, duplicate: false, stored: 1, fingerprint }; - } catch (error) { - if (stackKey) { - await deleteAttemptBlob(env, stackKey); - } - throw error; - } -} - export async function storeBug( payload: NormalizedBugPayload, env: DiagnosticsEnv, @@ -201,26 +72,22 @@ export async function storeBug( export async function runRetention(env: DiagnosticsEnv): Promise { const retentionDays = boundedPositiveInt(env.RETENTION_DAYS, 90, 1, 3_650); const cutoff = Date.now() - retentionDays * 86_400_000; - // Eight full passes plus the backlog check use at most 43 of D1's 50 queries per invocation. + // Eight passes plus the backlog check stay well within D1's 50 queries per invocation. for (let pass = 0; pass < 8; pass += 1) { const hasFullBatch = await runRetentionPass(env, cutoff); if (!hasFullBatch) return; } - const [events, crashes, bugs] = await env.DB.batch<{ count: number }>([ - env.DB.prepare("SELECT COUNT(*) AS count FROM event_batches WHERE received_at < ?").bind( - cutoff, - ), - env.DB.prepare("SELECT COUNT(*) AS count FROM crashes WHERE received_at < ?").bind(cutoff), - env.DB.prepare("SELECT COUNT(*) AS count FROM bugs WHERE received_at < ?").bind(cutoff), - ]); + const bugs = await env.DB.prepare( + "SELECT COUNT(*) AS count FROM bugs WHERE received_at < ?", + ) + .bind(cutoff) + .first<{ count: number }>(); console.warn( JSON.stringify({ message: "diagnostics retention reached its per-run pass limit", cutoff, backlog: { - eventBatches: events.results[0]?.count ?? 0, - crashes: crashes.results[0]?.count ?? 0, - bugs: bugs.results[0]?.count ?? 0, + bugs: bugs?.count ?? 0, }, }), ); @@ -228,28 +95,17 @@ export async function runRetention(env: DiagnosticsEnv): Promise { async function runRetentionPass(env: DiagnosticsEnv, cutoff: number): Promise { const reportBatchSize = 900; - const eventBatchSize = 1_000; - const [crashes, bugs] = await Promise.all([ - expiredBlobRows(env.DB, "crashes", "stack_r2_key", cutoff, reportBatchSize), - expiredBlobRows(env.DB, "bugs", "logs_r2_key", cutoff, reportBatchSize), - ]); + const bugs = await expiredBlobRows(env.DB, "bugs", "logs_r2_key", cutoff, reportBatchSize); - const blobKeys = [...crashes, ...bugs] + const blobKeys = bugs .map((row) => row.blobKey) .filter((key): key is string => key !== null); for (let offset = 0; offset < blobKeys.length; offset += 1_000) { await env.BLOBS.delete(blobKeys.slice(offset, offset + 1_000)); } - const statements = [retentionStatement(env.DB, "event_batches", cutoff, eventBatchSize)]; - if (crashes.length > 0) statements.push(deleteRowsById(env.DB, "crashes", crashes)); - if (bugs.length > 0) statements.push(deleteRowsById(env.DB, "bugs", bugs)); - const [eventsResult] = await env.DB.batch(statements); - return ( - eventsResult.meta.changes === eventBatchSize || - crashes.length === reportBatchSize || - bugs.length === reportBatchSize - ); + if (bugs.length > 0) await deleteRowsById(env.DB, "bugs", bugs).run(); + return bugs.length === reportBatchSize; } interface ExpiredBlobRow { @@ -259,8 +115,8 @@ interface ExpiredBlobRow { async function expiredBlobRows( database: D1Database, - table: "crashes" | "bugs", - column: "stack_r2_key" | "logs_r2_key", + table: "bugs", + column: "logs_r2_key", cutoff: number, batchSize: number, ): Promise { @@ -277,25 +133,9 @@ async function expiredBlobRows( return result.results; } -function retentionStatement( - database: D1Database, - table: "event_batches" | "crashes" | "bugs", - cutoff: number, - batchSize: number, -): D1PreparedStatement { - return database - .prepare( - `DELETE FROM ${table} - WHERE rowid IN ( - SELECT rowid FROM ${table} WHERE received_at < ? ORDER BY received_at LIMIT ? - )`, - ) - .bind(cutoff, batchSize); -} - function deleteRowsById( database: D1Database, - table: "crashes" | "bugs", + table: "bugs", rows: ExpiredBlobRow[], ): D1PreparedStatement { return database @@ -303,18 +143,6 @@ function deleteRowsById( .bind(JSON.stringify(rows.map((row) => row.id))); } -async function crashFingerprint(exceptionType: string, stackTrace: string): Promise { - const topFrames = stackTrace - .split("\n") - .map((line) => line.trim()) - .filter(Boolean) - .slice(0, 4) - .join("\n"); - const bytes = new TextEncoder().encode(`${exceptionType}\n${topFrames}`); - const digest = new Uint8Array(await crypto.subtle.digest("SHA-256", bytes)); - return Array.from(digest, (byte) => byte.toString(16).padStart(2, "0")).join(""); -} - async function deleteAttemptBlob(env: DiagnosticsEnv, key: string): Promise { try { await env.BLOBS.delete(key); diff --git a/services/diagnostics-api/test/input.test.ts b/services/diagnostics-api/test/input.test.ts index 6390cfe..5d343f2 100644 --- a/services/diagnostics-api/test/input.test.ts +++ b/services/diagnostics-api/test/input.test.ts @@ -4,8 +4,6 @@ import { MAX_DEVICE_JSON_BYTES, MAX_LOG_BYTES, normalizeBug, - normalizeCrash, - normalizeEvents, readJsonObject, } from "../src/input"; @@ -57,7 +55,7 @@ describe("readJsonObject", () => { }); it("requires application/json with a UTF-8 charset", async () => { - const missing = new Request("https://example.test/v1/events", { + const missing = new Request("https://example.test/v1/bugs", { method: "POST", body: "{}", }); @@ -108,18 +106,6 @@ describe("readJsonObject", () => { describe("normalizers", () => { it("preserves false booleans and rejects their string representation", () => { - const crash = crashPayload(false); - const normalizedCrash = normalizeCrash(crash); - expect(normalizedCrash.ok).toBe(true); - if (normalizedCrash.ok) { - expect(normalizedCrash.value.diagnosticsEnabledAtCapture).toBe(false); - } - expect(normalizeCrash(crashPayload("false"))).toEqual({ - ok: false, - status: 400, - error: "invalid_diagnostics_enabled", - }); - const bug = bugPayload({ include_logs: false, logs: "discard me" }); const normalizedBug = normalizeBug(bug); expect(normalizedBug.ok).toBe(true); @@ -171,42 +157,31 @@ describe("normalizers", () => { }); it("requires stable report IDs and validates supplied IDs and schema versions", () => { - const result = normalizeEvents({ - events: [{ name: "opened", ts: 1, schema_version: 1 }], - }); - expect(result).toEqual({ ok: false, status: 400, error: "invalid_batch_id" }); - const legacyInstall = normalizeEvents({ - batch_id: ID, - install_id: "legacy-test-install", - events: [{ name: "opened", ts: 1 }], - }); - expect(legacyInstall.ok && legacyInstall.value.installId).toBe("legacy-test-install"); - const missingInstall = normalizeEvents({ - batch_id: ID, - events: [{ name: "opened", ts: 1 }], - }); - expect(missingInstall.ok && missingInstall.value.installId).toBe("unknown"); - expect( - normalizeEvents({ - batch_id: ID, - install_id: "bad\u0000install", - events: [{ name: "opened", ts: 1 }], - }), - ).toEqual({ ok: false, status: 400, error: "invalid_install_id" }); + const missingId = normalizeBug(bugPayload({ id: undefined })); + expect(missingId).toEqual({ ok: false, status: 400, error: "invalid_id" }); - expect( - normalizeEvents({ - batch_id: "not-a-uuid", - events: [{ name: "opened", timestamp_millis: 1 }], - }), - ).toEqual({ ok: false, status: 400, error: "invalid_batch_id" }); - expect( - normalizeEvents({ - batch_id: ID, - install_id: INSTALL_ID, - events: [{ name: "opened", timestamp_millis: 1, schema_version: 2 }], - }), - ).toEqual({ ok: false, status: 400, error: "unsupported_schema_version" }); + const legacyInstall = normalizeBug(bugPayload({ install_id: "legacy-test-install" })); + expect(legacyInstall.ok && legacyInstall.value.installId).toBe("legacy-test-install"); + + const missingInstall = normalizeBug(bugPayload({ install_id: undefined })); + expect(missingInstall.ok && missingInstall.value.installId).toBe("unknown"); + + expect(normalizeBug(bugPayload({ install_id: "bad\u0000install" }))).toEqual({ + ok: false, + status: 400, + error: "invalid_install_id", + }); + + expect(normalizeBug(bugPayload({ id: "not-a-uuid" }))).toEqual({ + ok: false, + status: 400, + error: "invalid_id", + }); + expect(normalizeBug(bugPayload({ schema_version: 2 }))).toEqual({ + ok: false, + status: 400, + error: "unsupported_schema_version", + }); }); }); @@ -215,7 +190,7 @@ function chunkedJsonRequest( contentType = "application/json; charset=utf-8", contentLength?: string, ): Request { - return new Request("https://example.test/v1/events", { + return new Request("https://example.test/v1/bugs", { method: "POST", headers: { "content-type": contentType, @@ -230,24 +205,8 @@ function chunkedJsonRequest( }); } -function crashPayload(diagnosticsEnabled: unknown): Record { - return { - id: ID, - install_id: INSTALL_ID, - app_version: "1.0", - platform: "test", - exception_type: "ExampleError", - exception_message: "message", - stack_trace: "stack", - occurred_at: 1, - diagnostics_enabled: diagnosticsEnabled, - schema_version: 1, - breadcrumbs: [], - }; -} - function bugPayload(overrides: Record = {}): Record { - return { + const payload: Record = { id: ID, install_id: INSTALL_ID, app_version: "1.0", @@ -263,4 +222,9 @@ function bugPayload(overrides: Record = {}): Record { expect(unknown.status).toBe(404); const unauthorized = await exports.default.fetch( - jsonRequest("/v1/events", eventPayload(uuid(1)), "wrong-key"), + jsonRequest("/v1/bugs", bugPayload(uuid(1), "logs"), "wrong-key"), ); expect(unauthorized.status).toBe(401); expect(await unauthorized.json()).toEqual({ error: "unauthorized" }); @@ -44,7 +34,7 @@ describe("diagnostics Worker", () => { ); const preflight = await exports.default.fetch( - new Request("https://diagnostics.test/v1/events", { method: "OPTIONS" }), + new Request("https://diagnostics.test/v1/bugs", { method: "OPTIONS" }), ); expect(preflight.status).toBe(204); expect(preflight.headers.get("access-control-allow-origin")).toBeNull(); @@ -68,7 +58,7 @@ describe("diagnostics Worker", () => { const context = createExecutionContext(); const response = await worker.fetch( - jsonRequest("/v1/events", eventPayload(uuid(3)), "wrong-key", "198.51.100.3"), + jsonRequest("/v1/bugs", bugPayload(uuid(3), "logs"), "wrong-key", "198.51.100.3"), limitedEnv, context, ); @@ -80,7 +70,7 @@ describe("diagnostics Worker", () => { }); it("returns structured errors for invalid bodies and asynchronous storage failures", async () => { - const invalid = await exports.default.fetch(jsonRequest("/v1/events", null)); + const invalid = await exports.default.fetch(jsonRequest("/v1/bugs", null)); expect(invalid.status).toBe(400); expect(await invalid.json()).toEqual({ error: "invalid_body" }); @@ -92,6 +82,8 @@ describe("diagnostics Worker", () => { }; const rejectingDatabase = { prepare: () => statement, + batch: async () => Promise.reject(rejection), + withSession: () => ({ prepare: () => statement }), } as unknown as D1Database; const rejectingEnv: DiagnosticsEnv = { ...env, DB: rejectingDatabase }; @@ -106,7 +98,7 @@ describe("diagnostics Worker", () => { const ingestContext = createExecutionContext(); const failedIngest = await worker.fetch( - jsonRequest("/v1/events", eventPayload(uuid(2)), env.INGEST_KEY, "198.51.100.2"), + jsonRequest("/v1/bugs", bugPayload(uuid(2), "logs"), env.INGEST_KEY, "198.51.100.2"), rejectingEnv, ingestContext, ); @@ -114,101 +106,6 @@ describe("diagnostics Worker", () => { expect(await failedIngest.json()).toEqual({ error: "internal" }); }); - it("deduplicates event batches using the client batch ID", async () => { - const id = uuid(10); - const first = await exports.default.fetch(jsonRequest("/v1/events", eventPayload(id))); - const second = await exports.default.fetch(jsonRequest("/v1/events", eventPayload(id))); - - expect(first.status).toBe(202); - expect(await first.json()).toMatchObject({ - ok: true, - id, - stored: 1, - duplicate: false, - }); - expect(second.status).toBe(202); - expect(await second.json()).toMatchObject({ - ok: true, - id, - stored: 0, - duplicate: true, - }); - - const row = await env.DB.prepare( - "SELECT event_count AS eventCount, payload_json AS payloadJson FROM event_batches WHERE id = ?", - ) - .bind(id) - .first<{ eventCount: number; payloadJson: string }>(); - expect(row?.eventCount).toBe(1); - expect(JSON.parse(row?.payloadJson ?? "null")).toEqual([ - { - name: "app_open", - timestampMillis: 1, - properties: { screen: "home" }, - schemaVersion: 1, - }, - ]); - }); - - it("keeps D1 idempotency when the optional analytics index is enabled", async () => { - const points: AnalyticsEngineDataPoint[] = []; - const analytics = { - writeDataPoint: (point: AnalyticsEngineDataPoint) => points.push(point), - } as AnalyticsEngineDataset; - const analyticsEnv: DiagnosticsEnv = { ...env, AE: analytics }; - const payload: NormalizedEventsPayload = { - batchId: uuid(11), - installId: INSTALL_ID, - appVersion: "1.0", - platform: "test", - events: [ - { - name: "indexed", - timestampMillis: 1, - properties: {}, - schemaVersion: 1, - }, - ], - }; - - expect(await storeEvents(payload, analyticsEnv)).toMatchObject({ duplicate: false, stored: 1 }); - expect(await storeEvents(payload, analyticsEnv)).toMatchObject({ duplicate: true, stored: 0 }); - expect(points).toHaveLength(1); - }); - - it("keeps the accepted crash blob when a duplicate request arrives", async () => { - const id = uuid(20); - const first = await exports.default.fetch( - jsonRequest("/v1/crashes", crashPayload(id, "first stack")), - ); - const second = await exports.default.fetch( - jsonRequest("/v1/crashes", crashPayload(id, "second stack")), - ); - - expect(first.status).toBe(202); - const firstBody = await first.json<{ fingerprint: string }>(); - expect(firstBody).toMatchObject({ ok: true, id, duplicate: false }); - expect(second.status).toBe(202); - const secondBody = await second.json<{ fingerprint: string }>(); - expect(secondBody).toMatchObject({ ok: true, id, duplicate: true }); - - const row = await env.DB.prepare( - `SELECT stack_r2_key AS stackKey, breadcrumbs_json AS breadcrumbsJson, - fingerprint - FROM crashes WHERE id = ?`, - ) - .bind(id) - .first<{ stackKey: string; breadcrumbsJson: string; fingerprint: string }>(); - expect(row?.stackKey).toMatch(new RegExp(`^crashes/${id}/[0-9a-f-]+/stack\\.txt$`)); - expect(firstBody.fingerprint).toBe(row?.fingerprint); - expect(secondBody.fingerprint).toBe(row?.fingerprint); - expect(JSON.parse(row?.breadcrumbsJson ?? "null")).toEqual([]); - expect(await (await env.BLOBS.get(row?.stackKey ?? "missing"))?.text()).toBe("first stack"); - - const objects = await env.BLOBS.list({ prefix: `crashes/${id}/` }); - expect(objects.objects.map((object) => object.key)).toEqual([row?.stackKey]); - }); - it("stores bug metadata as JSON and cleans the duplicate upload attempt", async () => { const id = uuid(30); const payload = bugPayload(id, "first logs"); @@ -252,9 +149,7 @@ describe("diagnostics Worker", () => { }); it("acknowledges known report IDs without touching an unavailable blob store", async () => { - const crash = normalizedCrash(uuid(31), "accepted stack"); const bug = normalizedBug(uuid(32), "accepted logs"); - const firstCrash = await storeCrash(crash, env); await storeBug(bug, env); let blobWrites = 0; const unavailableBlobs = { @@ -265,14 +160,6 @@ describe("diagnostics Worker", () => { } as unknown as R2Bucket; const unavailableEnv: DiagnosticsEnv = { ...env, BLOBS: unavailableBlobs }; - await expect( - storeCrash({ ...crash, stackTrace: "retry stack" }, unavailableEnv), - ).resolves.toEqual({ - id: crash.id, - duplicate: true, - stored: 0, - fingerprint: firstCrash.fingerprint, - }); await expect( storeBug({ ...bug, logs: "retry logs" }, unavailableEnv), ).resolves.toEqual({ id: bug.id, duplicate: true, stored: 0 }); @@ -286,59 +173,27 @@ describe("diagnostics Worker", () => { async () => Promise.reject(rejection), ); const rejectingEnv: DiagnosticsEnv = { ...env, DB: rejectingDatabase }; - const crashId = uuid(33); const bugId = uuid(34); - const crashResponse = await worker.fetch( - jsonRequest("/v1/crashes", crashPayload(crashId, "orphan candidate"), env.INGEST_KEY, "198.51.100.33"), - rejectingEnv, - createExecutionContext(), - ); const bugResponse = await worker.fetch( jsonRequest("/v1/bugs", bugPayload(bugId, "orphan candidate"), env.INGEST_KEY, "198.51.100.34"), rejectingEnv, createExecutionContext(), ); - expect(crashResponse.status).toBe(500); - expect(await crashResponse.json()).toEqual({ error: "internal" }); expect(bugResponse.status).toBe(500); expect(await bugResponse.json()).toEqual({ error: "internal" }); - expect((await env.BLOBS.list({ prefix: `crashes/${crashId}/` })).objects).toEqual([]); expect((await env.BLOBS.list({ prefix: `bugs/${bugId}/` })).objects).toEqual([]); }); it("removes expired rows and their exact R2 objects while preserving current data", async () => { - const oldEventId = uuid(40); - const oldCrashId = uuid(41); const oldBugId = uuid(42); - const currentEventId = uuid(43); - const oldCrashKey = `crashes/${oldCrashId}/retention/stack.txt`; + const currentBugId = uuid(43); const oldBugKey = `bugs/${oldBugId}/retention/logs.txt`; const oldReceivedAt = Date.now() - 100 * 86_400_000; - await Promise.all([ - env.BLOBS.put(oldCrashKey, "expired crash"), - env.BLOBS.put(oldBugKey, "expired logs"), - ]); + await env.BLOBS.put(oldBugKey, "expired logs"); await env.DB.batch([ - env.DB.prepare( - `INSERT INTO event_batches - (id, received_at, install_id, app_version, platform, event_count, payload_json) - VALUES (?, ?, ?, '', '', 1, '[]')`, - ).bind(oldEventId, oldReceivedAt, INSTALL_ID), - env.DB.prepare( - `INSERT INTO event_batches - (id, received_at, install_id, app_version, platform, event_count, payload_json) - VALUES (?, ?, ?, '', '', 1, '[]')`, - ).bind(currentEventId, Date.now(), INSTALL_ID), - env.DB.prepare( - `INSERT INTO crashes - (id, received_at, occurred_at, install_id, app_version, platform, - exception_type, exception_message, fingerprint, diagnostics_enabled, - stack_r2_key, breadcrumbs_json, schema_version) - VALUES (?, ?, ?, ?, '', '', 'Error', '', 'fingerprint', 1, ?, '[]', 1)`, - ).bind(oldCrashId, oldReceivedAt, oldReceivedAt, INSTALL_ID, oldCrashKey), env.DB.prepare( `INSERT INTO bugs (id, received_at, occurred_at, install_id, app_version, platform, @@ -346,28 +201,28 @@ describe("diagnostics Worker", () => { device_json, breadcrumbs_json, status, schema_version) VALUES (?, ?, ?, ?, '', '', 'failed', 'worked', '', '', ?, '{}', '[]', 'open', 1)`, ).bind(oldBugId, oldReceivedAt, oldReceivedAt, INSTALL_ID, oldBugKey), + env.DB.prepare( + `INSERT INTO bugs + (id, received_at, occurred_at, install_id, app_version, platform, + what_happened, expected, steps, contact, logs_r2_key, + device_json, breadcrumbs_json, status, schema_version) + VALUES (?, ?, ?, ?, '', '', 'failed', 'worked', '', '', NULL, '{}', '[]', 'open', 1)`, + ).bind(currentBugId, Date.now(), Date.now(), INSTALL_ID), ]); await runRetention(env); - for (const [table, id] of [ - ["event_batches", oldEventId], - ["crashes", oldCrashId], - ["bugs", oldBugId], - ] as const) { - const row = await env.DB.prepare(`SELECT id FROM ${table} WHERE id = ?`).bind(id).first(); - expect(row).toBeNull(); - } - expect(await env.BLOBS.head(oldCrashKey)).toBeNull(); + expect( + await env.DB.prepare("SELECT id FROM bugs WHERE id = ?").bind(oldBugId).first(), + ).toBeNull(); expect(await env.BLOBS.head(oldBugKey)).toBeNull(); expect( - await env.DB.prepare("SELECT id FROM event_batches WHERE id = ?").bind(currentEventId).first(), + await env.DB.prepare("SELECT id FROM bugs WHERE id = ?").bind(currentBugId).first(), ).not.toBeNull(); }); it("bounds a full retention run below the D1 per-invocation query limit", async () => { let queryCount = 0; - let batchCalls = 0; const blobDeleteBatchSizes: number[] = []; const rows = Array.from({ length: 900 }, (_, index) => ({ id: `expired-${index}`, @@ -381,17 +236,17 @@ describe("diagnostics Worker", () => { queryCount += 1; return d1Result(rows, 0); }, + run: async () => { + queryCount += 1; + return d1Result([], rows.length); + }, + first: async () => { + queryCount += 1; + return { count: rows.length }; + }, }; return statement; }, - batch: async (statements: D1PreparedStatement[]) => { - batchCalls += 1; - queryCount += statements.length; - if (batchCalls === 9) { - return statements.map(() => d1Result([{ count: 1 }], 0)); - } - return statements.map((_, index) => d1Result([], index === 0 ? 1_000 : 900)); - }, } as unknown as D1Database; const warning = vi.spyOn(console, "warn").mockImplementation(() => undefined); const blobs = { @@ -406,9 +261,10 @@ describe("diagnostics Worker", () => { warning.mockRestore(); } - expect(queryCount).toBe(43); - expect(blobDeleteBatchSizes).toHaveLength(16); - expect(Math.max(...blobDeleteBatchSizes)).toBe(1_000); + // Eight passes (one SELECT + one DELETE each) plus the final backlog SELECT. + expect(queryCount).toBe(17); + expect(blobDeleteBatchSizes).toHaveLength(8); + expect(Math.max(...blobDeleteBatchSizes)).toBe(900); }); it("converges an expired report backlog across bounded retention runs", async () => { @@ -424,13 +280,13 @@ describe("diagnostics Worker", () => { CROSS JOIN digits AS ones WHERE thousands.value * 1000 + hundreds.value * 100 + tens.value * 10 + ones.value < 7201 ) - INSERT INTO crashes ( + INSERT INTO bugs ( id, received_at, occurred_at, install_id, app_version, platform, - exception_type, exception_message, fingerprint, diagnostics_enabled, - stack_r2_key, breadcrumbs_json, schema_version + what_happened, expected, steps, contact, logs_r2_key, + device_json, breadcrumbs_json, status, schema_version ) SELECT 'retention-backlog-' || printf('%04d', value), ?, ?, ?, '', '', - 'Error', '', 'fingerprint-' || value, 0, NULL, '[]', 1 + 'failed', 'worked', '', '', NULL, '{}', '[]', 'open', 1 FROM sequence`, ) .bind(oldReceivedAt, oldReceivedAt, INSTALL_ID) @@ -440,13 +296,13 @@ describe("diagnostics Worker", () => { try { await runRetention(env); const afterFirstRun = await env.DB.prepare( - "SELECT COUNT(*) AS count FROM crashes WHERE id LIKE 'retention-backlog-%'", + "SELECT COUNT(*) AS count FROM bugs WHERE id LIKE 'retention-backlog-%'", ).first<{ count: number }>(); expect(afterFirstRun?.count).toBe(1); await runRetention(env); const afterSecondRun = await env.DB.prepare( - "SELECT COUNT(*) AS count FROM crashes WHERE id LIKE 'retention-backlog-%'", + "SELECT COUNT(*) AS count FROM bugs WHERE id LIKE 'retention-backlog-%'", ).first<{ count: number }>(); expect(afterSecondRun?.count).toBe(0); } finally { @@ -482,39 +338,6 @@ function healthRequest(key = env.INGEST_KEY, source = "198.51.100.1"): Request { }); } -function eventPayload(batchId: string): Record { - return { - batchId, - installId: INSTALL_ID, - appVersion: "1.0", - platform: "test", - events: [ - { - name: "app_open", - timestampMillis: 1, - properties: { screen: "home" }, - schemaVersion: 1, - }, - ], - }; -} - -function crashPayload(id: string, stackTrace: string): Record { - return { - id, - installId: INSTALL_ID, - appVersion: "1.0", - platform: "test", - exceptionType: "TestError", - exceptionMessage: "failed", - stackTrace, - timestampMillis: 2, - diagnosticsEnabledAtCapture: true, - breadcrumbs: [], - schemaVersion: 1, - }; -} - function bugPayload(id: string, logs: string): Record { return { id, @@ -540,22 +363,6 @@ function bugPayload(id: string, logs: string): Record { }; } -function normalizedCrash(id: string, stackTrace: string): NormalizedCrashPayload { - return { - id, - installId: INSTALL_ID, - appVersion: "1.0", - platform: "test", - exceptionType: "TestError", - exceptionMessage: "failed", - stackTrace, - occurredAt: 2, - diagnosticsEnabledAtCapture: true, - breadcrumbs: [], - schemaVersion: 1, - }; -} - function normalizedBug(id: string, logs: string): NormalizedBugPayload { return { id, diff --git a/services/diagnostics-api/worker-configuration.d.ts b/services/diagnostics-api/worker-configuration.d.ts index 0481e2f..0d0ca95 100644 --- a/services/diagnostics-api/worker-configuration.d.ts +++ b/services/diagnostics-api/worker-configuration.d.ts @@ -1,5 +1,5 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types` (hash: e2953336d40e96c4b5125a5de01f7cac) +// Generated by Wrangler by running `wrangler types` (hash: 9c27cfd9219ba0d4efc153db78cad4c3) // Runtime types generated with workerd@1.20260708.1 2026-07-14 nodejs_compat interface __BaseEnv_Env { BLOBS: R2Bucket; @@ -7,7 +7,6 @@ interface __BaseEnv_Env { INSTALL_RATE_LIMITER: RateLimit; SOURCE_RATE_LIMITER: RateLimit; MAX_BODY_BYTES: "262144"; - MAX_EVENTS_PER_BATCH: "50"; RETENTION_DAYS: "90"; } declare namespace Cloudflare { @@ -21,7 +20,7 @@ type StringifyValues> = { [Binding in keyof EnvType]: EnvType[Binding] extends string ? EnvType[Binding] : string; }; declare namespace NodeJS { - interface ProcessEnv extends StringifyValues> {} + interface ProcessEnv extends StringifyValues> {} } // Begin runtime types diff --git a/services/diagnostics-api/wrangler.jsonc b/services/diagnostics-api/wrangler.jsonc index 665477c..1f2a09d 100644 --- a/services/diagnostics-api/wrangler.jsonc +++ b/services/diagnostics-api/wrangler.jsonc @@ -47,13 +47,8 @@ }, }, ], - // Optional: bind Analytics Engine as `AE` when event volume justifies it. - // "analytics_engine_datasets": [ - // { "binding": "AE", "dataset": "vnidrop_events" }, - // ], "vars": { "MAX_BODY_BYTES": "262144", - "MAX_EVENTS_PER_BATCH": "50", "RETENTION_DAYS": "90", }, }