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
This commit is contained in:
2026-08-25 21:38:08 -07:00
parent e2863c1f6c
commit 06fe1f4ed7
12 changed files with 1009 additions and 11 deletions
+11 -6
View File
@@ -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);
+11 -5
View File
@@ -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" },
}
);
},
+76
View File
@@ -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());
+171
View File
@@ -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<TelemetryListResult<ClientEvent>>(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<TelemetryListResult<ClientError>>(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<ClientErrorGroup[]>(async () => {
const { from, to, includeResolved = false, limit = 50 } = data;
const rollups = await pbAdmin.getRollupRange({
metrics: ["client_error.count"],
from,
to,
});
const countByGroup = new Map<string, number>();
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<RollupRow[]>(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<TelemetryListResult<TelemetryAlert>>(async () =>
pbAdmin.listTelemetryAlerts(data.page, data.perPage)
)
);
export const getTelemetryRuntimeStatus = createServerFn()
.middleware([superTokensAdminFunctionMiddleware])
.handler(async () =>
toServerResult<TelemetryRuntimeStatus>(async () => {
const health = await getHealthSnapshot();
return {
health: { status: health.status, checks: health.checks },
scheduler: getSchedulerStatus(),
activeSseConnections: getActiveConnectionCount(),
};
})
);
+38
View File
@@ -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;
}
+14
View File
@@ -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(
+170
View File
@@ -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<string[]> => {
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<typeof sendPushToPlayer>[1]["title"];
body: Parameters<typeof sendPushToPlayer>[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<void> => {
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<string, { count: number; message: string }>();
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<string, number>();
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`,
});
}
};
+56
View File
@@ -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<HealthSnapshot> | null = null;
const ping = async (url: string): Promise<boolean> => {
try {
const response = await fetch(url, { signal: AbortSignal.timeout(CHECK_TIMEOUT_MS) });
return response.ok;
} catch {
return false;
}
};
const runChecks = async (): Promise<HealthSnapshot> => {
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<HealthSnapshot> => {
if (cached && Date.now() - Date.parse(cached.checkedAt) < CHECK_TTL_MS) {
return cached;
}
if (!refreshing) {
refreshing = runChecks().finally(() => {
refreshing = null;
});
}
return cached ?? refreshing;
};
+113
View File
@@ -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<typeof buildDayRollups>, 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);
});
});
+175
View File
@@ -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<string>;
perFn: Map<string, { count: number; errors: number; durations: number[] }>;
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<string>;
sessions: Set<string>;
pageViewsByRoute: Map<string, number>;
vitals: Map<string, number[]>;
}
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<string, number>;
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<string, unknown>) => {
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;
};
+103
View File
@@ -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<number> => {
const players = new Set<string>();
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<number> => {
await pbAdmin.authPromise;
const activities = newActivityAccumulator();
const clientEvents = newClientEventAccumulator();
const clientErrors = newClientErrorAccumulator();
await pbAdmin.pageCollection<ActivityRow>(
"activities",
{ filter: dayFilter(dateStr), fields: "name,player,duration,success,error,created" },
(rows) => accumulateActivities(activities, rows)
);
await pbAdmin.pageCollection<ClientEventRow>(
"client_events",
{ filter: dayFilter(dateStr), fields: "name,player,session_id,route_id,value,created" },
(rows) => accumulateClientEvents(clientEvents, rows)
);
await pbAdmin.pageCollection<ClientErrorRow>(
"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<string> => {
const today = utcDay(new Date());
if (lastRunDay && lastRunDay !== today) {
await runRollups(lastRunDay);
}
await runRollups(today);
return today;
};
+71
View File
@@ -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 });