| 1 | // Ingest + dashboard for desktop crash/feedback/performance reports and the |
| 2 | // anonymous launch ping. Frontend reports are user-initiated; native fatal and |
| 3 | // lifecycle reports are sent on the next launch under the same opt-out desktop |
| 4 | // telemetry gate as pings. |
| 5 | import { z } from "zod"; |
| 6 | import type { Env } from "./env"; |
| 7 | import { html, redirect } from "./shell"; |
| 8 | import { renderGroup, renderStats, type Group, type StatsModule } from "./stats"; |
| 9 | import { renderAccount } from "./auth_pages"; |
| 10 | import { renderUsers, renderAudit, type UserRow, type AuditRow } from "./admin"; |
| 11 | import { |
| 12 | atLeast, |
| 13 | currentUser, |
| 14 | loginUrl, |
| 15 | logAction, |
| 16 | sameOrigin, |
| 17 | sharedLogout, |
| 18 | type Role, |
| 19 | type User, |
| 20 | } from "./auth"; |
| 21 | import registryApp from "./registry/app"; |
| 22 | import type { Bindings as RegistryBindings } from "./registry/env"; |
| 23 | import { PackageRepo } from "./registry/db/packages"; |
| 24 | import { EventRepo } from "./registry/db/events"; |
| 25 | import { renderCommunity } from "./community"; |
| 26 | import { |
| 27 | cliReleaseChannel, |
| 28 | desktopReleaseChannel, |
| 29 | handleCLIRelease, |
| 30 | handleDesktopReleaseManifest, |
| 31 | handleReleaseGatewayRequest, |
| 32 | } from "./desktop_release"; |
| 33 | |
| 34 | const MAX_BODY_BYTES = 96 * 1024; |
| 35 | const LATEST_SAMPLES_PER_GROUP = 5; |
| 36 | const DEVELOPMENT_FINGERPRINT_PREFIX = "dev:"; |
| 37 | const GROUP_PATH_RE = /^\/stats\/group\/((?:dev:)?[0-9a-f]{64})$/; |
| 38 | |
| 39 | const Device = z |
| 40 | .object({ |
| 41 | osVersion: z.string().max(128), |
| 42 | cpu: z.string().max(128), |
| 43 | cores: z.number().int().min(0).max(4096), |
| 44 | ramGb: z.number().min(0).max(65536), |
| 45 | }) |
| 46 | .partial(); |
| 47 | |
| 48 | const Report = z.object({ |
| 49 | kind: z.enum(["crash", "exception", "feedback", "performance", "bot"]), |
| 50 | version: z.string().min(1).max(64), |
| 51 | os: z.string().min(1).max(32), |
| 52 | arch: z.string().min(1).max(32), |
| 53 | message: z.string().min(1).max(16 * 1024), |
| 54 | device: Device.optional(), |
| 55 | schemaVersion: z.number().int().min(1).max(10).optional(), |
| 56 | source: z.string().trim().min(1).max(32).regex(/^[a-z0-9_.-]+$/).optional(), |
| 57 | label: z.string().max(64).optional(), |
| 58 | errorType: z.string().max(128).optional(), |
| 59 | errorMessage: z.string().max(4 * 1024).optional(), |
| 60 | stack: z.string().max(16 * 1024).optional(), |
| 61 | componentStack: z.string().max(16 * 1024).optional(), |
| 62 | topFrame: z.string().max(300).optional(), |
| 63 | fingerprintHint: z.string().max(300).optional(), |
| 64 | buildCommit: z.string().max(64).optional(), |
| 65 | channel: z.string().max(32).optional(), |
| 66 | language: z.string().max(64).optional(), |
| 67 | view: z.string().max(200).optional(), |
| 68 | breadcrumbs: z |
| 69 | .array( |
| 70 | z.object({ |
| 71 | t: z.number().int().optional(), |
| 72 | cat: z.string().max(64).optional(), |
| 73 | msg: z.string().max(240).optional(), |
| 74 | }), |
| 75 | ) |
| 76 | .max(30) |
| 77 | .optional(), |
| 78 | occurredAt: z.string().max(64).optional(), |
| 79 | }); |
| 80 | type ReportPayload = z.infer<typeof Report>; |
| 81 | |
| 82 | const ClientSurface = z.enum(["desktop", "cli"]); |
| 83 | type ClientSurfaceName = z.infer<typeof ClientSurface>; |
| 84 | |
| 85 | type TelemetryTableNames = { |
| 86 | pings: "pings" | "cli_pings"; |
| 87 | metrics: "metrics" | "cli_metrics"; |
| 88 | metricUsers: "metric_users" | "cli_metric_users"; |
| 89 | }; |
| 90 | |
| 91 | const TELEMETRY_TABLES: Record<ClientSurfaceName, TelemetryTableNames> = { |
| 92 | desktop: { pings: "pings", metrics: "metrics", metricUsers: "metric_users" }, |
| 93 | cli: { pings: "cli_pings", metrics: "cli_metrics", metricUsers: "cli_metric_users" }, |
| 94 | }; |
| 95 | |
| 96 | export function telemetryTableNames(surface: ClientSurfaceName): TelemetryTableNames { |
| 97 | return TELEMETRY_TABLES[surface]; |
| 98 | } |
| 99 | |
| 100 | export const CLI_TELEMETRY_SCHEMA_SQL = [ |
| 101 | `CREATE TABLE IF NOT EXISTS cli_pings ( |
| 102 | date TEXT NOT NULL, |
| 103 | install_id TEXT NOT NULL, |
| 104 | version TEXT NOT NULL, |
| 105 | os TEXT NOT NULL, |
| 106 | arch TEXT NOT NULL, |
| 107 | os_version TEXT NOT NULL DEFAULT '', |
| 108 | opens INTEGER NOT NULL DEFAULT 1, |
| 109 | PRIMARY KEY (date, install_id) |
| 110 | )`, |
| 111 | `CREATE TABLE IF NOT EXISTS cli_metrics ( |
| 112 | date TEXT NOT NULL, |
| 113 | version TEXT NOT NULL, |
| 114 | os TEXT NOT NULL, |
| 115 | signal TEXT NOT NULL, |
| 116 | bucket TEXT NOT NULL, |
| 117 | count INTEGER NOT NULL DEFAULT 0, |
| 118 | PRIMARY KEY (date, version, os, signal, bucket) |
| 119 | )`, |
| 120 | `CREATE TABLE IF NOT EXISTS cli_metric_users ( |
| 121 | date TEXT NOT NULL, |
| 122 | signal TEXT NOT NULL, |
| 123 | bucket TEXT NOT NULL, |
| 124 | install_id TEXT NOT NULL, |
| 125 | version TEXT NOT NULL, |
| 126 | os TEXT NOT NULL, |
| 127 | PRIMARY KEY (date, signal, bucket, install_id) |
| 128 | )`, |
| 129 | // No secondary indexes: each primary key already leads with `date`, which is |
| 130 | // what every dashboard query filters on. See migrate-window-index-fix.sql. |
| 131 | ] as const; |
| 132 | |
| 133 | const cliTelemetrySchemaPromises = new WeakMap<object, Promise<void>>(); |
| 134 | |
| 135 | export function ensureCLITelemetrySchema(env: Pick<Env, "DB">): Promise<void> { |
| 136 | const key = env.DB as unknown as object; |
| 137 | const existing = cliTelemetrySchemaPromises.get(key); |
| 138 | if (existing) return existing; |
| 139 | const creation = env.DB |
| 140 | .batch(CLI_TELEMETRY_SCHEMA_SQL.map((sql) => env.DB.prepare(sql))) |
| 141 | .then(() => undefined) |
| 142 | .catch((err) => { |
| 143 | cliTelemetrySchemaPromises.delete(key); |
| 144 | throw err; |
| 145 | }); |
| 146 | cliTelemetrySchemaPromises.set(key, creation); |
| 147 | return creation; |
| 148 | } |
| 149 | |
| 150 | export const Ping = z.object({ |
| 151 | installId: z.string().regex(/^[0-9a-f]{32}$/), |
| 152 | version: z.string().min(1).max(64), |
| 153 | os: z.string().min(1).max(32), |
| 154 | arch: z.string().min(1).max(32), |
| 155 | osVersion: z.string().max(128).optional(), |
| 156 | surface: ClientSurface.default("desktop"), |
| 157 | }); |
| 158 | |
| 159 | // Opt-in aggregate client metrics: a per-launch snapshot of (signal, bucket) |
| 160 | // counters. The optional surface-specific random install id deduplicates DAU; |
| 161 | // there is no user content. Unknown signals are discarded before storage so |
| 162 | // older workers can accept batches from newer clients safely. |
| 163 | const METRIC_SIGNALS = [ |
| 164 | "finish_reason", |
| 165 | "empty_final", |
| 166 | "provider_error", |
| 167 | "cache_hit", |
| 168 | "tool_error", |
| 169 | "updater_error", |
| 170 | "updater_event", |
| 171 | "compaction", |
| 172 | "turns", |
| 173 | "desktop_hang", |
| 174 | "desktop_hang_age", |
| 175 | "desktop_exit", |
| 176 | "desktop_exit_phase", |
| 177 | "desktop_uptime", |
| 178 | "desktop_install", |
| 179 | "desktop_update_transition", |
| 180 | "desktop_restore", |
| 181 | "desktop_webview2_failure", |
| 182 | "cli_mode", |
| 183 | "cli_profile", |
| 184 | "cli_permission_mode", |
| 185 | "cli_session_mode", |
| 186 | "cli_turn_latency", |
| 187 | "cli_exit", |
| 188 | "recovery_failure", |
| 189 | "recovery_rule_continue", |
| 190 | "recovery_review_continue", |
| 191 | "recovery_human_prompt", |
| 192 | "recovery_human_continue", |
| 193 | "recovery_human_revise", |
| 194 | "recovery_review_error", |
| 195 | "recovery_repeat_prompt", |
| 196 | "recovery_review_latency", |
| 197 | "client_surface", |
| 198 | "client_version", |
| 199 | "settings_language", |
| 200 | "settings_desktop_layout", |
| 201 | "settings_theme", |
| 202 | "settings_theme_style", |
| 203 | "settings_close_behavior", |
| 204 | "settings_display_mode", |
| 205 | "settings_status_bar_style", |
| 206 | "settings_status_bar_items_count", |
| 207 | "settings_check_updates", |
| 208 | "settings_default_model", |
| 209 | "settings_planner_model", |
| 210 | "settings_subagent_model", |
| 211 | "settings_subagent_effort", |
| 212 | "settings_reasoning_language", |
| 213 | "settings_provider_count", |
| 214 | "settings_provider_access_count", |
| 215 | "settings_provider_access", |
| 216 | "settings_bot_enabled", |
| 217 | "settings_bot_model", |
| 218 | "settings_bot_tool_approval", |
| 219 | "settings_bot_allowlist", |
| 220 | "settings_bot_allow_all", |
| 221 | "settings_bot_qq_enabled", |
| 222 | "settings_bot_feishu_enabled", |
| 223 | "settings_bot_weixin_enabled", |
| 224 | "settings_bot_connection_count", |
| 225 | "settings_bot_connection_provider", |
| 226 | "settings_bot_connection_enabled", |
| 227 | "settings_bot_connection_status", |
| 228 | "settings_bot_connection_model", |
| 229 | "settings_bot_connection_approval", |
| 230 | ] as const; |
| 231 | |
| 232 | type MetricSignal = (typeof METRIC_SIGNALS)[number]; |
| 233 | |
| 234 | const METRIC_SIGNAL_SET: ReadonlySet<string> = new Set(METRIC_SIGNALS); |
| 235 | |
| 236 | const KnownMetricCounter = z.object({ |
| 237 | signal: z.enum(METRIC_SIGNALS), |
| 238 | bucket: z |
| 239 | .string() |
| 240 | .min(1) |
| 241 | .max(96) |
| 242 | .regex(/^[a-z0-9_]+$/), |
| 243 | count: z.number().int().min(1).max(1_000_000), |
| 244 | }); |
| 245 | |
| 246 | const UnknownMetricCounter = z |
| 247 | .object({ |
| 248 | signal: z |
| 249 | .string() |
| 250 | .min(1) |
| 251 | .max(96) |
| 252 | .refine((signal) => !METRIC_SIGNAL_SET.has(signal)), |
| 253 | }) |
| 254 | .passthrough() |
| 255 | .transform(() => null); |
| 256 | |
| 257 | export const Metrics = z.object({ |
| 258 | installId: z |
| 259 | .string() |
| 260 | .regex(/^[0-9a-f]{32}$/) |
| 261 | .optional(), |
| 262 | version: z.string().min(1).max(64), |
| 263 | os: z.string().min(1).max(32), |
| 264 | surface: ClientSurface.default("desktop"), |
| 265 | counters: z |
| 266 | .array(z.union([KnownMetricCounter, UnknownMetricCounter])) |
| 267 | .min(1) |
| 268 | .max(128) |
| 269 | .transform((counters) => |
| 270 | counters.filter( |
| 271 | (counter): counter is z.infer<typeof KnownMetricCounter> & { signal: MetricSignal } => counter !== null, |
| 272 | ), |
| 273 | ), |
| 274 | }); |
| 275 | |
| 276 | type FingerprintInput = { |
| 277 | kind: string; |
| 278 | message: string; |
| 279 | source?: string; |
| 280 | label?: string; |
| 281 | errorType?: string; |
| 282 | errorMessage?: string; |
| 283 | topFrame?: string; |
| 284 | fingerprintHint?: string; |
| 285 | }; |
| 286 | |
| 287 | export function scrubSensitiveText(input: string): string { |
| 288 | return input |
| 289 | .replace(/([A-Z]:\\Users\\)[^/\\:\s"']+/gi, "$1_") |
| 290 | .replace(/(\/(?:home|Users)\/)[^/\\:\s"']+/g, "$1_") |
| 291 | .replace(/\b[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}\b/g, "[redacted-email]") |
| 292 | .replace(/\bBearer\s+[A-Za-z0-9._~+/=-]{16,}/gi, "Bearer [redacted]") |
| 293 | .replace( |
| 294 | /\b(api[_-]?key|access[_-]?token|refresh[_-]?token|id[_-]?token|authorization|secret|password|passwd|pwd|token)\b\s*[:=]\s*(?:Bearer\s+)?['"]?[^'"\s,;]+['"]?/gi, |
| 295 | "$1=[redacted]", |
| 296 | ) |
| 297 | .replace(/\beyJ[A-Za-z0-9_-]+\.[A-Za-z0-9_-]+\.[A-Za-z0-9_-]+\b/g, "[redacted-jwt]") |
| 298 | .replace(/\b(?:sk|rk)-(?:proj-)?[A-Za-z0-9_-]{16,}\b/g, "[redacted-key]") |
| 299 | .replace(/\b[0-9a-fA-F]{32,}\b/g, "[redacted-hex]") |
| 300 | .replace(/[A-Za-z0-9+/]{40,}={0,2}/g, "[redacted-token]") |
| 301 | .replace(/\b[A-Za-z0-9_-]{48,}\b/g, "[redacted-token]"); |
| 302 | } |
| 303 | |
| 304 | function normalizeStackFrame(frame: string): string { |
| 305 | return frame |
| 306 | .replace(/[A-Za-z]:\\[^\s)('"]+/g, "<path>") |
| 307 | .replace(/\/(?:home|Users)\/[^\s)('"]+/g, "/<home>") |
| 308 | .replace(/(?:wails|https?|file):\/\/[^\s)('"]+/g, "<url>") |
| 309 | .replace(/0x[0-9a-fA-F]+/g, "<addr>") |
| 310 | .replace(/:\d+(?::\d+)?/g, ":<n>"); |
| 311 | } |
| 312 | |
| 313 | function normalizeFingerprintText(text: string): string { |
| 314 | return text |
| 315 | .replace(/[A-Za-z]:\\[^\s)('"]+/g, "<path>") |
| 316 | .replace(/(?:wails|https?|file):\/\/[^\s)('"]+/g, "<url>") |
| 317 | .replace(/0x[0-9a-fA-F]+/g, "<addr>") |
| 318 | .replace(/^build [0-9a-f]+$/gm, "build <commit>") |
| 319 | .replace(/:\d+(?::\d+)?/g, ":<n>"); |
| 320 | } |
| 321 | |
| 322 | export function normalizeForFingerprint(inputOrKind: FingerprintInput | string, legacyMessage = ""): string { |
| 323 | if (typeof inputOrKind === "string") { |
| 324 | const head = legacyMessage.split("\n").slice(0, 12).join("\n"); |
| 325 | return inputOrKind + "\n" + normalizeFingerprintText(head); |
| 326 | } |
| 327 | const input = inputOrKind; |
| 328 | const messageBasis = input.errorMessage || input.message; |
| 329 | const head = messageBasis.split("\n").slice(0, 6).join("\n"); |
| 330 | return ( |
| 331 | input.kind + |
| 332 | "\n" + |
| 333 | (input.source || "legacy") + |
| 334 | "\n" + |
| 335 | (input.label || "") + |
| 336 | "\n" + |
| 337 | (input.errorType || "") + |
| 338 | "\n" + |
| 339 | normalizeStackFrame(input.topFrame || "") + |
| 340 | "\n" + |
| 341 | (input.fingerprintHint ? `${input.fingerprintHint}\n` : "") + |
| 342 | normalizeFingerprintText(head) |
| 343 | ); |
| 344 | } |
| 345 | |
| 346 | function hasStructuredCrashFields(r: ReportPayload): boolean { |
| 347 | return Boolean( |
| 348 | r.schemaVersion || |
| 349 | r.source || |
| 350 | r.label || |
| 351 | r.errorType || |
| 352 | r.errorMessage || |
| 353 | r.stack || |
| 354 | r.componentStack || |
| 355 | r.topFrame || |
| 356 | r.fingerprintHint || |
| 357 | r.buildCommit || |
| 358 | r.channel || |
| 359 | r.language || |
| 360 | r.view || |
| 361 | r.breadcrumbs?.length || |
| 362 | r.occurredAt, |
| 363 | ); |
| 364 | } |
| 365 | |
| 366 | // One-line human summary for the dashboard list. Frontend reports are formatted |
| 367 | // "[label]\n\n<detail>", so a bare label alone is folded together with its detail. |
| 368 | export function crashTitle(message: string): string { |
| 369 | const lines = message |
| 370 | .split("\n") |
| 371 | .map((l) => l.trim()) |
| 372 | .filter(Boolean); |
| 373 | let head = lines[0] ?? ""; |
| 374 | if (/^\[[^\]]+\]$/.test(head) && lines[1]) head = `${head} ${lines[1]}`; |
| 375 | return head.slice(0, 200); |
| 376 | } |
| 377 | |
| 378 | type SeverityInput = { |
| 379 | kind: string; |
| 380 | version?: string; |
| 381 | source: string; |
| 382 | label: string; |
| 383 | errorType: string; |
| 384 | errorMessage: string; |
| 385 | topFrame: string; |
| 386 | channel?: string; |
| 387 | }; |
| 388 | |
| 389 | const RESIZE_OBSERVER_NOTICE_RE = /^ResizeObserver loop (?:limit exceeded|completed with undelivered notifications\.?)$/; |
| 390 | |
| 391 | export function isDevelopmentReport(input: SeverityInput): boolean { |
| 392 | return input.channel?.trim().toLowerCase() === "dev" || input.version?.trim().toLowerCase().startsWith("dev") === true; |
| 393 | } |
| 394 | |
| 395 | export function namespaceReportFingerprint(hash: string, development: boolean): string { |
| 396 | return development ? `${DEVELOPMENT_FINGERPRINT_PREFIX}${hash}` : hash; |
| 397 | } |
| 398 | |
| 399 | export function groupFingerprintFromPath(path: string): string | null { |
| 400 | return path.match(GROUP_PATH_RE)?.[1] ?? null; |
| 401 | } |
| 402 | |
| 403 | export function isKnownNonCrashDiagnostic(input: SeverityInput): boolean { |
| 404 | const message = input.errorMessage.trim(); |
| 405 | return ( |
| 406 | RESIZE_OBSERVER_NOTICE_RE.test(message) || |
| 407 | /Minified React error #520\b/.test(message) || |
| 408 | message.includes("additional File object is not a file on the disk") |
| 409 | ); |
| 410 | } |
| 411 | |
| 412 | export function isOpaqueScriptErrorReport(input: SeverityInput): boolean { |
| 413 | return ( |
| 414 | input.kind === "crash" && |
| 415 | input.source === "frontend.global" && |
| 416 | input.label === "window.error" && |
| 417 | input.errorType === "string" && |
| 418 | input.errorMessage.trim() === "Script error." && |
| 419 | input.topFrame.trim() === "" |
| 420 | ); |
| 421 | } |
| 422 | |
| 423 | function severityForKind(kind: string): string { |
| 424 | if (kind === "crash") return "high"; |
| 425 | if (kind === "performance") return "medium"; |
| 426 | if (kind === "bot") return "medium"; |
| 427 | if (kind === "exception") return "medium"; |
| 428 | return "low"; |
| 429 | } |
| 430 | |
| 431 | export function severityForReport(input: SeverityInput): string { |
| 432 | if (isDevelopmentReport(input) || isOpaqueScriptErrorReport(input) || isKnownNonCrashDiagnostic(input)) return "low"; |
| 433 | return severityForKind(input.kind); |
| 434 | } |
| 435 | |
| 436 | async function sha256Hex(s: string): Promise<string> { |
| 437 | const digest = await crypto.subtle.digest("SHA-256", new TextEncoder().encode(s)); |
| 438 | return [...new Uint8Array(digest)].map((b) => b.toString(16).padStart(2, "0")).join(""); |
| 439 | } |
| 440 | |
| 441 | async function readJSON(request: Request): Promise<unknown | Response> { |
| 442 | const length = Number(request.headers.get("content-length") ?? "0"); |
| 443 | if (!length || length > MAX_BODY_BYTES) return new Response("payload too large", { status: 413 }); |
| 444 | try { |
| 445 | return JSON.parse(await request.text()); |
| 446 | } catch { |
| 447 | return new Response("bad request", { status: 400 }); |
| 448 | } |
| 449 | } |
| 450 | |
| 451 | // Ingest writes fail together when the database is unhealthy (e.g. the D1 |
| 452 | // size cap: every INSERT throws while reads stay fine). Surface that as a |
| 453 | // deliberate 503 with a loud log instead of an opaque worker exception, so |
| 454 | // `wrangler tail` / observability show the root cause and clients see a |
| 455 | // retryable status. |
| 456 | function storageUnavailable(op: string, err: unknown): Response { |
| 457 | console.error(`${op}: D1 write failed`, err); |
| 458 | return new Response("storage unavailable", { status: 503 }); |
| 459 | } |
| 460 | |
| 461 | async function handleReport(request: Request, env: Env): Promise<Response> { |
| 462 | const ip = request.headers.get("cf-connecting-ip") ?? "unknown"; |
| 463 | const { success } = await env.RATE_LIMITER.limit({ key: ip }); |
| 464 | if (!success) return new Response("rate limited", { status: 429 }); |
| 465 | |
| 466 | const raw = await readJSON(request); |
| 467 | if (raw instanceof Response) return raw; |
| 468 | const parsed = Report.safeParse(raw); |
| 469 | if (!parsed.success) return new Response("bad request", { status: 400 }); |
| 470 | const r = parsed.data; |
| 471 | const message = scrubSensitiveText(r.message); |
| 472 | const errorMessage = scrubSensitiveText(r.errorMessage ?? ""); |
| 473 | const stack = scrubSensitiveText(r.stack ?? ""); |
| 474 | const componentStack = scrubSensitiveText(r.componentStack ?? ""); |
| 475 | const topFrame = scrubSensitiveText(r.topFrame ?? ""); |
| 476 | const fingerprintHint = scrubSensitiveText(r.fingerprintHint ?? ""); |
| 477 | const view = scrubSensitiveText(r.view ?? ""); |
| 478 | const breadcrumbs = (r.breadcrumbs ?? []).map((b) => ({ |
| 479 | ...b, |
| 480 | msg: b.msg ? scrubSensitiveText(b.msg) : b.msg, |
| 481 | })); |
| 482 | |
| 483 | const fingerprintBasis = hasStructuredCrashFields(r) |
| 484 | ? normalizeForFingerprint({ |
| 485 | kind: r.kind, |
| 486 | message, |
| 487 | source: r.source, |
| 488 | label: r.label, |
| 489 | errorType: r.errorType, |
| 490 | errorMessage, |
| 491 | topFrame, |
| 492 | fingerprintHint, |
| 493 | }) |
| 494 | : normalizeForFingerprint(r.kind, message); |
| 495 | const now = new Date().toISOString(); |
| 496 | const title = crashTitle(message); |
| 497 | const source = r.source ?? "legacy"; |
| 498 | const label = r.label ?? ""; |
| 499 | const errorType = r.errorType ?? ""; |
| 500 | const buildCommit = r.buildCommit ?? ""; |
| 501 | const channel = r.channel ?? ""; |
| 502 | const severityInput = { |
| 503 | kind: r.kind, |
| 504 | version: r.version, |
| 505 | source, |
| 506 | label, |
| 507 | errorType, |
| 508 | errorMessage, |
| 509 | topFrame, |
| 510 | channel, |
| 511 | }; |
| 512 | const development = isDevelopmentReport(severityInput); |
| 513 | const fingerprint = namespaceReportFingerprint(await sha256Hex(fingerprintBasis), development); |
| 514 | const severity = severityForReport(severityInput); |
| 515 | try { |
| 516 | const prior = await env.DB.prepare("SELECT status FROM groups WHERE fingerprint = ?1") |
| 517 | .bind(fingerprint) |
| 518 | .first<{ status: string }>(); |
| 519 | const regressedAt = prior?.status === "resolved" ? now : ""; |
| 520 | |
| 521 | await env.DB.prepare( |
| 522 | `INSERT INTO groups ( |
| 523 | fingerprint, kind, count, first_seen, last_seen, first_version, last_version, |
| 524 | status, title, source, label, error_type, top_frame, severity, |
| 525 | last_os, last_arch, last_build_commit, last_channel, last_sample_at, regressed_at |
| 526 | ) |
| 527 | VALUES (?1, ?2, 1, ?3, ?3, ?4, ?4, 'open', ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?3, ?15) |
| 528 | ON CONFLICT (fingerprint) DO UPDATE SET |
| 529 | count = count + 1, |
| 530 | last_seen = ?3, |
| 531 | last_version = ?4, |
| 532 | title = ?5, |
| 533 | source = ?6, |
| 534 | label = ?7, |
| 535 | error_type = ?8, |
| 536 | top_frame = ?9, |
| 537 | severity = CASE WHEN severity = 'critical' THEN severity WHEN ?10 = 'low' THEN 'low' ELSE severity END, |
| 538 | last_os = ?11, |
| 539 | last_arch = ?12, |
| 540 | last_build_commit = ?13, |
| 541 | last_channel = ?14, |
| 542 | last_sample_at = ?3, |
| 543 | status = CASE WHEN status = 'resolved' THEN 'open' ELSE status END, |
| 544 | regressed_at = CASE WHEN status = 'resolved' THEN ?3 ELSE regressed_at END`, |
| 545 | ) |
| 546 | .bind(fingerprint, r.kind, now, r.version, title, source, label, errorType, topFrame, severity, r.os, r.arch, buildCommit, channel, regressedAt) |
| 547 | .run(); |
| 548 | |
| 549 | await env.DB.prepare( |
| 550 | `INSERT INTO reports ( |
| 551 | fingerprint, kind, version, os, arch, message, device, created_at, |
| 552 | source, label, error_type, error_message, top_frame, build_commit, channel, |
| 553 | language, view, breadcrumbs, component_stack, stack, occurred_at |
| 554 | ) |
| 555 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21)`, |
| 556 | ) |
| 557 | .bind( |
| 558 | fingerprint, |
| 559 | r.kind, |
| 560 | r.version, |
| 561 | r.os, |
| 562 | r.arch, |
| 563 | message, |
| 564 | JSON.stringify(r.device ?? {}), |
| 565 | now, |
| 566 | source, |
| 567 | label, |
| 568 | errorType, |
| 569 | errorMessage, |
| 570 | topFrame, |
| 571 | buildCommit, |
| 572 | channel, |
| 573 | r.language ?? "", |
| 574 | view, |
| 575 | JSON.stringify(breadcrumbs), |
| 576 | componentStack, |
| 577 | stack, |
| 578 | r.occurredAt ?? "", |
| 579 | ) |
| 580 | .run(); |
| 581 | |
| 582 | await env.DB.prepare( |
| 583 | `DELETE FROM reports |
| 584 | WHERE fingerprint = ?1 |
| 585 | AND id NOT IN ( |
| 586 | SELECT id FROM (SELECT id FROM reports WHERE fingerprint = ?1 ORDER BY id ASC LIMIT 1) |
| 587 | UNION |
| 588 | SELECT id FROM (SELECT id FROM reports WHERE fingerprint = ?1 ORDER BY id DESC LIMIT ?2) |
| 589 | )`, |
| 590 | ) |
| 591 | .bind(fingerprint, LATEST_SAMPLES_PER_GROUP) |
| 592 | .run(); |
| 593 | } catch (err) { |
| 594 | return storageUnavailable("report", err); |
| 595 | } |
| 596 | |
| 597 | return new Response("ok", { status: 202 }); |
| 598 | } |
| 599 | |
| 600 | async function handlePing(request: Request, env: Env): Promise<Response> { |
| 601 | const ip = request.headers.get("cf-connecting-ip") ?? "unknown"; |
| 602 | const { success } = await env.PING_LIMITER.limit({ key: ip }); |
| 603 | if (!success) return new Response("rate limited", { status: 429 }); |
| 604 | |
| 605 | const raw = await readJSON(request); |
| 606 | if (raw instanceof Response) return raw; |
| 607 | const parsed = Ping.safeParse(raw); |
| 608 | if (!parsed.success) return new Response("bad request", { status: 400 }); |
| 609 | const p = parsed.data; |
| 610 | const tables = telemetryTableNames(p.surface); |
| 611 | |
| 612 | try { |
| 613 | if (p.surface === "cli") await ensureCLITelemetrySchema(env); |
| 614 | await env.DB.prepare( |
| 615 | `INSERT INTO ${tables.pings} (date, install_id, version, os, arch, os_version, opens) |
| 616 | VALUES (date('now'), ?1, ?2, ?3, ?4, ?5, 1) |
| 617 | ON CONFLICT (date, install_id) DO UPDATE SET |
| 618 | opens = opens + 1, version = ?2, os_version = ?5`, |
| 619 | ) |
| 620 | .bind(p.installId, p.version, p.os, p.arch, p.osVersion ?? "") |
| 621 | .run(); |
| 622 | } catch (err) { |
| 623 | return storageUnavailable("ping", err); |
| 624 | } |
| 625 | |
| 626 | return new Response("ok", { status: 202 }); |
| 627 | } |
| 628 | |
| 629 | async function handleMetrics(request: Request, env: Env): Promise<Response> { |
| 630 | const ip = request.headers.get("cf-connecting-ip") ?? "unknown"; |
| 631 | const { success } = await env.METRICS_LIMITER.limit({ key: ip }); |
| 632 | if (!success) return new Response("rate limited", { status: 429 }); |
| 633 | |
| 634 | const raw = await readJSON(request); |
| 635 | if (raw instanceof Response) return raw; |
| 636 | const parsed = Metrics.safeParse(raw); |
| 637 | if (!parsed.success) return new Response("bad request", { status: 400 }); |
| 638 | const m = parsed.data; |
| 639 | if (m.counters.length === 0) return new Response("ok", { status: 202 }); |
| 640 | const tables = telemetryTableNames(m.surface); |
| 641 | |
| 642 | try { |
| 643 | if (m.surface === "cli") await ensureCLITelemetrySchema(env); |
| 644 | const upsert = env.DB.prepare( |
| 645 | `INSERT INTO ${tables.metrics} (date, version, os, signal, bucket, count) |
| 646 | VALUES (date('now'), ?1, ?2, ?3, ?4, ?5) |
| 647 | ON CONFLICT (date, version, os, signal, bucket) DO UPDATE SET |
| 648 | count = count + ?5`, |
| 649 | ); |
| 650 | await env.DB.batch(m.counters.map((c) => upsert.bind(m.version, m.os, c.signal, c.bucket, c.count))); |
| 651 | } catch (err) { |
| 652 | return storageUnavailable("metrics", err); |
| 653 | } |
| 654 | if (m.installId) { |
| 655 | const userUpsert = env.DB.prepare( |
| 656 | `INSERT INTO ${tables.metricUsers} (date, version, os, signal, bucket, install_id) |
| 657 | VALUES (date('now'), ?1, ?2, ?3, ?4, ?5) |
| 658 | ON CONFLICT (date, signal, bucket, install_id) DO UPDATE SET |
| 659 | version = ?1, os = ?2`, |
| 660 | ); |
| 661 | try { |
| 662 | await env.DB.batch(m.counters.map((c) => userUpsert.bind(m.version, m.os, c.signal, c.bucket, m.installId))); |
| 663 | } catch (err) { |
| 664 | console.warn("metric_users write failed", err); |
| 665 | } |
| 666 | } |
| 667 | |
| 668 | return new Response("ok", { status: 202 }); |
| 669 | } |
| 670 | |
| 671 | const UserAction = z.object({ |
| 672 | action: z.enum(["role", "delete"]), |
| 673 | userId: z.coerce.number().int().positive(), |
| 674 | role: z.enum(["pending", "viewer", "admin"]).optional(), |
| 675 | }); |
| 676 | |
| 677 | const GroupAction = z.object({ |
| 678 | action: z.enum(["status", "delete", "note", "resolution", "severity"]), |
| 679 | status: z.enum(["open", "resolved", "ignored"]).optional(), |
| 680 | note: z.string().max(500).optional(), |
| 681 | resolvedIn: z.string().max(64).optional(), |
| 682 | severity: z.enum(["low", "medium", "high", "critical"]).optional(), |
| 683 | }); |
| 684 | |
| 685 | async function formObject(request: Request): Promise<Record<string, string>> { |
| 686 | const form = await request.formData(); |
| 687 | const out: Record<string, string> = {}; |
| 688 | for (const [k, v] of form) out[k] = typeof v === "string" ? v : ""; |
| 689 | return out; |
| 690 | } |
| 691 | |
| 692 | type StatsFilters = { |
| 693 | surface: "desktop" | "cli"; |
| 694 | status: string; |
| 695 | source: string; |
| 696 | version: string; |
| 697 | os: string; |
| 698 | platform: string; |
| 699 | newLatest: boolean; |
| 700 | regressed: boolean; |
| 701 | windowDays: 7 | 30; |
| 702 | preferenceMode: "users" | "opens"; |
| 703 | }; |
| 704 | |
| 705 | function statsFilters(url: URL): StatsFilters { |
| 706 | const status = url.searchParams.get("status") ?? ""; |
| 707 | const surface = url.searchParams.get("surface") ?? "desktop"; |
| 708 | const windowParam = url.searchParams.get("window") ?? ""; |
| 709 | return { |
| 710 | surface: surface === "cli" ? "cli" : "desktop", |
| 711 | status: ["open", "resolved", "ignored"].includes(status) ? status : "", |
| 712 | source: (url.searchParams.get("source") ?? "").slice(0, 32), |
| 713 | version: (url.searchParams.get("version") ?? "").slice(0, 64), |
| 714 | os: (url.searchParams.get("os") ?? "").slice(0, 32), |
| 715 | platform: (url.searchParams.get("platform") ?? "").slice(0, 80), |
| 716 | newLatest: url.searchParams.get("new") === "latest", |
| 717 | regressed: url.searchParams.get("regressed") === "1", |
| 718 | windowDays: windowParam === "7d" ? 7 : 30, |
| 719 | preferenceMode: url.searchParams.get("prefs") === "opens" ? "opens" : "users", |
| 720 | }; |
| 721 | } |
| 722 | |
| 723 | async function crashGroups(env: Env, filters: StatsFilters, latestVersion: string) { |
| 724 | const where: string[] = [diagnosticWindowWhere(filters.windowDays)]; |
| 725 | const binds: unknown[] = []; |
| 726 | const add = (sql: string, value?: unknown) => { |
| 727 | where.push(sql.replace("?", `?${binds.length + 1}`)); |
| 728 | if (value !== undefined) binds.push(value); |
| 729 | }; |
| 730 | if (filters.status) add("status = ?", filters.status); |
| 731 | if (filters.source) add("source = ?", filters.source); |
| 732 | if (filters.version) add("last_version = ?", filters.version); |
| 733 | if (filters.os) add("last_os = ?", filters.os); |
| 734 | if (filters.platform) add("last_os || ' ' || last_arch = ?", filters.platform); |
| 735 | if (filters.newLatest && latestVersion) add("first_version = ?", latestVersion); |
| 736 | if (filters.regressed) where.push("regressed_at <> ''"); |
| 737 | let latestOrder = ""; |
| 738 | if (latestVersion) { |
| 739 | latestOrder = `CASE WHEN first_version = ?${binds.length + 1} THEN 0 ELSE 1 END,`; |
| 740 | binds.push(latestVersion); |
| 741 | } |
| 742 | const sql = `SELECT fingerprint, kind, count, first_version, last_version, substr(last_seen, 1, 10) AS seen, |
| 743 | status, title, source, label, error_type, top_frame, severity, last_os, last_arch, last_channel, regressed_at |
| 744 | FROM groups ${where.length ? `WHERE ${where.join(" AND ")}` : ""} |
| 745 | ORDER BY |
| 746 | CASE WHEN status = 'open' THEN 0 ELSE 1 END, |
| 747 | CASE |
| 748 | WHEN severity = 'critical' THEN 0 |
| 749 | WHEN ${developmentGroupSQL} |
| 750 | OR title = '[window.error] Script error.' |
| 751 | OR title LIKE '%ResizeObserver loop %' |
| 752 | OR title LIKE '%Minified React error #520%' |
| 753 | OR title LIKE '%additional File object is not a file on the disk%' |
| 754 | THEN 3 |
| 755 | WHEN severity = 'high' THEN 1 |
| 756 | WHEN severity = 'medium' THEN 2 |
| 757 | ELSE 3 |
| 758 | END, |
| 759 | ${latestOrder} |
| 760 | CASE WHEN regressed_at <> '' THEN 0 ELSE 1 END, |
| 761 | count DESC, |
| 762 | last_seen DESC |
| 763 | LIMIT 50`; |
| 764 | const stmt = env.DB.prepare(sql); |
| 765 | const query = binds.length ? stmt.bind(...binds) : stmt; |
| 766 | const result = await query.all<{ |
| 767 | fingerprint: string; |
| 768 | kind: string; |
| 769 | count: number; |
| 770 | first_version: string; |
| 771 | last_version: string; |
| 772 | seen: string; |
| 773 | status: string; |
| 774 | title: string; |
| 775 | source: string; |
| 776 | label: string; |
| 777 | error_type: string; |
| 778 | top_frame: string; |
| 779 | severity: string; |
| 780 | last_os: string; |
| 781 | last_arch: string; |
| 782 | last_channel: string; |
| 783 | regressed_at: string; |
| 784 | }>(); |
| 785 | result.results = result.results |
| 786 | .map((row) => ({ ...row, severity: effectiveGroupSeverity(row), development: isDevelopmentGroup(row) })) |
| 787 | .sort((a, b) => compareDiagnosticPriority(a, b, latestVersion)); |
| 788 | return result; |
| 789 | } |
| 790 | |
| 791 | type GroupPriorityRow = { |
| 792 | fingerprint: string; |
| 793 | status: string; |
| 794 | severity: string; |
| 795 | regressed_at: string; |
| 796 | first_version: string; |
| 797 | count: number; |
| 798 | seen: string; |
| 799 | title: string; |
| 800 | last_version: string; |
| 801 | last_channel: string; |
| 802 | }; |
| 803 | |
| 804 | const developmentGroupSQL = `fingerprint LIKE 'dev:%'`; |
| 805 | |
| 806 | export function isDevelopmentGroup(row: Pick<GroupPriorityRow, "fingerprint">): boolean { |
| 807 | // Legacy hashes can contain release observations between retained development |
| 808 | // samples, so only the explicit namespace proves that a group is dev-only. |
| 809 | return row.fingerprint.startsWith(DEVELOPMENT_FINGERPRINT_PREFIX); |
| 810 | } |
| 811 | |
| 812 | export function effectiveGroupSeverity( |
| 813 | row: Pick<GroupPriorityRow, "fingerprint" | "severity" | "title">, |
| 814 | ): string { |
| 815 | if (row.severity === "critical") return row.severity; |
| 816 | if (isDevelopmentGroup(row)) return "low"; |
| 817 | if ( |
| 818 | row.title === "[window.error] Script error." || |
| 819 | row.title.includes("ResizeObserver loop ") || |
| 820 | row.title.includes("Minified React error #520") || |
| 821 | row.title.includes("additional File object is not a file on the disk") |
| 822 | ) { |
| 823 | return "low"; |
| 824 | } |
| 825 | return row.severity; |
| 826 | } |
| 827 | |
| 828 | function compareDiagnosticPriority(a: GroupPriorityRow, b: GroupPriorityRow, latestVersion: string): number { |
| 829 | const statusRank = (value: string) => (value === "open" ? 0 : 1); |
| 830 | const severityRank = (value: string) => ({ critical: 0, high: 1, medium: 2, low: 3 })[value] ?? 4; |
| 831 | return ( |
| 832 | statusRank(a.status) - statusRank(b.status) || |
| 833 | severityRank(a.severity) - severityRank(b.severity) || |
| 834 | Number(b.first_version === latestVersion) - Number(a.first_version === latestVersion) || |
| 835 | Number(Boolean(b.regressed_at)) - Number(Boolean(a.regressed_at)) || |
| 836 | b.count - a.count || |
| 837 | b.seen.localeCompare(a.seen) |
| 838 | ); |
| 839 | } |
| 840 | |
| 841 | type ParsedVersion = { |
| 842 | version: string; |
| 843 | major: number; |
| 844 | minor: number; |
| 845 | patch: number; |
| 846 | }; |
| 847 | |
| 848 | function parseReleaseVersion(version: string): ParsedVersion | null { |
| 849 | // The dashboard's "latest" lane is for shipped stable builds. Development, |
| 850 | // prerelease, and build-metadata values remain visible in version facets but |
| 851 | // must not become the release baseline used for regression triage. |
| 852 | const m = version.trim().match(/^v?(\d+)\.(\d+)\.(\d+)$/); |
| 853 | if (!m) return null; |
| 854 | return { |
| 855 | version, |
| 856 | major: Number(m[1]), |
| 857 | minor: Number(m[2]), |
| 858 | patch: Number(m[3]), |
| 859 | }; |
| 860 | } |
| 861 | |
| 862 | export function newestReleaseVersion(versions: string[]): string { |
| 863 | const parsed = versions |
| 864 | .filter((v) => v && v.toLowerCase() !== "dev") |
| 865 | .map(parseReleaseVersion) |
| 866 | .filter((v): v is ParsedVersion => v !== null); |
| 867 | parsed.sort( |
| 868 | (a, b) => |
| 869 | b.major - a.major || |
| 870 | b.minor - a.minor || |
| 871 | b.patch - a.patch || |
| 872 | b.version.localeCompare(a.version), |
| 873 | ); |
| 874 | return parsed[0]?.version ?? ""; |
| 875 | } |
| 876 | |
| 877 | async function latestObservedVersion(env: Env, surface: ClientSurfaceName): Promise<string> { |
| 878 | const table = telemetryTableNames(surface).pings; |
| 879 | // Require independent installations and use pings as the sole source of |
| 880 | // release truth. A single synthetic diagnostic must never promote v9.9.9 (or |
| 881 | // a prerelease) to "latest" for every report group. |
| 882 | const sql = `SELECT version FROM ${table} |
| 883 | WHERE date >= date('now', '-29 day') AND version <> '' |
| 884 | GROUP BY version HAVING COUNT(DISTINCT install_id) >= 2`; |
| 885 | const rows = await env.DB.prepare(sql).all<{ version: string }>(); |
| 886 | return newestReleaseVersion(rows.results.map((r) => r.version)); |
| 887 | } |
| 888 | |
| 889 | type OverviewCounts = { |
| 890 | latestAdoptionPct: number | null; |
| 891 | openReports: number; |
| 892 | newLatestReports: number; |
| 893 | regressedReports: number; |
| 894 | criticalOpenReports: number; |
| 895 | }; |
| 896 | |
| 897 | async function latestAdoptionPct(env: Env, latestVersion: string, days: 7 | 30, surface: ClientSurfaceName): Promise<number | null> { |
| 898 | if (!latestVersion) return null; |
| 899 | const table = telemetryTableNames(surface).pings; |
| 900 | const row = await env.DB.prepare( |
| 901 | `SELECT |
| 902 | COUNT(DISTINCT install_id) AS total_installs, |
| 903 | COUNT(DISTINCT CASE WHEN version = ?1 THEN install_id END) AS latest_installs |
| 904 | FROM ${table} WHERE date >= date('now', '${currentWindowSince(days)}')`, |
| 905 | ) |
| 906 | .bind(latestVersion) |
| 907 | .first<{ total_installs: number; latest_installs: number }>(); |
| 908 | const total = Number(row?.total_installs ?? 0); |
| 909 | if (!total) return null; |
| 910 | return (Number(row?.latest_installs ?? 0) / total) * 100; |
| 911 | } |
| 912 | |
| 913 | async function diagnosticOverview(env: Env, latestVersion: string, days: 7 | 30, surface: ClientSurfaceName): Promise<OverviewCounts> { |
| 914 | if (surface === "cli") { |
| 915 | return { |
| 916 | latestAdoptionPct: await latestAdoptionPct(env, latestVersion, days, surface), |
| 917 | openReports: 0, |
| 918 | newLatestReports: 0, |
| 919 | regressedReports: 0, |
| 920 | criticalOpenReports: 0, |
| 921 | }; |
| 922 | } |
| 923 | // Keep the overview's red state aligned with the effective severity used by |
| 924 | // the diagnostics list. Historical rows retain their stored severity, so |
| 925 | // known browser notices and development builds must be discounted here too. |
| 926 | const criticalActionable = `(severity = 'critical' OR ( |
| 927 | severity = 'high' |
| 928 | AND kind <> 'performance' |
| 929 | AND NOT ${developmentGroupSQL} |
| 930 | AND title <> '[window.error] Script error.' |
| 931 | AND title NOT LIKE '%ResizeObserver loop %' |
| 932 | AND title NOT LIKE '%Minified React error #520%' |
| 933 | AND title NOT LIKE '%additional File object is not a file on the disk%' |
| 934 | ))`; |
| 935 | const diagnosticCounts = latestVersion |
| 936 | ? env.DB.prepare( |
| 937 | `SELECT |
| 938 | SUM(CASE WHEN status = 'open' THEN 1 ELSE 0 END) AS open_reports, |
| 939 | SUM(CASE WHEN first_version = ?1 THEN 1 ELSE 0 END) AS new_latest_reports, |
| 940 | SUM(CASE WHEN regressed_at <> '' THEN 1 ELSE 0 END) AS regressed_reports, |
| 941 | SUM(CASE WHEN status = 'open' AND ${criticalActionable} THEN 1 ELSE 0 END) AS critical_open_reports |
| 942 | FROM groups WHERE ${diagnosticWindowWhere(days)}`, |
| 943 | ) |
| 944 | .bind(latestVersion) |
| 945 | .first<{ open_reports: number; new_latest_reports: number; regressed_reports: number; critical_open_reports: number }>() |
| 946 | : env.DB.prepare( |
| 947 | `SELECT |
| 948 | SUM(CASE WHEN status = 'open' THEN 1 ELSE 0 END) AS open_reports, |
| 949 | 0 AS new_latest_reports, |
| 950 | SUM(CASE WHEN regressed_at <> '' THEN 1 ELSE 0 END) AS regressed_reports, |
| 951 | SUM(CASE WHEN status = 'open' AND ${criticalActionable} THEN 1 ELSE 0 END) AS critical_open_reports |
| 952 | FROM groups WHERE ${diagnosticWindowWhere(days)}`, |
| 953 | ).first<{ open_reports: number; new_latest_reports: number; regressed_reports: number; critical_open_reports: number }>(); |
| 954 | const [row, adoptionPct] = await Promise.all([ |
| 955 | diagnosticCounts, |
| 956 | latestAdoptionPct(env, latestVersion, days, surface), |
| 957 | ]); |
| 958 | return { |
| 959 | latestAdoptionPct: adoptionPct, |
| 960 | openReports: Number(row?.open_reports ?? 0), |
| 961 | newLatestReports: Number(row?.new_latest_reports ?? 0), |
| 962 | regressedReports: Number(row?.regressed_reports ?? 0), |
| 963 | criticalOpenReports: Number(row?.critical_open_reports ?? 0), |
| 964 | }; |
| 965 | } |
| 966 | |
| 967 | function currentWindowSince(days: 7 | 30): string { |
| 968 | return `-${days - 1} day`; |
| 969 | } |
| 970 | |
| 971 | export function diagnosticWindowWhere(days: 7 | 30): string { |
| 972 | return `date(last_seen) >= date('now', '${currentWindowSince(days)}')`; |
| 973 | } |
| 974 | |
| 975 | function previousWindowSince(days: 7 | 30): string { |
| 976 | return `-${days * 2 - 1} day`; |
| 977 | } |
| 978 | |
| 979 | function previousWindowUntil(days: 7 | 30): string { |
| 980 | return currentWindowSince(days); |
| 981 | } |
| 982 | |
| 983 | async function metricRows(env: Env, days: 7 | 30, surface: ClientSurfaceName, previous = false): Promise<{ signal: string; bucket: string; total: number }[]> { |
| 984 | const where = previous |
| 985 | ? `date >= date('now', '${previousWindowSince(days)}') AND date < date('now', '${previousWindowUntil(days)}')` |
| 986 | : `date >= date('now', '${currentWindowSince(days)}')`; |
| 987 | const table = telemetryTableNames(surface).metrics; |
| 988 | const rows = await env.DB.prepare( |
| 989 | `SELECT signal, bucket, SUM(count) AS total FROM ${table} WHERE ${where} GROUP BY signal, bucket ORDER BY signal, total DESC`, |
| 990 | ).all<{ signal: string; bucket: string; total: number }>(); |
| 991 | return rows.results; |
| 992 | } |
| 993 | |
| 994 | // The 30-day desktop window is served from the cron-built rollup: computing it |
| 995 | // live exceeds what D1 spends on one query (see refreshMetricUserRollup). Null |
| 996 | // here means "not computed yet", which the dashboard says out loud rather than |
| 997 | // rendering as an empty result. |
| 998 | async function rollupMetricUserRows( |
| 999 | env: Env, |
| 1000 | ): Promise<{ rows: { signal: string; bucket: string; total: number }[]; computedAt: string } | null> { |
| 1001 | try { |
| 1002 | await ensureRollupSchema(env); |
| 1003 | const rows = await env.DB.prepare( |
| 1004 | `SELECT signal, bucket, total, computed_at FROM metric_user_rollup WHERE window_days = ?1 ORDER BY signal, total DESC`, |
| 1005 | ) |
| 1006 | .bind(ROLLUP_WINDOW_DAYS) |
| 1007 | .all<{ signal: string; bucket: string; total: number; computed_at: string }>(); |
| 1008 | if (!rows.results.length) return null; |
| 1009 | // Oldest wins: the cursor refreshes a slice at a time, so this is how far |
| 1010 | // behind the least recently recomputed signal is. |
| 1011 | const computedAt = rows.results.reduce((min, r) => (r.computed_at < min ? r.computed_at : min), rows.results[0].computed_at); |
| 1012 | return { rows: rows.results, computedAt }; |
| 1013 | } catch (err) { |
| 1014 | console.warn("metric_user_rollup read failed", err); |
| 1015 | return null; |
| 1016 | } |
| 1017 | } |
| 1018 | |
| 1019 | // Null means the query did not complete, which is distinct from "no rows": at |
| 1020 | // ~1M rows a day, the 30-day COUNT(DISTINCT install_id) exceeds what D1 will |
| 1021 | // spend on one query and comes back as a CPU-limit reset. Rendering that as an |
| 1022 | // empty dashboard would read as "nobody uses these settings". |
| 1023 | async function metricUserRows( |
| 1024 | env: Env, |
| 1025 | days: 7 | 30, |
| 1026 | surface: ClientSurfaceName, |
| 1027 | ): Promise<{ rows: { signal: string; bucket: string; total: number }[]; computedAt: string } | null> { |
| 1028 | if (surface === "desktop" && days === ROLLUP_WINDOW_DAYS) return rollupMetricUserRows(env); |
| 1029 | try { |
| 1030 | const table = telemetryTableNames(surface).metricUsers; |
| 1031 | const rows = await env.DB.prepare( |
| 1032 | `SELECT signal, bucket, COUNT(DISTINCT install_id) AS total FROM ${table} WHERE date >= date('now', '${currentWindowSince(days)}') GROUP BY signal, bucket ORDER BY signal, total DESC`, |
| 1033 | ).all<{ signal: string; bucket: string; total: number }>(); |
| 1034 | return { rows: rows.results, computedAt: "" }; |
| 1035 | } catch (err) { |
| 1036 | console.warn("metric_users query failed", err); |
| 1037 | return null; |
| 1038 | } |
| 1039 | } |
| 1040 | |
| 1041 | type Bar = { label: string; users: number }; |
| 1042 | type MetricTotals = { signal: string; bucket: string; total: number }[]; |
| 1043 | |
| 1044 | // Each stats module renders only its own section, so a page load should query |
| 1045 | // only what that section shows — the 30-day COUNT(DISTINCT) over metric_users, |
| 1046 | // the heaviest query, is read solely by the preferences module. |
| 1047 | async function handleStats(request: Request, env: Env, user: User, activeModule: StatsModule): Promise<Response> { |
| 1048 | const url = new URL(request.url); |
| 1049 | const filters = statsFilters(url); |
| 1050 | const days = filters.windowDays; |
| 1051 | const since = currentWindowSince(days); |
| 1052 | const surface = activeModule === "diagnostics" ? "desktop" : filters.surface; |
| 1053 | if (activeModule === "diagnostics") filters.surface = "desktop"; |
| 1054 | if (surface === "cli") await ensureCLITelemetrySchema(env); |
| 1055 | const pingsTable = telemetryTableNames(surface).pings; |
| 1056 | const bars = (sql: string) => env.DB.prepare(sql).all<Bar>().then((r) => r.results); |
| 1057 | const pingVersions = () => |
| 1058 | bars(`SELECT version AS label, COUNT(DISTINCT install_id) AS users FROM ${pingsTable} WHERE date >= date('now', '${since}') GROUP BY label ORDER BY users DESC LIMIT 15`); |
| 1059 | const pingPlatforms = () => |
| 1060 | bars(`SELECT os || ' ' || arch AS label, COUNT(DISTINCT install_id) AS users FROM ${pingsTable} WHERE date >= date('now', '${since}') GROUP BY label ORDER BY users DESC`); |
| 1061 | |
| 1062 | let daily: { date: string; users: number; opens: number }[] = []; |
| 1063 | let versions: Bar[] = []; |
| 1064 | let platforms: Bar[] = []; |
| 1065 | let crashes: Awaited<ReturnType<typeof crashGroups>>["results"] = []; |
| 1066 | let metrics: MetricTotals = []; |
| 1067 | let previousMetrics: MetricTotals = []; |
| 1068 | let metricUsers: MetricTotals = []; |
| 1069 | let metricUsersUnavailable = false; |
| 1070 | let metricUsersComputedAt = ""; |
| 1071 | let sources: Bar[] = []; |
| 1072 | let overview: OverviewCounts = { |
| 1073 | latestAdoptionPct: null, |
| 1074 | openReports: 0, |
| 1075 | newLatestReports: 0, |
| 1076 | regressedReports: 0, |
| 1077 | criticalOpenReports: 0, |
| 1078 | }; |
| 1079 | let latestVersion = ""; |
| 1080 | |
| 1081 | if (activeModule === "usage") { |
| 1082 | latestVersion = await latestObservedVersion(env, surface); |
| 1083 | const [dailyR, versionsR, platformsR, metricsR, overviewR] = await Promise.all([ |
| 1084 | env.DB.prepare( |
| 1085 | `SELECT date, COUNT(*) AS users, SUM(opens) AS opens FROM ${pingsTable} WHERE date >= date('now', '${since}') GROUP BY date`, |
| 1086 | ).all<{ date: string; users: number; opens: number }>(), |
| 1087 | pingVersions(), |
| 1088 | pingPlatforms(), |
| 1089 | metricRows(env, days, surface), |
| 1090 | diagnosticOverview(env, latestVersion, days, surface), |
| 1091 | ]); |
| 1092 | daily = dailyR.results; |
| 1093 | versions = versionsR; |
| 1094 | platforms = platformsR; |
| 1095 | metrics = metricsR; |
| 1096 | overview = overviewR; |
| 1097 | } else if (activeModule === "diagnostics") { |
| 1098 | latestVersion = await latestObservedVersion(env, "desktop"); |
| 1099 | const [crashesR, sourcesR, versionsR, platformsR] = await Promise.all([ |
| 1100 | crashGroups(env, filters, latestVersion), |
| 1101 | bars(`SELECT source AS label, COUNT(*) AS users FROM groups WHERE ${diagnosticWindowWhere(days)} GROUP BY source ORDER BY users DESC`), |
| 1102 | pingVersions(), |
| 1103 | pingPlatforms(), |
| 1104 | ]); |
| 1105 | crashes = crashesR.results; |
| 1106 | sources = sourcesR; |
| 1107 | versions = versionsR; |
| 1108 | platforms = platformsR; |
| 1109 | } else if (activeModule === "preferences") { |
| 1110 | const [metricsR, usersR] = await Promise.all([metricRows(env, days, surface), metricUserRows(env, days, surface)]); |
| 1111 | metrics = metricsR; |
| 1112 | metricUsersUnavailable = usersR === null; |
| 1113 | metricUsers = usersR?.rows ?? []; |
| 1114 | metricUsersComputedAt = usersR?.computedAt ?? ""; |
| 1115 | } else { |
| 1116 | const [metricsR, previousMetricsR, usersR] = await Promise.all([ |
| 1117 | metricRows(env, days, surface), |
| 1118 | metricRows(env, days, surface, true), |
| 1119 | metricUserRows(env, days, surface), |
| 1120 | ]); |
| 1121 | metrics = metricsR; |
| 1122 | previousMetrics = previousMetricsR; |
| 1123 | metricUsersUnavailable = usersR === null; |
| 1124 | metricUsers = usersR?.rows ?? []; |
| 1125 | metricUsersComputedAt = usersR?.computedAt ?? ""; |
| 1126 | } |
| 1127 | |
| 1128 | return html( |
| 1129 | renderStats( |
| 1130 | { daily, versions, platforms, crashes, metrics, previousMetrics, metricUsers, metricUsersUnavailable, metricUsersComputedAt, sources, overview, latestVersion, filters }, |
| 1131 | user, |
| 1132 | activeModule, |
| 1133 | ), |
| 1134 | ); |
| 1135 | } |
| 1136 | |
| 1137 | async function handleGroup(env: Env, fingerprint: string, user: User): Promise<Response> { |
| 1138 | const group = await env.DB.prepare("SELECT * FROM groups WHERE fingerprint = ?1").bind(fingerprint).first<Group>(); |
| 1139 | if (!group) return new Response("not found", { status: 404 }); |
| 1140 | group.severity = effectiveGroupSeverity(group); |
| 1141 | const reports = await env.DB.prepare( |
| 1142 | `SELECT version, os, arch, message, device, created_at, source, label, error_type, error_message, |
| 1143 | top_frame, build_commit, channel, language, view, breadcrumbs, component_stack, stack, occurred_at |
| 1144 | FROM reports WHERE fingerprint = ?1 ORDER BY id DESC`, |
| 1145 | ) |
| 1146 | .bind(fingerprint) |
| 1147 | .all<{ |
| 1148 | version: string; |
| 1149 | os: string; |
| 1150 | arch: string; |
| 1151 | message: string; |
| 1152 | device: string; |
| 1153 | created_at: string; |
| 1154 | source: string; |
| 1155 | label: string; |
| 1156 | error_type: string; |
| 1157 | error_message: string; |
| 1158 | top_frame: string; |
| 1159 | build_commit: string; |
| 1160 | channel: string; |
| 1161 | language: string; |
| 1162 | view: string; |
| 1163 | breadcrumbs: string; |
| 1164 | component_stack: string; |
| 1165 | stack: string; |
| 1166 | occurred_at: string; |
| 1167 | }>(); |
| 1168 | return html(renderGroup(group, reports.results, user)); |
| 1169 | } |
| 1170 | |
| 1171 | async function handleGroupAction(request: Request, env: Env, admin: User, fingerprint: string): Promise<Response> { |
| 1172 | if (!sameOrigin(request)) return new Response("forbidden", { status: 403 }); |
| 1173 | const parsed = GroupAction.safeParse(await formObject(request)); |
| 1174 | if (!parsed.success) return redirect(`/stats/group/${fingerprint}`); |
| 1175 | const a = parsed.data; |
| 1176 | |
| 1177 | if (a.action === "delete") { |
| 1178 | await env.DB.batch([ |
| 1179 | env.DB.prepare("DELETE FROM reports WHERE fingerprint = ?1").bind(fingerprint), |
| 1180 | env.DB.prepare("DELETE FROM groups WHERE fingerprint = ?1").bind(fingerprint), |
| 1181 | ]); |
| 1182 | await logAction(env, admin, "delete_group", fingerprint.slice(0, 8)); |
| 1183 | return redirect("/stats"); |
| 1184 | } |
| 1185 | if (a.action === "status") { |
| 1186 | const status = a.status ?? "open"; |
| 1187 | await env.DB.prepare( |
| 1188 | "UPDATE groups SET status = ?1, resolved_at = CASE WHEN ?1 = 'resolved' THEN ?3 ELSE resolved_at END WHERE fingerprint = ?2", |
| 1189 | ) |
| 1190 | .bind(status, fingerprint, new Date().toISOString()) |
| 1191 | .run(); |
| 1192 | await logAction(env, admin, "set_status", fingerprint.slice(0, 8), status); |
| 1193 | return redirect(`/stats/group/${fingerprint}`); |
| 1194 | } |
| 1195 | if (a.action === "resolution") { |
| 1196 | await env.DB.prepare("UPDATE groups SET resolved_in = ?1 WHERE fingerprint = ?2") |
| 1197 | .bind(a.resolvedIn ?? "", fingerprint) |
| 1198 | .run(); |
| 1199 | await logAction(env, admin, "set_resolved_in", fingerprint.slice(0, 8), a.resolvedIn ?? ""); |
| 1200 | return redirect(`/stats/group/${fingerprint}`); |
| 1201 | } |
| 1202 | if (a.action === "severity") { |
| 1203 | await env.DB.prepare("UPDATE groups SET severity = ?1 WHERE fingerprint = ?2") |
| 1204 | .bind(a.severity ?? "medium", fingerprint) |
| 1205 | .run(); |
| 1206 | await logAction(env, admin, "set_severity", fingerprint.slice(0, 8), a.severity ?? "medium"); |
| 1207 | return redirect(`/stats/group/${fingerprint}`); |
| 1208 | } |
| 1209 | await env.DB.prepare("UPDATE groups SET note = ?1 WHERE fingerprint = ?2").bind(a.note ?? "", fingerprint).run(); |
| 1210 | await logAction(env, admin, "set_note", fingerprint.slice(0, 8)); |
| 1211 | return redirect(`/stats/group/${fingerprint}`); |
| 1212 | } |
| 1213 | |
| 1214 | async function handleAdminUsers(request: Request, env: Env, admin: User): Promise<Response> { |
| 1215 | if (!sameOrigin(request)) return new Response("forbidden", { status: 403 }); |
| 1216 | const parsed = UserAction.safeParse(await formObject(request)); |
| 1217 | if (!parsed.success) return redirect("/admin"); |
| 1218 | const a = parsed.data; |
| 1219 | if (a.userId === admin.id) return redirect("/admin"); |
| 1220 | |
| 1221 | const target = await env.DB.prepare("SELECT email, role FROM access WHERE id = ?1") |
| 1222 | .bind(a.userId) |
| 1223 | .first<{ email: string; role: Role }>(); |
| 1224 | if (!target) return redirect("/admin"); |
| 1225 | |
| 1226 | if (a.action === "delete") { |
| 1227 | await env.DB.prepare("DELETE FROM access WHERE id = ?1").bind(a.userId).run(); |
| 1228 | await logAction(env, admin, "delete_user", target.email); |
| 1229 | return redirect("/admin"); |
| 1230 | } |
| 1231 | |
| 1232 | const role: Role = a.role ?? "pending"; |
| 1233 | const now = new Date().toISOString(); |
| 1234 | await env.DB.prepare("UPDATE access SET role = ?1, approved_at = ?2, approved_by = ?3 WHERE id = ?4") |
| 1235 | .bind(role, role === "pending" ? null : now, admin.email, a.userId) |
| 1236 | .run(); |
| 1237 | await logAction(env, admin, "set_role", target.email, `${target.role} → ${role}`); |
| 1238 | return redirect("/admin"); |
| 1239 | } |
| 1240 | |
| 1241 | async function handleAdminList(env: Env, admin: User): Promise<Response> { |
| 1242 | const users = await env.DB.prepare( |
| 1243 | "SELECT id, email, role, created_at, approved_at FROM access ORDER BY (role = 'pending') DESC, created_at DESC", |
| 1244 | ).all<UserRow>(); |
| 1245 | return html(renderUsers(admin, users.results)); |
| 1246 | } |
| 1247 | |
| 1248 | async function handleAdminAudit(env: Env, admin: User): Promise<Response> { |
| 1249 | const rows = await env.DB.prepare( |
| 1250 | "SELECT at, actor_email, action, target, detail FROM audit_log ORDER BY id DESC LIMIT 200", |
| 1251 | ).all<AuditRow>(); |
| 1252 | return html(renderAudit(admin, rows.results)); |
| 1253 | } |
| 1254 | |
| 1255 | function requireViewer(user: User | null, login: string): Response | null { |
| 1256 | if (!user) return redirect(login); |
| 1257 | if (!atLeast(user.role, "viewer")) return redirect("/account"); |
| 1258 | return null; |
| 1259 | } |
| 1260 | |
| 1261 | // The folded registry API runs against its own database and resolves identity |
| 1262 | // itself; hand it the second binding plus the account/site origins it expects. |
| 1263 | function registryBindings(env: Env): RegistryBindings { |
| 1264 | return { |
| 1265 | DB: env.REGISTRY_DB, |
| 1266 | WRITE_LIMITER: env.WRITE_LIMITER, |
| 1267 | ACCOUNTS_ORIGIN: env.ID_ORIGIN ?? "https://id.reasonix.io", |
| 1268 | APP_ORIGIN: env.APP_ORIGIN ?? "https://reasonix.io", |
| 1269 | ALLOWED_ORIGINS: env.ALLOWED_ORIGINS ?? "https://reasonix.io,https://www.reasonix.io", |
| 1270 | }; |
| 1271 | } |
| 1272 | |
| 1273 | function communityStatus(url: URL): string { |
| 1274 | const s = url.searchParams.get("status") ?? "pending"; |
| 1275 | return ["pending", "active", "hidden", "rejected"].includes(s) ? s : "pending"; |
| 1276 | } |
| 1277 | |
| 1278 | async function handleCommunityList(env: Env, admin: User, status: string): Promise<Response> { |
| 1279 | const rows = await new PackageRepo(env.REGISTRY_DB).listByStatus(status, 200); |
| 1280 | return html(renderCommunity(admin, rows, status)); |
| 1281 | } |
| 1282 | |
| 1283 | async function handleCommunityAction( |
| 1284 | request: Request, |
| 1285 | env: Env, |
| 1286 | admin: User, |
| 1287 | handle: string, |
| 1288 | name: string, |
| 1289 | action: string, |
| 1290 | ): Promise<Response> { |
| 1291 | if (!sameOrigin(request)) return new Response("forbidden", { status: 403 }); |
| 1292 | const form = await formObject(request); |
| 1293 | const backStatus = ["pending", "active", "hidden", "rejected"].includes(form.status) ? form.status : "pending"; |
| 1294 | const back = redirect(`/community?status=${backStatus}`); |
| 1295 | const slug = `${handle}/${name}`; |
| 1296 | const repo = new PackageRepo(env.REGISTRY_DB); |
| 1297 | const now = new Date().toISOString(); |
| 1298 | |
| 1299 | if (action === "verify" || action === "unverify") { |
| 1300 | await repo.setVerified(slug, action === "verify", now); |
| 1301 | await logAction(env, admin, `pkg_${action}`, slug); |
| 1302 | return back; |
| 1303 | } |
| 1304 | if (action === "approve") { |
| 1305 | const expectedStatus = ["pending", "hidden", "rejected"].includes(form.expectedStatus) |
| 1306 | ? form.expectedStatus |
| 1307 | : ""; |
| 1308 | if (!form.expectedVersion || !form.expectedUpdatedAt || !expectedStatus) { |
| 1309 | return new Response("Package review revision is missing. Refresh the review page and try again.", { |
| 1310 | status: 409, |
| 1311 | }); |
| 1312 | } |
| 1313 | const row = await repo.setStatusIfCurrent( |
| 1314 | slug, |
| 1315 | "active", |
| 1316 | form.expectedVersion, |
| 1317 | form.expectedUpdatedAt, |
| 1318 | expectedStatus, |
| 1319 | now, |
| 1320 | ); |
| 1321 | if (!row) { |
| 1322 | return new Response("Package changed since it was reviewed. Refresh and review the latest version.", { |
| 1323 | status: 409, |
| 1324 | }); |
| 1325 | } |
| 1326 | // Emit the publish event only after the reviewed revision becomes public. |
| 1327 | await new EventRepo(env.REGISTRY_DB).log({ |
| 1328 | type: "publish", |
| 1329 | packageId: row.id, |
| 1330 | actorHandle: row.scope_handle, |
| 1331 | summary: `published ${row.slug}@${row.latest_version}`, |
| 1332 | now, |
| 1333 | }); |
| 1334 | await logAction(env, admin, "pkg_approve", slug); |
| 1335 | return back; |
| 1336 | } |
| 1337 | await repo.setStatus(slug, action === "reject" ? "rejected" : "hidden", now); |
| 1338 | await logAction(env, admin, `pkg_${action}`, slug); |
| 1339 | return back; |
| 1340 | } |
| 1341 | |
| 1342 | // Time-series retention, run by the daily cron trigger. Every dashboard query |
| 1343 | // against the per-install tables reads at most the current window (-29 day), |
| 1344 | // while the aggregate `metrics` table also serves the 30d view's |
| 1345 | // previous-window delta (back to -59 day), so it keeps a doubled horizon. |
| 1346 | // `reports`/`groups` are excluded on purpose: they are the triage queue and |
| 1347 | // the regression baseline, are not date-partitioned, and their growth is |
| 1348 | // already bounded by per-group sampling. Without this purge the database |
| 1349 | // grows until D1's size cap, at which point every ingest write starts |
| 1350 | // throwing (all of /v1/ping, /v1/metrics and /v1/report 500 while reads keep |
| 1351 | // working — exactly the 2026-07-03 stats blackout). |
| 1352 | const RETENTION = [ |
| 1353 | { table: "pings", keepDays: 30 }, |
| 1354 | { table: "metrics", keepDays: 60 }, |
| 1355 | { table: "metric_users", keepDays: 30 }, |
| 1356 | { table: "cli_pings", keepDays: 30 }, |
| 1357 | { table: "cli_metrics", keepDays: 60 }, |
| 1358 | { table: "cli_metric_users", keepDays: 30 }, |
| 1359 | ] as const; |
| 1360 | // Deletes run in rowid chunks so a run never holds one giant transaction. |
| 1361 | // Steady state is one expired day per table; the chunk cap is a backstop that |
| 1362 | // still drains ~2M rows per table per run after an ingest outage or backlog. |
| 1363 | const RETENTION_CHUNK_ROWS = 10_000; |
| 1364 | const RETENTION_MAX_CHUNKS = 200; |
| 1365 | |
| 1366 | // Must match the sentinel entry in wrangler.toml [triggers] exactly — the |
| 1367 | // scheduled handler dispatches on controller.cron; every other trigger |
| 1368 | // (the retention cron, manual runs) falls through to the purge. |
| 1369 | const SENTINEL_CRON = "17 1,7,13,19 * * *"; |
| 1370 | const ROLLUP_CRON = "23 * * * *"; |
| 1371 | |
| 1372 | // The preferences module's 30-day COUNT(DISTINCT install_id) spans ~28M rows |
| 1373 | // and D1 abandons it mid-query. It cannot be summed from per-day totals either: |
| 1374 | // an install active on twelve days would count twelve times. So the window is |
| 1375 | // computed here instead, one signal at a time — a single signal takes ~1s, and |
| 1376 | // the cursor spreads the ~57 of them across hourly runs rather than blowing one |
| 1377 | // invocation's CPU budget. |
| 1378 | const ROLLUP_WINDOW_DAYS = 30; |
| 1379 | const ROLLUP_SIGNALS_PER_RUN = 8; |
| 1380 | |
| 1381 | const ROLLUP_SCHEMA_SQL = [ |
| 1382 | `CREATE TABLE IF NOT EXISTS metric_user_rollup ( |
| 1383 | window_days INTEGER NOT NULL, |
| 1384 | signal TEXT NOT NULL, |
| 1385 | bucket TEXT NOT NULL, |
| 1386 | total INTEGER NOT NULL, |
| 1387 | computed_at TEXT NOT NULL, |
| 1388 | PRIMARY KEY (window_days, signal, bucket) |
| 1389 | )`, |
| 1390 | `CREATE TABLE IF NOT EXISTS metric_user_rollup_state ( |
| 1391 | id INTEGER PRIMARY KEY CHECK (id = 1), |
| 1392 | next_signal INTEGER NOT NULL, |
| 1393 | updated_at TEXT NOT NULL |
| 1394 | )`, |
| 1395 | ] as const; |
| 1396 | |
| 1397 | const rollupSchemaPromises = new WeakMap<object, Promise<void>>(); |
| 1398 | |
| 1399 | function ensureRollupSchema(env: Pick<Env, "DB">): Promise<void> { |
| 1400 | const key = env.DB as unknown as object; |
| 1401 | const existing = rollupSchemaPromises.get(key); |
| 1402 | if (existing) return existing; |
| 1403 | const creation = env.DB |
| 1404 | .batch(ROLLUP_SCHEMA_SQL.map((sql) => env.DB.prepare(sql))) |
| 1405 | .then(() => undefined) |
| 1406 | .catch((err) => { |
| 1407 | rollupSchemaPromises.delete(key); |
| 1408 | throw err; |
| 1409 | }); |
| 1410 | rollupSchemaPromises.set(key, creation); |
| 1411 | return creation; |
| 1412 | } |
| 1413 | |
| 1414 | export async function refreshMetricUserRollup(env: Env, signalsPerRun = ROLLUP_SIGNALS_PER_RUN): Promise<void> { |
| 1415 | await ensureRollupSchema(env); |
| 1416 | const state = await env.DB.prepare("SELECT next_signal FROM metric_user_rollup_state WHERE id = 1").first<{ |
| 1417 | next_signal: number; |
| 1418 | }>(); |
| 1419 | const start = Number(state?.next_signal ?? 0) % METRIC_SIGNALS.length; |
| 1420 | const now = new Date().toISOString(); |
| 1421 | let advanced = 0; |
| 1422 | |
| 1423 | for (let i = 0; i < signalsPerRun && i < METRIC_SIGNALS.length; i++) { |
| 1424 | const signal = METRIC_SIGNALS[(start + i) % METRIC_SIGNALS.length]; |
| 1425 | try { |
| 1426 | const rows = await env.DB.prepare( |
| 1427 | `SELECT bucket, COUNT(DISTINCT install_id) AS total FROM metric_users |
| 1428 | WHERE date >= date('now', '-${ROLLUP_WINDOW_DAYS - 1} day') AND signal = ?1 |
| 1429 | GROUP BY bucket`, |
| 1430 | ) |
| 1431 | .bind(signal) |
| 1432 | .all<{ bucket: string; total: number }>(); |
| 1433 | // Delete and insert in one batch so a reader never sees a signal |
| 1434 | // half-replaced; an empty result still clears the previous window's rows. |
| 1435 | await env.DB.batch([ |
| 1436 | env.DB |
| 1437 | .prepare("DELETE FROM metric_user_rollup WHERE window_days = ?1 AND signal = ?2") |
| 1438 | .bind(ROLLUP_WINDOW_DAYS, signal), |
| 1439 | ...rows.results.map((r) => |
| 1440 | env.DB |
| 1441 | .prepare( |
| 1442 | `INSERT INTO metric_user_rollup (window_days, signal, bucket, total, computed_at) |
| 1443 | VALUES (?1, ?2, ?3, ?4, ?5)`, |
| 1444 | ) |
| 1445 | .bind(ROLLUP_WINDOW_DAYS, signal, r.bucket, r.total, now), |
| 1446 | ), |
| 1447 | ]); |
| 1448 | advanced++; |
| 1449 | } catch (err) { |
| 1450 | // One signal timing out must not strand the cursor on it forever. |
| 1451 | console.error(`rollup: ${signal} failed`, err); |
| 1452 | advanced++; |
| 1453 | } |
| 1454 | } |
| 1455 | |
| 1456 | await env.DB |
| 1457 | .prepare( |
| 1458 | `INSERT INTO metric_user_rollup_state (id, next_signal, updated_at) VALUES (1, ?1, ?2) |
| 1459 | ON CONFLICT(id) DO UPDATE SET next_signal = ?1, updated_at = ?2`, |
| 1460 | ) |
| 1461 | .bind((start + advanced) % METRIC_SIGNALS.length, now) |
| 1462 | .run(); |
| 1463 | console.log(`rollup: refreshed ${advanced} signals from index ${start}`); |
| 1464 | } |
| 1465 | |
| 1466 | // Ingest sentinel. The 2026-07-03 blackout went unnoticed for ten days because |
| 1467 | // clients swallow ping failures by design and nothing watched the write path. |
| 1468 | // Four times a day (hours chosen so the UTC day always has >1h of traffic; |
| 1469 | // ~14k DAU means a healthy hour is never empty) this probes the two failure |
| 1470 | // shapes independently: |
| 1471 | // 1. canary write into `pings` (immediately deleted) — catches writes |
| 1472 | // throwing, e.g. the D1 size cap, regardless of traffic; |
| 1473 | // 2. today's real ping and open totals compared with the previous run — |
| 1474 | // catches ingest dying upstream of the worker (edge blocking, client |
| 1475 | // regression) even after the UTC day already has traffic. |
| 1476 | // Alerts go to the optional ALERT_WEBHOOK secret; without it they still land |
| 1477 | // in the worker logs. While broken this fires at most 4 alerts/day. |
| 1478 | const CANARY_INSTALL_ID = "ffffffffffffffffffffffffffffffff"; |
| 1479 | |
| 1480 | function errText(err: unknown): string { |
| 1481 | return err instanceof Error ? err.message : String(err); |
| 1482 | } |
| 1483 | |
| 1484 | async function sendAlert(env: Env, text: string): Promise<void> { |
| 1485 | if (!env.ALERT_WEBHOOK) return; |
| 1486 | try { |
| 1487 | const webhook = new URL(env.ALERT_WEBHOOK); |
| 1488 | const feishu = webhook.hostname === "open.feishu.cn" || webhook.hostname === "open.larksuite.com"; |
| 1489 | const body = feishu ? { msg_type: "text", content: { text } } : { text }; |
| 1490 | const res = await fetch(webhook.toString(), { |
| 1491 | method: "POST", |
| 1492 | headers: { "content-type": "application/json" }, |
| 1493 | body: JSON.stringify(body), |
| 1494 | }); |
| 1495 | if (!res.ok) console.error(`alert webhook responded ${res.status}`); |
| 1496 | } catch (err) { |
| 1497 | console.error("alert webhook unreachable", err); |
| 1498 | } |
| 1499 | } |
| 1500 | |
| 1501 | async function runIngestSentinel(env: Env): Promise<void> { |
| 1502 | const problems: string[] = []; |
| 1503 | try { |
| 1504 | await env.DB.prepare( |
| 1505 | `INSERT INTO pings (date, install_id, version, os, arch, opens) |
| 1506 | VALUES (date('now'), ?1, 'canary', 'canary', 'canary', 0) |
| 1507 | ON CONFLICT (date, install_id) DO NOTHING`, |
| 1508 | ) |
| 1509 | .bind(CANARY_INSTALL_ID) |
| 1510 | .run(); |
| 1511 | // Also removes any leftover canary from a run that died mid-way. |
| 1512 | await env.DB.prepare("DELETE FROM pings WHERE install_id = ?1").bind(CANARY_INSTALL_ID).run(); |
| 1513 | } catch (err) { |
| 1514 | problems.push(`canary write failed: ${errText(err)}`); |
| 1515 | } |
| 1516 | try { |
| 1517 | // Auto-create the one-row checkpoint so existing databases do not need a |
| 1518 | // manual migration before this worker version is deployed. |
| 1519 | await env.DB.prepare( |
| 1520 | `CREATE TABLE IF NOT EXISTS ingest_sentinel_state ( |
| 1521 | id INTEGER PRIMARY KEY CHECK (id = 1), |
| 1522 | day TEXT NOT NULL, |
| 1523 | ping_count INTEGER NOT NULL, |
| 1524 | open_count INTEGER NOT NULL, |
| 1525 | checked_at TEXT NOT NULL |
| 1526 | )`, |
| 1527 | ).run(); |
| 1528 | const row = await env.DB.prepare( |
| 1529 | `SELECT date('now') AS day, |
| 1530 | COUNT(*) AS ping_count, |
| 1531 | COALESCE(SUM(opens), 0) AS open_count |
| 1532 | FROM pings |
| 1533 | WHERE date = date('now') AND install_id <> ?1`, |
| 1534 | ) |
| 1535 | .bind(CANARY_INSTALL_ID) |
| 1536 | .first<{ day: string; ping_count: number; open_count: number }>(); |
| 1537 | const day = row?.day ?? ""; |
| 1538 | const pingCount = Number(row?.ping_count ?? 0); |
| 1539 | const openCount = Number(row?.open_count ?? 0); |
| 1540 | const previous = await env.DB.prepare( |
| 1541 | "SELECT day, ping_count, open_count, checked_at FROM ingest_sentinel_state WHERE id = 1", |
| 1542 | ).first<{ day: string; ping_count: number; open_count: number; checked_at: string }>(); |
| 1543 | if (!pingCount) { |
| 1544 | problems.push("no launch pings recorded today (UTC)"); |
| 1545 | } else if ( |
| 1546 | previous?.day === day && |
| 1547 | pingCount <= Number(previous.ping_count) && |
| 1548 | openCount <= Number(previous.open_count) |
| 1549 | ) { |
| 1550 | problems.push( |
| 1551 | `launch ping totals unchanged since ${previous.checked_at} UTC (${pingCount} install rows, ${openCount} opens)`, |
| 1552 | ); |
| 1553 | } |
| 1554 | await env.DB.prepare( |
| 1555 | `INSERT INTO ingest_sentinel_state (id, day, ping_count, open_count, checked_at) |
| 1556 | VALUES (1, ?1, ?2, ?3, datetime('now')) |
| 1557 | ON CONFLICT (id) DO UPDATE SET |
| 1558 | day = ?1, ping_count = ?2, open_count = ?3, checked_at = datetime('now')`, |
| 1559 | ) |
| 1560 | .bind(day, pingCount, openCount) |
| 1561 | .run(); |
| 1562 | } catch (err) { |
| 1563 | problems.push(`ping progress check failed: ${errText(err)}`); |
| 1564 | } |
| 1565 | if (!problems.length) return; |
| 1566 | const message = `crash.reasonix.io ingest sentinel: ${problems.join("; ")} — https://crash.reasonix.io/stats`; |
| 1567 | console.error(message); |
| 1568 | await sendAlert(env, message); |
| 1569 | } |
| 1570 | |
| 1571 | async function purgeExpiredStatsRows(env: Env): Promise<void> { |
| 1572 | try { |
| 1573 | await ensureCLITelemetrySchema(env); |
| 1574 | } catch (err) { |
| 1575 | console.error("retention: CLI telemetry schema unavailable", err); |
| 1576 | } |
| 1577 | for (const { table, keepDays } of RETENTION) { |
| 1578 | // Keep exactly the newest `keepDays` dates: today plus keepDays-1 back, |
| 1579 | // matching the `date >= date('now', '-{keepDays-1} day')` reads. |
| 1580 | const cutoff = `-${keepDays - 1} day`; |
| 1581 | let purged = 0; |
| 1582 | try { |
| 1583 | for (let i = 0; i < RETENTION_MAX_CHUNKS; i++) { |
| 1584 | const res = await env.DB.prepare( |
| 1585 | `DELETE FROM ${table} WHERE rowid IN ( |
| 1586 | SELECT rowid FROM ${table} WHERE date < date('now', ?1) LIMIT ${RETENTION_CHUNK_ROWS} |
| 1587 | )`, |
| 1588 | ) |
| 1589 | .bind(cutoff) |
| 1590 | .run(); |
| 1591 | const changes = res.meta.changes ?? 0; |
| 1592 | purged += changes; |
| 1593 | if (changes < RETENTION_CHUNK_ROWS) break; |
| 1594 | } |
| 1595 | console.log(`retention: purged ${purged} rows from ${table} (keep ${keepDays}d)`); |
| 1596 | } catch (err) { |
| 1597 | // One broken table must not stop the others; the cron retries tomorrow. |
| 1598 | console.error(`retention: purge failed for ${table} after ${purged} rows`, err); |
| 1599 | } |
| 1600 | } |
| 1601 | } |
| 1602 | |
| 1603 | export default { |
| 1604 | async fetch(request: Request, env: Env): Promise<Response> { |
| 1605 | const url = new URL(request.url); |
| 1606 | const path = url.pathname; |
| 1607 | const method = request.method; |
| 1608 | |
| 1609 | const desktopRelease = desktopReleaseChannel(path); |
| 1610 | if (desktopRelease) { |
| 1611 | return handleReleaseGatewayRequest(method, () => handleDesktopReleaseManifest(desktopRelease)); |
| 1612 | } |
| 1613 | const cliRelease = cliReleaseChannel(path); |
| 1614 | if (cliRelease) { |
| 1615 | return handleReleaseGatewayRequest(method, () => handleCLIRelease(cliRelease)); |
| 1616 | } |
| 1617 | |
| 1618 | if (path === "/v1/report" && method === "POST") return handleReport(request, env); |
| 1619 | if (path === "/v1/ping" && method === "POST") return handlePing(request, env); |
| 1620 | if (path === "/v1/metrics" && method === "POST") return handleMetrics(request, env); |
| 1621 | |
| 1622 | // Skill/MCP registry API — the folded Hono app handles its own auth, CORS |
| 1623 | // and rate limiting against the registry database (public reads + publish, |
| 1624 | // plus the JSON /v1/admin the site's moderation panel calls). |
| 1625 | if (path.startsWith("/v1/packages") || path === "/v1/activity" || path.startsWith("/v1/admin")) { |
| 1626 | return registryApp.fetch(request, registryBindings(env)); |
| 1627 | } |
| 1628 | |
| 1629 | const login = loginUrl(env, request); |
| 1630 | |
| 1631 | // Authentication moved to id.reasonix.io; these paths just bounce there. |
| 1632 | if ((path === "/login" || path === "/register") && method === "GET") return redirect(login); |
| 1633 | if (path === "/logout" && method === "POST") return redirect(login, await sharedLogout(request, env)); |
| 1634 | |
| 1635 | const user = await currentUser(request, env); |
| 1636 | |
| 1637 | if (path === "/") return redirect(user ? (atLeast(user.role, "viewer") ? "/stats" : "/account") : login); |
| 1638 | |
| 1639 | if (path === "/account" && method === "GET") return user ? html(renderAccount(user)) : redirect(login); |
| 1640 | |
| 1641 | const groupFingerprint = groupFingerprintFromPath(path); |
| 1642 | const statsModuleMatch = path.match(/^\/stats\/(diagnostics|usage|preferences|health)$/); |
| 1643 | if ((path === "/stats" || statsModuleMatch) && method === "GET") |
| 1644 | return requireViewer(user, login) ?? handleStats(request, env, user as User, (statsModuleMatch?.[1] as StatsModule | undefined) ?? "usage"); |
| 1645 | if (groupFingerprint && method === "GET") return requireViewer(user, login) ?? handleGroup(env, groupFingerprint, user as User); |
| 1646 | if (groupFingerprint && method === "POST") { |
| 1647 | if (user?.role !== "admin") return new Response("forbidden", { status: 403 }); |
| 1648 | return handleGroupAction(request, env, user, groupFingerprint); |
| 1649 | } |
| 1650 | |
| 1651 | if (path === "/admin" && method === "GET") { |
| 1652 | if (!user) return redirect(login); |
| 1653 | return user.role === "admin" ? handleAdminList(env, user) : redirect("/account"); |
| 1654 | } |
| 1655 | if (path === "/admin/audit" && method === "GET") { |
| 1656 | if (!user) return redirect(login); |
| 1657 | return user.role === "admin" ? handleAdminAudit(env, user) : redirect("/account"); |
| 1658 | } |
| 1659 | if (path === "/admin/users" && method === "POST") { |
| 1660 | if (user?.role !== "admin") return new Response("forbidden", { status: 403 }); |
| 1661 | return handleAdminUsers(request, env, user); |
| 1662 | } |
| 1663 | |
| 1664 | if (path === "/community" && method === "GET") { |
| 1665 | if (!user) return redirect(login); |
| 1666 | return user.role === "admin" ? handleCommunityList(env, user, communityStatus(url)) : redirect("/account"); |
| 1667 | } |
| 1668 | const pkgActionMatch = path.match(/^\/community\/([^/]+)\/([^/]+)\/(approve|reject|hide|verify|unverify)$/); |
| 1669 | if (pkgActionMatch && method === "POST") { |
| 1670 | if (user?.role !== "admin") return new Response("forbidden", { status: 403 }); |
| 1671 | return handleCommunityAction(request, env, user, pkgActionMatch[1], pkgActionMatch[2], pkgActionMatch[3]); |
| 1672 | } |
| 1673 | |
| 1674 | if ( |
| 1675 | path === "/v1/report" || |
| 1676 | path === "/v1/ping" || |
| 1677 | path === "/v1/metrics" || |
| 1678 | path === "/login" || |
| 1679 | path === "/register" || |
| 1680 | path === "/logout" || |
| 1681 | path === "/account" || |
| 1682 | path.startsWith("/stats") || |
| 1683 | path.startsWith("/admin") || |
| 1684 | path.startsWith("/community") |
| 1685 | ) { |
| 1686 | return new Response("method not allowed", { status: 405 }); |
| 1687 | } |
| 1688 | return new Response("not found", { status: 404 }); |
| 1689 | }, |
| 1690 | |
| 1691 | async scheduled(controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> { |
| 1692 | if (controller.cron === SENTINEL_CRON) { |
| 1693 | ctx.waitUntil(runIngestSentinel(env)); |
| 1694 | return; |
| 1695 | } |
| 1696 | if (controller.cron === ROLLUP_CRON) { |
| 1697 | ctx.waitUntil(refreshMetricUserRollup(env)); |
| 1698 | return; |
| 1699 | } |
| 1700 | ctx.waitUntil(purgeExpiredStatsRows(env)); |
| 1701 | }, |
| 1702 | }; |
| 1703 |