From 06fe1f4ed796b3a68ba565912f90a6f488fa5ce3 Mon Sep 17 00:00:00 2001 From: yohlo Date: Tue, 25 Aug 2026 21:38:08 -0700 Subject: [PATCH] feat(telemetry): rollups, alerts, scheduler, health checks, admin server fns - pure rollup accumulators with tests; daily upserts keyed date/metric/dim - alert checks with db-backed cooldown and push to admin players - in-process scheduler started from health probes - health endpoint checks pb + supertokens, stale-while-revalidate, always 200 - sse connection counter moved to emitter - telemetry feature server fns: events/errors/groups/rollups/alerts/status --- src/app/routes/api/events.$.ts | 17 ++- src/app/routes/api/health.ts | 16 ++- src/features/telemetry/queries.ts | 76 +++++++++++ src/features/telemetry/server.ts | 171 +++++++++++++++++++++++++ src/features/telemetry/types.ts | 38 ++++++ src/lib/events/emitter.ts | 14 +++ src/lib/telemetry/alerts.server.ts | 170 +++++++++++++++++++++++++ src/lib/telemetry/health.server.ts | 56 +++++++++ src/lib/telemetry/rollup-core.test.ts | 113 +++++++++++++++++ src/lib/telemetry/rollup-core.ts | 175 ++++++++++++++++++++++++++ src/lib/telemetry/rollup.server.ts | 103 +++++++++++++++ src/lib/telemetry/scheduler.server.ts | 71 +++++++++++ 12 files changed, 1009 insertions(+), 11 deletions(-) create mode 100644 src/features/telemetry/queries.ts create mode 100644 src/features/telemetry/server.ts create mode 100644 src/features/telemetry/types.ts create mode 100644 src/lib/telemetry/alerts.server.ts create mode 100644 src/lib/telemetry/health.server.ts create mode 100644 src/lib/telemetry/rollup-core.test.ts create mode 100644 src/lib/telemetry/rollup-core.ts create mode 100644 src/lib/telemetry/rollup.server.ts create mode 100644 src/lib/telemetry/scheduler.server.ts diff --git a/src/app/routes/api/events.$.ts b/src/app/routes/api/events.$.ts index ff488fa..5fa0689 100644 --- a/src/app/routes/api/events.$.ts +++ b/src/app/routes/api/events.$.ts @@ -1,9 +1,14 @@ import { createFileRoute } from "@tanstack/react-router"; -import { serverEvents, EVENT_TYPES, type ServerEvent } from "@/lib/events/emitter"; +import { + serverEvents, + EVENT_TYPES, + incrementConnections, + decrementConnections, + type ServerEvent, +} from "@/lib/events/emitter"; import { logger } from "@/lib/logger"; import { superTokensRequestMiddleware } from "@/utils/supertokens"; -let activeConnections = 0; const encoder = new TextEncoder(); export const Route = createFileRoute("/api/events/$")({ @@ -11,9 +16,9 @@ export const Route = createFileRoute("/api/events/$")({ middleware: [superTokensRequestMiddleware], handlers: { GET: ({ request }) => { - activeConnections++; + const active = incrementConnections(); const connectionId = `conn_${Date.now()}_${Math.random().toString(36).substring(2, 11)}`; - logger.info(`ServerEvents | New connection ${connectionId}. Active: ${activeConnections}`); + logger.info(`ServerEvents | New connection ${connectionId}. Active: ${active}`); let cleanedUp = false; let cleanup = () => {}; @@ -41,12 +46,12 @@ export const Route = createFileRoute("/api/events/$")({ cleanup = () => { if (cleanedUp) return; cleanedUp = true; - activeConnections--; + const remaining = decrementConnections(); for (const type of EVENT_TYPES) { serverEvents.off(type, handleEvent); } clearInterval(pingInterval); - logger.info(`ServerEvents | Connection ${connectionId} cleanup completed. Active: ${activeConnections}`); + logger.info(`ServerEvents | Connection ${connectionId} cleanup completed. Active: ${remaining}`); }; request.signal?.addEventListener("abort", cleanup); diff --git a/src/app/routes/api/health.ts b/src/app/routes/api/health.ts index 38c7e5a..0226fed 100644 --- a/src/app/routes/api/health.ts +++ b/src/app/routes/api/health.ts @@ -1,19 +1,25 @@ import { createFileRoute } from "@tanstack/react-router"; +import { ensureTelemetryScheduler } from "@/lib/telemetry/scheduler.server"; +import { getHealthSnapshot } from "@/lib/telemetry/health.server"; export const Route = createFileRoute("/api/health")({ server: { handlers: { - GET: () => { + GET: async () => { + ensureTelemetryScheduler(); + const snapshot = await getHealthSnapshot(); + + // Always 200: a PB/SuperTokens outage must not fail the k8s probes; + // restarting the app pod cannot fix a dependency. return new Response( JSON.stringify({ - status: "ok", + status: snapshot.status, + checks: snapshot.checks, timestamp: new Date().toISOString(), }), { status: 200, - headers: { - "Content-Type": "application/json", - }, + headers: { "Content-Type": "application/json" }, } ); }, diff --git a/src/features/telemetry/queries.ts b/src/features/telemetry/queries.ts new file mode 100644 index 0000000..6f41bed --- /dev/null +++ b/src/features/telemetry/queries.ts @@ -0,0 +1,76 @@ +import { useServerSuspenseQuery } from "@/lib/tanstack-query/hooks"; +import { + getClientErrorGroups, + getRollupRange, + getTelemetryRuntimeStatus, + listTelemetryAlerts, + searchClientErrors, + searchClientEvents, +} from "./server"; +import type { ClientErrorSearchParams, ClientEventSearchParams } from "./types"; + +interface ErrorGroupParams { + from: string; + to: string; + includeResolved?: boolean; + limit?: number; +} + +interface RollupRangeParams { + metrics: string[]; + from: string; + to: string; + dim?: string; +} + +export const telemetryKeys = { + all: ["telemetry"] as const, + events: (params: ClientEventSearchParams) => ["telemetry", "events", params] as const, + errors: (params: ClientErrorSearchParams) => ["telemetry", "errors", params] as const, + errorGroups: (params: ErrorGroupParams) => ["telemetry", "errorGroups", params] as const, + rollups: (params: RollupRangeParams) => ["telemetry", "rollups", params] as const, + alerts: (page: number) => ["telemetry", "alerts", page] as const, + runtimeStatus: ["telemetry", "runtimeStatus"] as const, +}; + +export const telemetryQueries = { + events: (params: ClientEventSearchParams = {}) => ({ + queryKey: telemetryKeys.events(params), + queryFn: () => searchClientEvents({ data: params }), + }), + errors: (params: ClientErrorSearchParams = {}) => ({ + queryKey: telemetryKeys.errors(params), + queryFn: () => searchClientErrors({ data: params }), + }), + errorGroups: (params: ErrorGroupParams) => ({ + queryKey: telemetryKeys.errorGroups(params), + queryFn: () => getClientErrorGroups({ data: params }), + }), + rollups: (params: RollupRangeParams) => ({ + queryKey: telemetryKeys.rollups(params), + queryFn: () => getRollupRange({ data: params }), + }), + alerts: (page = 1) => ({ + queryKey: telemetryKeys.alerts(page), + queryFn: () => listTelemetryAlerts({ data: { page } }), + }), + runtimeStatus: () => ({ + queryKey: telemetryKeys.runtimeStatus, + queryFn: () => getTelemetryRuntimeStatus(), + }), +}; + +export const useClientEvents = (params: ClientEventSearchParams = {}) => + useServerSuspenseQuery(telemetryQueries.events(params)); + +export const useClientErrors = (params: ClientErrorSearchParams = {}) => + useServerSuspenseQuery(telemetryQueries.errors(params)); + +export const useClientErrorGroups = (params: ErrorGroupParams) => + useServerSuspenseQuery(telemetryQueries.errorGroups(params)); + +export const useRollupRange = (params: RollupRangeParams) => + useServerSuspenseQuery(telemetryQueries.rollups(params)); + +export const useTelemetryRuntimeStatus = () => + useServerSuspenseQuery(telemetryQueries.runtimeStatus()); diff --git a/src/features/telemetry/server.ts b/src/features/telemetry/server.ts new file mode 100644 index 0000000..a94eb54 --- /dev/null +++ b/src/features/telemetry/server.ts @@ -0,0 +1,171 @@ +import { createServerFn } from "@tanstack/react-start"; +import { z } from "zod"; +import { superTokensAdminFunctionMiddleware } from "@/utils/supertokens"; +import { pbAdmin } from "@/lib/pocketbase/client"; +import { toServerResult } from "@/lib/tanstack-query/utils/to-server-result"; +import { + transformClientError, + transformClientEvent, +} from "@/lib/pocketbase/util/transform-types"; +import { getSchedulerStatus } from "@/lib/telemetry/scheduler.server"; +import { getHealthSnapshot } from "@/lib/telemetry/health.server"; +import { getActiveConnectionCount } from "@/lib/events/emitter"; +import type { + ClientError, + ClientErrorGroup, + ClientEvent, + RollupRow, + TelemetryAlert, + TelemetryListResult, + TelemetryRuntimeStatus, +} from "./types"; + +const clientEventSearchSchema = z.object({ + page: z.number().optional(), + perPage: z.number().max(200).optional(), + name: z.string().optional(), + player: z.string().optional(), + sessionId: z.string().optional(), + path: z.string().optional(), + from: z.string().optional(), + to: z.string().optional(), + sortBy: z.string().optional(), +}); + +export const searchClientEvents = createServerFn() + .validator(clientEventSearchSchema) + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async ({ data }) => + toServerResult>(async () => { + const result = await pbAdmin.searchClientEvents(data); + return { ...result, items: result.items.map(transformClientEvent) }; + }) + ); + +const clientErrorSearchSchema = z.object({ + page: z.number().optional(), + perPage: z.number().max(200).optional(), + groupHash: z.string().optional(), + resolved: z.boolean().optional(), + player: z.string().optional(), + from: z.string().optional(), + to: z.string().optional(), + sortBy: z.string().optional(), +}); + +export const searchClientErrors = createServerFn() + .validator(clientErrorSearchSchema) + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async ({ data }) => + toServerResult>(async () => { + const result = await pbAdmin.searchClientErrors(data); + return { ...result, items: result.items.map(transformClientError) }; + }) + ); + +const errorGroupsSchema = z.object({ + from: z.string(), + to: z.string(), + includeResolved: z.boolean().optional(), + limit: z.number().max(100).optional(), +}); + +export const getClientErrorGroups = createServerFn() + .validator(errorGroupsSchema) + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async ({ data }) => + toServerResult(async () => { + const { from, to, includeResolved = false, limit = 50 } = data; + + const rollups = await pbAdmin.getRollupRange({ + metrics: ["client_error.count"], + from, + to, + }); + + const countByGroup = new Map(); + for (const rollup of rollups) { + if (!rollup.dim) continue; + countByGroup.set(rollup.dim, (countByGroup.get(rollup.dim) ?? 0) + rollup.value); + } + + const topGroups = [...countByGroup.entries()] + .sort((a, b) => b[1] - a[1]) + .slice(0, limit); + + const groups: ClientErrorGroup[] = []; + for (const [groupHash, count] of topGroups) { + const latest = await pbAdmin.searchClientErrors({ groupHash, perPage: 1 }); + const sample = latest.items[0]; + if (!sample) continue; + + // A group counts as resolved when its most recent occurrence is + // resolved; a recurrence arrives resolved=false and surfaces again. + const resolved = !!sample.resolved; + if (resolved && !includeResolved) continue; + + groups.push({ + group_hash: groupHash, + count, + lastSeen: sample.created, + resolved, + sample: { + message: sample.message, + stack: sample.stack, + path: sample.path, + source: sample.source, + }, + }); + } + + return groups; + }) + ); + +export const resolveClientErrorGroup = createServerFn({ method: "POST" }) + .validator(z.object({ groupHash: z.string(), resolved: z.boolean().optional() })) + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async ({ data }) => + toServerResult<{ updated: number }>(async () => { + const updated = await pbAdmin.resolveErrorsByGroup(data.groupHash, data.resolved ?? true); + return { updated }; + }) + ); + +const rollupRangeSchema = z.object({ + metrics: z.array(z.string()).min(1).max(8), + from: z.string().regex(/^\d{4}-\d{2}-\d{2}$/), + to: z.string().regex(/^\d{4}-\d{2}-\d{2}$/), + dim: z.string().optional(), +}); + +export const getRollupRange = createServerFn() + .validator(rollupRangeSchema) + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async ({ data }) => + toServerResult(async () => pbAdmin.getRollupRange(data)) + ); + +export const listTelemetryAlerts = createServerFn() + .validator( + z.object({ page: z.number().optional(), perPage: z.number().max(200).optional() }) + ) + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async ({ data }) => + toServerResult>(async () => + pbAdmin.listTelemetryAlerts(data.page, data.perPage) + ) + ); + +export const getTelemetryRuntimeStatus = createServerFn() + .middleware([superTokensAdminFunctionMiddleware]) + .handler(async () => + toServerResult(async () => { + const health = await getHealthSnapshot(); + return { + health: { status: health.status, checks: health.checks }, + scheduler: getSchedulerStatus(), + activeSseConnections: getActiveConnectionCount(), + }; + }) + ); diff --git a/src/features/telemetry/types.ts b/src/features/telemetry/types.ts new file mode 100644 index 0000000..716f077 --- /dev/null +++ b/src/features/telemetry/types.ts @@ -0,0 +1,38 @@ +export type { + ClientEventRecord as ClientEvent, + ClientErrorRecord as ClientError, + ClientErrorSource, + ClientEventSearchParams, + ClientErrorSearchParams, + RollupRecord as RollupRow, + TelemetryAlertRecord as TelemetryAlert, + TelemetryListResult, +} from "@/lib/pocketbase/services/telemetry"; + +export interface ClientErrorGroup { + group_hash: string; + count: number; + lastSeen: string; + resolved: boolean; + sample: { + message: string; + stack?: string; + path?: string; + source?: string; + }; +} + +export interface TelemetryRuntimeStatus { + health: { + status: "ok" | "degraded"; + checks: { pocketbase: string; supertokens: string }; + }; + scheduler: { + running: boolean; + lastRollupAt?: string; + lastAlertCheckAt?: string; + lastRollupError?: string; + lastAlertError?: string; + }; + activeSseConnections: number; +} diff --git a/src/lib/events/emitter.ts b/src/lib/events/emitter.ts index 06ba7ae..4772a4b 100644 --- a/src/lib/events/emitter.ts +++ b/src/lib/events/emitter.ts @@ -68,6 +68,20 @@ export const emitServerEvent = (event: ServerEvent) => { serverEvents.emit(event.type, event); }; +let activeConnections = 0; + +export const incrementConnections = () => { + activeConnections++; + return activeConnections; +}; + +export const decrementConnections = () => { + activeConnections = Math.max(0, activeConnections - 1); + return activeConnections; +}; + +export const getActiveConnectionCount = () => activeConnections; + if (process.env.NODE_ENV === 'development') { setInterval(() => { const listenerCounts = Object.fromEntries( diff --git a/src/lib/telemetry/alerts.server.ts b/src/lib/telemetry/alerts.server.ts new file mode 100644 index 0000000..eb0863b --- /dev/null +++ b/src/lib/telemetry/alerts.server.ts @@ -0,0 +1,170 @@ +import { msg } from "@lingui/core/macro"; +import { pbAdmin } from "@/lib/pocketbase/client"; +import { sendPushToPlayer } from "@/lib/push"; +import { Logger } from "@/lib/logger"; +import { ADMIN_ROLE } from "@/features/core/utils/roles"; + +const logger = new Logger("Telemetry > Alerts"); + +const WINDOW_MS = 15 * 60 * 1000; +const COOLDOWN_MS = 60 * 60 * 1000; +const MAX_NEW_GROUPS_PER_RUN = 20; + +const errorSpikeThreshold = (): number => { + const raw = Number(process.env.TELEMETRY_ALERT_ERROR_COUNT); + return Number.isFinite(raw) && raw > 0 ? raw : 10; +}; + +const fnFailureThreshold = (): number => { + const raw = Number(process.env.TELEMETRY_ALERT_FN_FAILURES); + return Number.isFinite(raw) && raw > 0 ? raw : 5; +}; + +const pbTimestamp = (date: Date): string => date.toISOString().replace("T", " "); + +const ADMIN_CACHE_TTL_MS = 10 * 60 * 1000; +let adminPlayersCache: { playerIds: string[]; expiresAt: number } | null = null; + +const getAdminPlayerIds = async (): Promise => { + const now = Date.now(); + if (adminPlayersCache && adminPlayersCache.expiresAt > now) { + return adminPlayersCache.playerIds; + } + + const UserRoles = (await import("supertokens-node/recipe/userroles")).default; + const response = await UserRoles.getUsersThatHaveRole("public", ADMIN_ROLE); + const users = response.status === "OK" ? response.users : []; + const playerIds: string[] = []; + for (const userId of users) { + try { + const player = await pbAdmin.getPlayerByAuthId(userId); + if (player) playerIds.push(player.id); + } catch {} + } + + adminPlayersCache = { playerIds, expiresAt: now + ADMIN_CACHE_TTL_MS }; + return playerIds; +}; + +interface AlertCandidate { + kind: string; + dim: string; + message: string; + value: number; + title: Parameters[1]["title"]; + body: Parameters[1]["body"]; +} + +const fireAlert = async (candidate: AlertCandidate) => { + const since = new Date(Date.now() - COOLDOWN_MS).toISOString(); + const recent = await pbAdmin.findRecentAlert(candidate.kind, candidate.dim, pbTimestamp(new Date(since))); + if (recent) return; + + await pbAdmin.createTelemetryAlert({ + kind: candidate.kind, + dim: candidate.dim, + message: candidate.message, + value: candidate.value, + }); + + const admins = await getAdminPlayerIds(); + for (const playerId of admins) { + try { + await sendPushToPlayer(playerId, { + title: candidate.title, + body: candidate.body, + url: "/admin/telemetry/errors", + tag: `telemetry-${candidate.kind}`, + }); + } catch (error) { + logger.error("Failed to push telemetry alert", error); + } + } + + logger.info(`Alert fired: ${candidate.kind} (${candidate.dim || "-"}) = ${candidate.value}`); +}; + +export const runAlertChecks = async (): Promise => { + await pbAdmin.authPromise; + const windowStart = pbTimestamp(new Date(Date.now() - WINDOW_MS)); + + const recentErrors = await pbAdmin.searchClientErrors({ + from: windowStart, + perPage: 200, + sortBy: "-created", + }); + + if (recentErrors.totalItems >= errorSpikeThreshold()) { + await fireAlert({ + kind: "client_error_spike", + dim: "", + message: `${recentErrors.totalItems} client errors in 15 minutes`, + value: recentErrors.totalItems, + title: msg`Client error spike`, + body: msg`${recentErrors.totalItems} client errors in the last 15 minutes`, + }); + } + + const groups = new Map(); + for (const error of recentErrors.items) { + const entry = groups.get(error.group_hash) ?? { count: 0, message: error.message }; + entry.count++; + groups.set(error.group_hash, entry); + } + let checked = 0; + for (const [groupHash, info] of groups) { + if (checked++ >= MAX_NEW_GROUPS_PER_RUN) break; + const older = await pbAdmin.searchClientErrors({ groupHash, to: windowStart, perPage: 1 }); + if (older.totalItems === 0) { + await fireAlert({ + kind: "client_error_new_group", + dim: groupHash, + message: info.message, + value: info.count, + title: msg`New client error`, + body: msg`A new error appeared: ${info.message.slice(0, 120)}`, + }); + } + } + + const failedActivities = await pbAdmin.searchActivities({ + success: false, + perPage: 200, + sortBy: "-created", + }); + const failuresByFn = new Map(); + for (const activity of failedActivities.items) { + if (activity.created < windowStart) continue; + if (activity.error?.startsWith("FORBIDDEN")) { + failuresByFn.set("__denied__", (failuresByFn.get("__denied__") ?? 0) + 1); + continue; + } + failuresByFn.set(activity.name, (failuresByFn.get(activity.name) ?? 0) + 1); + } + + for (const [name, count] of failuresByFn) { + if (name === "__denied__") continue; + if (count >= fnFailureThreshold()) { + await fireAlert({ + kind: "server_fn_failures", + dim: name, + message: `${name} failed ${count} times in 15 minutes`, + value: count, + title: msg`Server function failing`, + body: msg`${name} failed ${count} times in the last 15 minutes`, + }); + } + } + + const deniedCount = failuresByFn.get("__denied__") ?? 0; + if (deniedCount >= 1) { + await fireAlert({ + kind: "denied_access", + dim: "", + message: `${deniedCount} denied admin access attempts in 15 minutes`, + value: deniedCount, + title: msg`Denied access attempt`, + body: msg`${deniedCount} denied admin access attempts in the last 15 minutes`, + }); + } +}; diff --git a/src/lib/telemetry/health.server.ts b/src/lib/telemetry/health.server.ts new file mode 100644 index 0000000..79fad37 --- /dev/null +++ b/src/lib/telemetry/health.server.ts @@ -0,0 +1,56 @@ +export interface HealthChecks { + pocketbase: "ok" | "down"; + supertokens: "ok" | "down"; +} + +export interface HealthSnapshot { + status: "ok" | "degraded"; + checks: HealthChecks; + checkedAt: string; +} + +const CHECK_TTL_MS = 10_000; +const CHECK_TIMEOUT_MS = 2_000; + +let cached: HealthSnapshot | null = null; +let refreshing: Promise | null = null; + +const ping = async (url: string): Promise => { + try { + const response = await fetch(url, { signal: AbortSignal.timeout(CHECK_TIMEOUT_MS) }); + return response.ok; + } catch { + return false; + } +}; + +const runChecks = async (): Promise => { + const [pocketbase, supertokens] = await Promise.all([ + ping(`${process.env.POCKETBASE_URL}/api/health`), + ping(`${process.env.SUPERTOKENS_URI}/hello`), + ]); + + cached = { + status: pocketbase && supertokens ? "ok" : "degraded", + checks: { + pocketbase: pocketbase ? "ok" : "down", + supertokens: supertokens ? "ok" : "down", + }, + checkedAt: new Date().toISOString(), + }; + return cached; +}; + +// Stale-while-revalidate: k8s probes get the cached snapshot instantly while a +// refresh runs in the background; only the very first request awaits a check. +export const getHealthSnapshot = async (): Promise => { + if (cached && Date.now() - Date.parse(cached.checkedAt) < CHECK_TTL_MS) { + return cached; + } + if (!refreshing) { + refreshing = runChecks().finally(() => { + refreshing = null; + }); + } + return cached ?? refreshing; +}; diff --git a/src/lib/telemetry/rollup-core.test.ts b/src/lib/telemetry/rollup-core.test.ts new file mode 100644 index 0000000..1267578 --- /dev/null +++ b/src/lib/telemetry/rollup-core.test.ts @@ -0,0 +1,113 @@ +import { describe, expect, it } from "vitest"; +import { + accumulateActivities, + accumulateClientErrors, + accumulateClientEvents, + buildDayRollups, + newActivityAccumulator, + newClientErrorAccumulator, + newClientEventAccumulator, + percentile, +} from "./rollup-core"; + +describe("percentile", () => { + it("handles empty, single, even, odd inputs", () => { + expect(percentile([], 50)).toBe(0); + expect(percentile([42], 50)).toBe(42); + expect(percentile([1, 2, 3, 4], 50)).toBe(2); + expect(percentile([1, 2, 3, 4, 5], 50)).toBe(3); + expect(percentile([10, 20, 30, 40, 50, 60, 70, 80, 90, 100], 95)).toBe(100); + }); + + it("does not mutate its input", () => { + const values = [3, 1, 2]; + percentile(values, 50); + expect(values).toEqual([3, 1, 2]); + }); +}); + +const rollupFor = (rollups: ReturnType, metric: string, dim = "") => + rollups.find((r) => r.metric === metric && r.dim === dim); + +describe("buildDayRollups", () => { + it("aggregates a representative day", () => { + const activities = newActivityAccumulator(); + accumulateActivities(activities, [ + { name: "updatePlayer", player: "p1", duration: 100, success: true }, + { name: "updatePlayer", player: "p2", duration: 300, success: true }, + { name: "updatePlayer", player: "p1", duration: 200, success: false, error: "DATABASE: nope" }, + { name: "endMatch", player: "p3", duration: 50, success: true }, + { name: "endMatch", player: "p1", duration: 0, success: false, error: "FORBIDDEN: Access denied" }, + { name: "auth.otp_sent", success: true }, + { name: "auth.otp_failed", success: false, error: "INCORRECT_USER_INPUT_CODE_ERROR" }, + ]); + + const clientEvents = newClientEventAccumulator(); + accumulateClientEvents(clientEvents, [ + { name: "page_view", player: "p1", session_id: "s1", route_id: "/home" }, + { name: "page_view", player: "p1", session_id: "s1", route_id: "/home" }, + { name: "page_view", player: "p4", session_id: "s2", route_id: "/tournaments" }, + { name: "vital.LCP", player: "p1", session_id: "s1", value: 1200 }, + { name: "vital.LCP", player: "p4", session_id: "s2", value: 2400 }, + { name: "predictions_submitted", player: "p4", session_id: "s2" }, + ]); + + const clientErrors = newClientErrorAccumulator(); + accumulateClientErrors(clientErrors, [ + { group_hash: "aaaa1111" }, + { group_hash: "aaaa1111" }, + { group_hash: "bbbb2222" }, + ]); + + const rollups = buildDayRollups("2026-08-25", { + activities, + clientEvents, + clientErrors, + wau: 7, + mau: 12, + partial: false, + }); + + expect(rollupFor(rollups, "dau")?.value).toBe(4); + expect(rollupFor(rollups, "wau")?.value).toBe(7); + expect(rollupFor(rollups, "mau")?.value).toBe(12); + + expect(rollupFor(rollups, "server_fn.count", "updatePlayer")?.value).toBe(3); + expect(rollupFor(rollups, "server_fn.error_count", "updatePlayer")?.value).toBe(1); + expect(rollupFor(rollups, "server_fn.duration.p50", "updatePlayer")?.value).toBe(200); + expect(rollupFor(rollups, "server_fn.count", "auth.otp_sent")).toBeUndefined(); + + expect(rollupFor(rollups, "page_view.count", "/home")?.value).toBe(2); + expect(rollupFor(rollups, "page_view.count", "")?.value).toBe(3); + expect(rollupFor(rollups, "page_view.sessions")?.value).toBe(2); + + expect(rollupFor(rollups, "client_error.count", "")?.value).toBe(3); + expect(rollupFor(rollups, "client_error.count", "aaaa1111")?.value).toBe(2); + + expect(rollupFor(rollups, "vital.p75", "LCP")?.value).toBe(2400); + expect(rollupFor(rollups, "denied.count")?.value).toBe(1); + expect(rollupFor(rollups, "auth.otp_sent.count")?.value).toBe(1); + expect(rollupFor(rollups, "auth.otp_failed.count")?.value).toBe(1); + + for (const rollup of rollups) { + expect(rollup.dim).toBeDefined(); + expect(rollup.meta?.partial).toBeUndefined(); + } + }); + + it("marks partial days and is deterministic across re-runs", () => { + const build = () => + buildDayRollups("2026-08-25", { + activities: newActivityAccumulator(), + clientEvents: newClientEventAccumulator(), + clientErrors: newClientErrorAccumulator(), + wau: 0, + mau: 0, + partial: true, + }); + + const first = build(); + expect(first.every((r) => r.meta?.partial === true)).toBe(true); + expect(build()).toEqual(first); + }); +}); diff --git a/src/lib/telemetry/rollup-core.ts b/src/lib/telemetry/rollup-core.ts new file mode 100644 index 0000000..af1d5b8 --- /dev/null +++ b/src/lib/telemetry/rollup-core.ts @@ -0,0 +1,175 @@ +import type { RollupInput } from "@/lib/pocketbase/services/telemetry"; + +export interface ActivityRow { + name: string; + player?: string; + duration?: number; + success?: boolean; + error?: string; +} + +export interface ClientEventRow { + name: string; + player?: string; + session_id?: string; + route_id?: string; + value?: number; +} + +export interface ClientErrorRow { + group_hash: string; +} + +// Nearest-rank percentile; values need not be pre-sorted. +export const percentile = (values: number[], p: number): number => { + if (values.length === 0) return 0; + const sorted = [...values].sort((a, b) => a - b); + const rank = Math.ceil((p / 100) * sorted.length); + return sorted[Math.min(sorted.length, Math.max(1, rank)) - 1]; +}; + +export interface ActivityAccumulator { + players: Set; + perFn: Map; + denied: number; + otpSent: number; + otpFailed: number; +} + +export const newActivityAccumulator = (): ActivityAccumulator => ({ + players: new Set(), + perFn: new Map(), + denied: 0, + otpSent: 0, + otpFailed: 0, +}); + +export const accumulateActivities = (acc: ActivityAccumulator, rows: ActivityRow[]) => { + for (const row of rows) { + if (row.player) acc.players.add(row.player); + if (row.error?.startsWith("FORBIDDEN")) acc.denied++; + if (row.name === "auth.otp_sent") acc.otpSent++; + if (row.name === "auth.otp_failed") acc.otpFailed++; + if (row.name.startsWith("auth.")) continue; + + const fn = acc.perFn.get(row.name) ?? { count: 0, errors: 0, durations: [] }; + fn.count++; + if (row.success === false) fn.errors++; + if (typeof row.duration === "number" && row.duration > 0) fn.durations.push(row.duration); + acc.perFn.set(row.name, fn); + } +}; + +export interface ClientEventAccumulator { + players: Set; + sessions: Set; + pageViewsByRoute: Map; + vitals: Map; +} + +export const newClientEventAccumulator = (): ClientEventAccumulator => ({ + players: new Set(), + sessions: new Set(), + pageViewsByRoute: new Map(), + vitals: new Map(), +}); + +export const accumulateClientEvents = (acc: ClientEventAccumulator, rows: ClientEventRow[]) => { + for (const row of rows) { + if (row.player) acc.players.add(row.player); + if (row.session_id) acc.sessions.add(row.session_id); + + if (row.name === "page_view") { + const route = row.route_id || "(unknown)"; + acc.pageViewsByRoute.set(route, (acc.pageViewsByRoute.get(route) ?? 0) + 1); + } else if (row.name.startsWith("vital.") && typeof row.value === "number") { + const metric = row.name.slice("vital.".length); + const values = acc.vitals.get(metric) ?? []; + values.push(row.value); + acc.vitals.set(metric, values); + } + } +}; + +export interface ClientErrorAccumulator { + byGroup: Map; + total: number; +} + +export const newClientErrorAccumulator = (): ClientErrorAccumulator => ({ + byGroup: new Map(), + total: 0, +}); + +export const accumulateClientErrors = (acc: ClientErrorAccumulator, rows: ClientErrorRow[]) => { + for (const row of rows) { + acc.total++; + acc.byGroup.set(row.group_hash, (acc.byGroup.get(row.group_hash) ?? 0) + 1); + } +}; + +export const buildDayRollups = ( + date: string, + input: { + activities: ActivityAccumulator; + clientEvents: ClientEventAccumulator; + clientErrors: ClientErrorAccumulator; + wau: number; + mau: number; + partial: boolean; + } +): RollupInput[] => { + const { activities, clientEvents, clientErrors, wau, mau, partial } = input; + const meta = partial ? { partial: true } : undefined; + const rollups: RollupInput[] = []; + const push = (metric: string, dim: string, value: number, extraMeta?: Record) => { + rollups.push({ + date, + metric, + dim, + value, + meta: extraMeta ? { ...meta, ...extraMeta } : meta, + }); + }; + + const dayPlayers = new Set([...activities.players, ...clientEvents.players]); + push("dau", "", dayPlayers.size); + push("wau", "", wau); + push("mau", "", mau); + + for (const [name, fn] of activities.perFn) { + push("server_fn.count", name, fn.count); + if (fn.errors > 0) push("server_fn.error_count", name, fn.errors); + if (fn.durations.length > 0) { + push("server_fn.duration.p50", name, percentile(fn.durations, 50), { + samples: fn.durations.length, + }); + push("server_fn.duration.p95", name, percentile(fn.durations, 95), { + samples: fn.durations.length, + }); + } + } + + let totalPageViews = 0; + for (const [route, count] of clientEvents.pageViewsByRoute) { + push("page_view.count", route, count); + totalPageViews += count; + } + push("page_view.count", "", totalPageViews); + push("page_view.sessions", "", clientEvents.sessions.size); + + push("client_error.count", "", clientErrors.total); + for (const [group, count] of clientErrors.byGroup) { + push("client_error.count", group, count); + } + + for (const [metric, values] of clientEvents.vitals) { + push("vital.p75", metric, percentile(values, 75), { samples: values.length }); + } + + push("denied.count", "", activities.denied); + push("auth.otp_sent.count", "", activities.otpSent); + push("auth.otp_failed.count", "", activities.otpFailed); + + return rollups; +}; diff --git a/src/lib/telemetry/rollup.server.ts b/src/lib/telemetry/rollup.server.ts new file mode 100644 index 0000000..f5de7b6 --- /dev/null +++ b/src/lib/telemetry/rollup.server.ts @@ -0,0 +1,103 @@ +import { pbAdmin } from "@/lib/pocketbase/client"; +import { Logger } from "@/lib/logger"; +import { + accumulateActivities, + accumulateClientErrors, + accumulateClientEvents, + buildDayRollups, + newActivityAccumulator, + newClientErrorAccumulator, + newClientEventAccumulator, + type ActivityRow, + type ClientErrorRow, + type ClientEventRow, +} from "./rollup-core"; + +const logger = new Logger("Telemetry > Rollup"); + +export const utcDay = (date: Date): string => date.toISOString().slice(0, 10); + +const addDays = (dateStr: string, days: number): string => { + const d = new Date(`${dateStr}T00:00:00.000Z`); + d.setUTCDate(d.getUTCDate() + days); + return utcDay(d); +}; + +// PB stores autodates as "YYYY-MM-DD HH:MM:SS.sssZ" UTC; boundaries in the +// same shape compare correctly in filters. +const bound = (dateStr: string): string => `${dateStr} 00:00:00.000Z`; + +const dayFilter = (dateStr: string): string => + `created >= '${bound(dateStr)}' && created < '${bound(addDays(dateStr, 1))}'`; + +const distinctPlayersInWindow = async (fromDate: string, toDateExclusive: string): Promise => { + const players = new Set(); + const filter = `created >= '${bound(fromDate)}' && created < '${bound(toDateExclusive)}' && player != ''`; + + await pbAdmin.pageCollection<{ player?: string }>("activities", { filter, fields: "player,created" }, (rows) => { + for (const row of rows) if (row.player) players.add(row.player); + }); + await pbAdmin.pageCollection<{ player?: string }>("client_events", { filter, fields: "player,created" }, (rows) => { + for (const row of rows) if (row.player) players.add(row.player); + }); + + return players.size; +}; + +export const runRollups = async (dateStr: string): Promise => { + await pbAdmin.authPromise; + + const activities = newActivityAccumulator(); + const clientEvents = newClientEventAccumulator(); + const clientErrors = newClientErrorAccumulator(); + + await pbAdmin.pageCollection( + "activities", + { filter: dayFilter(dateStr), fields: "name,player,duration,success,error,created" }, + (rows) => accumulateActivities(activities, rows) + ); + await pbAdmin.pageCollection( + "client_events", + { filter: dayFilter(dateStr), fields: "name,player,session_id,route_id,value,created" }, + (rows) => accumulateClientEvents(clientEvents, rows) + ); + await pbAdmin.pageCollection( + "client_errors", + { filter: dayFilter(dateStr), fields: "group_hash,created" }, + (rows) => accumulateClientErrors(clientErrors, rows) + ); + + const nextDay = addDays(dateStr, 1); + const [wau, mau] = await Promise.all([ + distinctPlayersInWindow(addDays(dateStr, -6), nextDay), + distinctPlayersInWindow(addDays(dateStr, -29), nextDay), + ]); + + const partial = dateStr >= utcDay(new Date()); + const rollups = buildDayRollups(dateStr, { + activities, + clientEvents, + clientErrors, + wau, + mau, + partial, + }); + + for (const rollup of rollups) { + await pbAdmin.upsertRollup(rollup); + } + + logger.info(`Rolled up ${dateStr}: ${rollups.length} metrics${partial ? " (partial)" : ""}`); + return rollups.length; +}; + +// The first run after UTC midnight recomputes yesterday once more so its +// partial rows become final. Idempotent via the (date, metric, dim) key. +export const runScheduledRollups = async (lastRunDay?: string): Promise => { + const today = utcDay(new Date()); + if (lastRunDay && lastRunDay !== today) { + await runRollups(lastRunDay); + } + await runRollups(today); + return today; +}; diff --git a/src/lib/telemetry/scheduler.server.ts b/src/lib/telemetry/scheduler.server.ts new file mode 100644 index 0000000..57649ba --- /dev/null +++ b/src/lib/telemetry/scheduler.server.ts @@ -0,0 +1,71 @@ +import { Logger } from "@/lib/logger"; +import { runScheduledRollups } from "./rollup.server"; +import { runAlertChecks } from "./alerts.server"; + +const logger = new Logger("Telemetry"); + +const ROLLUP_INTERVAL_MS = 15 * 60 * 1000; +const ALERT_INTERVAL_MS = 5 * 60 * 1000; + +interface SchedulerStatus { + running: boolean; + lastRollupAt?: string; + lastAlertCheckAt?: string; + lastRollupError?: string; + lastAlertError?: string; +} + +const status: SchedulerStatus = { running: false }; + +let initialized = false; +let rollupInFlight = false; +let alertsInFlight = false; +let lastRollupDay: string | undefined; + +const tickRollups = async () => { + if (rollupInFlight) return; + rollupInFlight = true; + try { + lastRollupDay = await runScheduledRollups(lastRollupDay); + status.lastRollupAt = new Date().toISOString(); + status.lastRollupError = undefined; + } catch (error) { + status.lastRollupError = error instanceof Error ? error.message : String(error); + logger.error("Rollup run failed", error); + } finally { + rollupInFlight = false; + } +}; + +const tickAlerts = async () => { + if (alertsInFlight) return; + alertsInFlight = true; + try { + await runAlertChecks(); + status.lastAlertCheckAt = new Date().toISOString(); + status.lastAlertError = undefined; + } catch (error) { + status.lastAlertError = error instanceof Error ? error.message : String(error); + logger.error("Alert check failed", error); + } finally { + alertsInFlight = false; + } +}; + +// Single-replica deployment makes in-process intervals a valid scheduler +// (same precedent as the badge migration job). Started lazily from request +// handlers (health probes guarantee a call within seconds of boot). +export const ensureTelemetryScheduler = () => { + if (initialized || typeof window !== "undefined") return; + initialized = true; + status.running = true; + + setTimeout(() => void tickRollups(), 15_000); + setTimeout(() => void tickAlerts(), 30_000); + setInterval(() => void tickRollups(), ROLLUP_INTERVAL_MS); + setInterval(() => void tickAlerts(), ALERT_INTERVAL_MS); + + logger.info("Telemetry scheduler started"); +}; + +export const getSchedulerStatus = (): SchedulerStatus => ({ ...status });