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); }; 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)}'`; 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, opts: { trailing?: boolean } = {} ): 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) ); // Trailing scans re-read the full 7/30-day windows; hourly at most. let wau: number | undefined; let mau: number | undefined; if (opts.trailing !== false) { const nextDay = addDays(dateStr, 1); [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; }; export const runScheduledRollups = async ( lastRunDay: string | undefined, opts: { trailing?: boolean } = {} ): Promise => { const today = utcDay(new Date()); if (lastRunDay && lastRunDay !== today) { await runRollups(lastRunDay, { trailing: true }); } await runRollups(today, opts); return today; };