Files
flxn-app/src/lib/pocketbase/services/telemetry.ts
T
kyle e4d1ba720a feat(telemetry): session journey view replaces raw events table
- sessions listed from client_events grouped by session id (7d window)
- timeline sheet merges page views, client events, errors, and server-fn
  activities joined by player within the session window
- vitals excluded from timelines; gap markers over 5 minutes
2026-08-25 22:56:09 -07:00

323 lines
9.4 KiB
TypeScript

import PocketBase from "pocketbase";
import { PlayerInfo } from "@/features/players/types";
import { pbFilter } from "../util/filter";
import { likePattern } from "../util/like-pattern";
export type ClientErrorSource = "window" | "unhandledrejection" | "error-boundary" | "sw";
export interface ClientEventRecord {
id: string;
name: string;
player?: string | PlayerInfo;
session_id?: string;
path?: string;
route_id?: string;
value?: number;
props?: any;
user_agent?: string;
created: string;
updated: string;
}
export interface ClientEventInput {
name: string;
player?: string;
session_id?: string;
path?: string;
route_id?: string;
value?: number;
props?: any;
user_agent?: string;
}
export interface ClientErrorRecord {
id: string;
message: string;
stack?: string;
source?: ClientErrorSource;
path?: string;
route_id?: string;
group_hash: string;
resolved: boolean;
player?: string | PlayerInfo;
session_id?: string;
user_agent?: string;
props?: any;
created: string;
updated: string;
}
export interface ClientErrorInput {
message: string;
stack?: string;
source?: ClientErrorSource;
path?: string;
route_id?: string;
group_hash: string;
resolved?: boolean;
player?: string;
session_id?: string;
user_agent?: string;
props?: any;
}
export interface RollupRecord {
id: string;
date: string;
metric: string;
dim: string;
value: number;
meta?: any;
created: string;
updated: string;
}
export interface RollupInput {
date: string;
metric: string;
dim: string;
value: number;
meta?: any;
}
export interface TelemetryAlertRecord {
id: string;
kind: string;
dim: string;
message?: string;
value?: number;
meta?: any;
created: string;
updated: string;
}
export interface TelemetryAlertInput {
kind: string;
dim: string;
message?: string;
value?: number;
meta?: any;
}
export interface TelemetryListResult<T> {
items: T[];
page: number;
perPage: number;
totalPages: number;
totalItems: number;
}
export interface ClientEventSearchParams {
page?: number;
perPage?: number;
name?: string;
player?: string;
sessionId?: string;
path?: string;
from?: string;
to?: string;
sortBy?: string;
}
export interface ClientErrorSearchParams {
page?: number;
perPage?: number;
groupHash?: string;
resolved?: boolean;
player?: string;
sessionId?: string;
from?: string;
to?: string;
sortBy?: string;
}
function isNotFound(error: unknown): boolean {
return (error as { status?: number })?.status === 404;
}
export function createTelemetryService(pb: PocketBase) {
const service = {
async createClientEvent(data: ClientEventInput): Promise<ClientEventRecord> {
return pb.collection("client_events").create<ClientEventRecord>(data);
},
async createClientError(data: ClientErrorInput): Promise<ClientErrorRecord> {
return pb.collection("client_errors").create<ClientErrorRecord>(data);
},
async searchClientEvents(
params: ClientEventSearchParams = {}
): Promise<TelemetryListResult<ClientEventRecord>> {
const { page = 1, perPage = 100, name, player, sessionId, path, from, to, sortBy = "-created" } = params;
const filters: string[] = [];
if (name) filters.push(pbFilter(pb, "name ~ {:name}", { name: likePattern(name) }));
if (player) filters.push(pbFilter(pb, "player = {:player}", { player }));
if (sessionId) filters.push(pbFilter(pb, "session_id = {:sessionId}", { sessionId }));
if (path) filters.push(pbFilter(pb, "path ~ {:path}", { path: likePattern(path) }));
if (from) filters.push(pbFilter(pb, "created >= {:from}", { from }));
if (to) filters.push(pbFilter(pb, "created < {:to}", { to }));
const result = await pb.collection("client_events").getList<ClientEventRecord>(page, perPage, {
filter: filters.join(" && "),
sort: sortBy,
expand: "player",
});
return {
items: result.items,
page: result.page,
perPage: result.perPage,
totalPages: result.totalPages,
totalItems: result.totalItems,
};
},
async searchClientErrors(
params: ClientErrorSearchParams = {}
): Promise<TelemetryListResult<ClientErrorRecord>> {
const { page = 1, perPage = 100, groupHash, resolved, player, sessionId, from, to, sortBy = "-created" } = params;
const filters: string[] = [];
if (groupHash) filters.push(pbFilter(pb, "group_hash = {:groupHash}", { groupHash }));
if (resolved !== undefined) filters.push(pbFilter(pb, "resolved = {:resolved}", { resolved }));
if (player) filters.push(pbFilter(pb, "player = {:player}", { player }));
if (sessionId) filters.push(pbFilter(pb, "session_id = {:sessionId}", { sessionId }));
if (from) filters.push(pbFilter(pb, "created >= {:from}", { from }));
if (to) filters.push(pbFilter(pb, "created < {:to}", { to }));
const result = await pb.collection("client_errors").getList<ClientErrorRecord>(page, perPage, {
filter: filters.join(" && "),
sort: sortBy,
expand: "player",
});
return {
items: result.items,
page: result.page,
perPage: result.perPage,
totalPages: result.totalPages,
totalItems: result.totalItems,
};
},
// dim must be "" (never null/undefined) or the unique index stops deduplicating.
async upsertRollup(input: RollupInput): Promise<RollupRecord> {
const data = { ...input, dim: input.dim ?? "" };
const filter = pbFilter(pb, "date = {:date} && metric = {:metric} && dim = {:dim}", {
date: data.date,
metric: data.metric,
dim: data.dim,
});
try {
const existing = await pb.collection("telemetry_rollups").getFirstListItem<RollupRecord>(filter);
return await pb.collection("telemetry_rollups").update<RollupRecord>(existing.id, data);
} catch (error) {
if (!isNotFound(error)) throw error;
}
try {
return await pb.collection("telemetry_rollups").create<RollupRecord>(data);
} catch {
const existing = await pb.collection("telemetry_rollups").getFirstListItem<RollupRecord>(filter);
return pb.collection("telemetry_rollups").update<RollupRecord>(existing.id, data);
}
},
async getRollupRange(params: {
metrics: string[];
from: string;
to: string;
dim?: string;
}): Promise<RollupRecord[]> {
const { metrics, from, to, dim } = params;
if (metrics.length === 0) return [];
const metricFilter = metrics
.map((metric, i) => pbFilter(pb, `metric = {:m${i}}`, { [`m${i}`]: metric }))
.join(" || ");
const filters = [
`(${metricFilter})`,
pbFilter(pb, "date >= {:from}", { from }),
pbFilter(pb, "date <= {:to}", { to }),
];
if (dim !== undefined) filters.push(pbFilter(pb, "dim = {:dim}", { dim }));
return pb.collection("telemetry_rollups").getFullList<RollupRecord>({
filter: filters.join(" && "),
sort: "date",
});
},
async createTelemetryAlert(data: TelemetryAlertInput): Promise<TelemetryAlertRecord> {
return pb.collection("telemetry_alerts").create<TelemetryAlertRecord>({ ...data, dim: data.dim ?? "" });
},
async findRecentAlert(kind: string, dim: string, sinceIso: string): Promise<TelemetryAlertRecord | null> {
try {
return await pb.collection("telemetry_alerts").getFirstListItem<TelemetryAlertRecord>(
pbFilter(pb, "kind = {:kind} && dim = {:dim} && created >= {:since}", { kind, dim, since: sinceIso })
);
} catch (error) {
if (isNotFound(error)) return null;
throw error;
}
},
async listTelemetryAlerts(page = 1, perPage = 50): Promise<TelemetryListResult<TelemetryAlertRecord>> {
const result = await pb.collection("telemetry_alerts").getList<TelemetryAlertRecord>(page, perPage, {
sort: "-created",
});
return {
items: result.items,
page: result.page,
perPage: result.perPage,
totalPages: result.totalPages,
totalItems: result.totalItems,
};
},
async resolveErrorsByGroup(groupHash: string, resolved = true): Promise<number> {
const MAX_PAGES = 25;
let updated = 0;
for (let i = 0; i < MAX_PAGES; i++) {
const batch = await pb.collection("client_errors").getList<ClientErrorRecord>(1, 200, {
filter: pbFilter(pb, "group_hash = {:groupHash} && resolved != {:resolved}", { groupHash, resolved }),
fields: "id",
skipTotal: true,
});
if (batch.items.length === 0) break;
for (const item of batch.items) {
await pb.collection("client_errors").update(item.id, { resolved });
updated++;
}
if (batch.items.length < 200) break;
}
return updated;
},
async pageCollection<T>(
collection: string,
opts: { filter?: string; fields?: string },
cb: (items: T[]) => void
): Promise<void> {
const PER_PAGE = 500;
for (let page = 1; ; page++) {
const result = await pb.collection(collection).getList<T>(page, PER_PAGE, {
filter: opts.filter ?? "",
fields: opts.fields,
sort: "created",
skipTotal: true,
});
if (result.items.length > 0) cb(result.items);
if (result.items.length < PER_PAGE) break;
}
},
};
return service;
}