refactor(diagnostics-api): drop telemetry and crash ingestion, keep bug reports

Remove the /v1/events and /v1/crashes routes, their normalizers and storage
paths, and simplify retention to the bugs table. Add a migration dropping the
now-unused event_batches and crashes tables, and regenerate worker types.
This commit is contained in:
2026-08-02 10:03:12 +02:00
parent b68d338097
commit 3e18378610
9 changed files with 122 additions and 730 deletions

View File

@@ -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<Response> {
@@ -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<Respon
}
try {
await env.DB.batch([
env.DB.prepare(
"SELECT id, received_at, install_id, payload_json FROM event_batches LIMIT 1",
),
env.DB.prepare(
"SELECT id, occurred_at, stack_r2_key, breadcrumbs_json FROM crashes LIMIT 1",
),
env.DB.prepare(
"SELECT id, occurred_at, logs_r2_key, device_json FROM bugs LIMIT 1",
),
@@ -216,8 +170,8 @@ async function installRateLimited(
return !result.success;
}
function isIngestPath(path: string): path is "/v1/events" | "/v1/crashes" | "/v1/bugs" {
return path === "/v1/events" || path === "/v1/crashes" || path === "/v1/bugs";
function isIngestPath(path: string): path is "/v1/bugs" {
return path === "/v1/bugs";
}
async function timingSafeEqual(provided: string, expected: string): Promise<boolean> {

View File

@@ -10,13 +10,6 @@ export type InputResult<T> = { ok: true; value: T } | InputFailure;
export type NormalizedProperties = Record<string, string>;
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<NormalizedEventsPayload> {
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<NormalizedCrashPayload> {
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<NormalizedBugPayload> {
if (!isPlainObject(body)) return failure(400, "invalid_body");
@@ -341,24 +199,6 @@ export function normalizeBug(body: JsonObject): InputResult<NormalizedBugPayload
});
}
function normalizeEvent(raw: unknown): InputResult<NormalizedEvent> {
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<NormalizedBreadcrumb[]> {
if (raw === MISSING) return success([]);
if (!Array.isArray(raw)) return failure(400, "invalid_breadcrumbs");

View File

@@ -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<StoreResult> {
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<StoreResult & { fingerprint: string }> {
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<void> {
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<void> {
async function runRetentionPass(env: DiagnosticsEnv, cutoff: number): Promise<boolean> {
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<ExpiredBlobRow[]> {
@@ -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<string> {
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<void> {
try {
await env.BLOBS.delete(key);