-
Notifications
You must be signed in to change notification settings - Fork 190
feat(worker): rebuild global stats from logs #3552
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,180 @@ | ||
| /* eslint-disable no-console */ | ||
| /** | ||
| * Rebuild global_model_stats / global_source_stats from `log`. | ||
| * | ||
| * Why this is possible at all: data retention never deletes log rows. The | ||
| * cleanup job (cleanupExpiredLogData) only nulls the verbose payload columns — | ||
| * messages, content, raw/upstream request and response, tools, customHeaders, | ||
| * userAgent, responsesApiData — and sets data_retention_cleaned_up. Every column | ||
| * the aggregator reads (used_model, used_provider, used_mode, source, | ||
| * organization_id, created_at, the token columns, the cost columns, has_error, | ||
| * cached, streamed, unified_finish_reason) survives untouched, so a rebuild is | ||
| * lossless. | ||
| * | ||
| * That makes this the correct fix for rows the mode/kind migration | ||
| * (1786136603_warm_omega_flight.sql) left as 'unknown': re-deriving from the | ||
| * source data gives an exact used_mode and org_kind for every row, with no | ||
| * proration and no scaling. | ||
| * | ||
| * Each day is handled by the aggregator's own recomputeDayFully(), which deletes | ||
| * that day's rows and re-aggregates its 24 hours. So the script is idempotent | ||
| * and resumable: re-running any day is safe, and a crashed run is continued by | ||
| * re-running with --from/--to for whatever is left. | ||
| * | ||
| * Ordering: newest day first by default, so the recent days people actually look | ||
| * at are correct within minutes instead of at the very end of the run. | ||
| * | ||
| * Today and yesterday are skipped. The incremental walker owns today, and the | ||
| * safety net wipes-and-recomputes yesterday; racing either of them risks | ||
| * interleaving a delete with an insert and double-counting a day. | ||
| * | ||
| * Usage: | ||
| * pnpm --filter worker rebuild-global-stats # dry run: report the plan | ||
| * pnpm --filter worker rebuild-global-stats --commit # rebuild everything | ||
| * pnpm --filter worker rebuild-global-stats --commit --limit=7 # rebuild the 7 newest days first | ||
| * pnpm --filter worker rebuild-global-stats --commit --from=2025-05-20 --to=2026-08-01 | ||
| * pnpm --filter worker rebuild-global-stats --commit --oldest-first | ||
| * | ||
| * Environment: | ||
| * DATABASE_URL - defaults to local postgres if unset | ||
| */ | ||
|
|
||
| import { recomputeDayFully } from "@/services/global-stats-aggregator.js"; | ||
|
|
||
| import { db, sql } from "@llmgateway/db"; | ||
|
|
||
| const DAY_MS = 24 * 60 * 60 * 1000; | ||
|
|
||
| function parseFlag(name: string): string | undefined { | ||
| const flag = `--${name}=`; | ||
| const arg = process.argv.find((a) => a.startsWith(flag)); | ||
| return arg ? arg.slice(flag.length) : undefined; | ||
| } | ||
|
|
||
| function hasFlag(name: string): boolean { | ||
| return process.argv.includes(`--${name}`); | ||
| } | ||
|
|
||
| function floorToUTCDay(d: Date): Date { | ||
| return new Date( | ||
| Date.UTC(d.getUTCFullYear(), d.getUTCMonth(), d.getUTCDate(), 0, 0, 0, 0), | ||
| ); | ||
| } | ||
|
|
||
| function toDateString(d: Date): string { | ||
| return d.toISOString().slice(0, 10); | ||
| } | ||
|
|
||
| /** | ||
| * Earliest day worth rebuilding: whichever of the first log and the first | ||
| * existing stats row is older, so a rebuild can also fill in days the walker | ||
| * never reached. | ||
| * | ||
| * `order by created_at limit 1` rather than `min(created_at)` — the log table is | ||
| * huge, and this form is guaranteed to ride the index on (created_at, ...) | ||
| * instead of risking a sequential scan. | ||
| */ | ||
| async function resolveStart(): Promise<Date | null> { | ||
| const result = await db.execute(sql` | ||
| select least( | ||
| (select created_at from log order by created_at limit 1), | ||
| (select min(day_timestamp) from global_model_stats), | ||
| (select min(day_timestamp) from global_source_stats) | ||
| ) as start | ||
| `); | ||
| const start = (result.rows[0] as { start: Date | string | null } | undefined) | ||
| ?.start; | ||
| return start ? floorToUTCDay(new Date(start)) : null; | ||
| } | ||
|
|
||
| async function main(): Promise<void> { | ||
| const commit = hasFlag("commit"); | ||
| const oldestFirst = hasFlag("oldest-first"); | ||
| const limit = parseFlag("limit") ? Number(parseFlag("limit")) : undefined; | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win Validate all CLI flags before use. The script converts
📍 Affects 1 file
🤖 Prompt for AI Agents |
||
|
|
||
| const fromFlag = parseFlag("from"); | ||
| const toFlag = parseFlag("to"); | ||
|
|
||
| const today = floorToUTCDay(new Date()); | ||
| // Today belongs to the walker, yesterday to the safety net. | ||
| const defaultEnd = new Date(today.getTime() - 2 * DAY_MS); // eslint-disable-line no-mixed-operators | ||
|
|
||
| const start = fromFlag | ||
| ? floorToUTCDay(new Date(`${fromFlag}T00:00:00Z`)) | ||
| : await resolveStart(); | ||
| if (!start) { | ||
| console.log("No logs and no global stats rows — nothing to rebuild."); | ||
| return; | ||
| } | ||
|
|
||
| let end = toFlag | ||
| ? floorToUTCDay(new Date(`${toFlag}T00:00:00Z`)) | ||
| : defaultEnd; | ||
| if (end.getTime() > defaultEnd.getTime()) { | ||
| console.log( | ||
| `Clamping --to to ${toDateString(defaultEnd)}: today and yesterday are owned by the walker and the safety net.`, | ||
| ); | ||
| end = defaultEnd; | ||
| } | ||
|
|
||
| if (end.getTime() < start.getTime()) { | ||
| console.log("Nothing to rebuild — the range is empty."); | ||
| return; | ||
| } | ||
|
|
||
| const days: Date[] = []; | ||
| for (let d = new Date(start); d <= end; d = new Date(d.getTime() + DAY_MS)) { | ||
| days.push(new Date(d)); | ||
| } | ||
| if (!oldestFirst) { | ||
| days.reverse(); | ||
| } | ||
| const planned = limit ? days.slice(0, limit) : days; | ||
|
|
||
| console.log( | ||
| `${commit ? "REBUILDING" : "DRY RUN"} — ${toDateString(start)} .. ${toDateString(end)} (${planned.length} of ${days.length} days, ${oldestFirst ? "oldest" : "newest"} first)`, | ||
| ); | ||
|
|
||
| if (!commit) { | ||
| const before = await db.execute(sql` | ||
| select | ||
| (select count(*) from global_model_stats | ||
| where used_mode = 'unknown' or org_kind = 'unknown') as model_rows_unattributed, | ||
| (select count(*) from global_source_stats | ||
| where used_mode = 'unknown' or org_kind = 'unknown') as source_rows_unattributed | ||
| `); | ||
| console.log(before.rows[0]); | ||
| console.log("\nDry run — nothing was written. Re-run with --commit."); | ||
| return; | ||
| } | ||
|
|
||
| const startedAt = Date.now(); | ||
| let done = 0; | ||
| for (const day of planned) { | ||
| const completed = await recomputeDayFully(day); | ||
| if (!completed) { | ||
| console.log(`Stop requested — halted after ${done} day(s).`); | ||
| break; | ||
| } | ||
|
Comment on lines
+153
to
+158
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win 🧩 Analysis chain🏁 Script executed: #!/bin/bash
# Description: Find where the stop flag used by isStopRequested is set.
rg -nP --type=ts -C5 '\b(isStopRequested|requestStop|stopRequested)\b'Repository: theopenco/llmgateway Length of output: 158 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- candidate file ---'
wc -l apps/worker/src/scripts/rebuild-global-stats.ts
sed -n '1,220p' apps/worker/src/scripts/rebuild-global-stats.ts
printf '%s\n' '--- stop-related definitions and usages ---'
rg -n -C4 'isStopRequested|requestStop|stopRequested|SIGINT|SIGTERM|shutdown|graceful' . \
-g '!node_modules' -g '!dist' -g '!build'Repository: theopenco/llmgateway Length of output: 50377 🏁 Script executed: #!/bin/bash
set -eu
service_file="$(fd -t f 'global-stats-aggregator' . | head -n 1)"
printf 'service_file=%s\n' "$service_file"
wc -l "$service_file"
ast-grep outline "$service_file"
printf '%s\n' '--- recomputeDayFully and direct stop references ---'
rg -n -C12 'recomputeDayFully|isStopRequested|requestStop|stopRequested' "$service_file"
printf '%s\n' '--- all exact symbol references, excluding generated/dependency files ---'
rg -n --hidden -g '!node_modules' -g '!dist' -g '!build' -g '!pnpm-lock.yaml' '\bisStopRequested\b|\brequestStop\b|\bstopRequested\b' .Repository: theopenco/llmgateway Length of output: 486 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- candidate service files ---'
fd -t f 'global-stats-aggregator' apps/worker
printf '%s\n' '--- implementation references ---'
rg -n -C16 'recomputeDayFully|isStopRequested|requestStop|stopRequested' apps/worker/src \
-g '*.ts' -g '!*.spec.ts'Repository: theopenco/llmgateway Length of output: 50376 🏁 Script executed: #!/bin/bash
set -eu
printf '%s\n' '--- shutdown state and signal registration ---'
cat -n apps/worker/src/shutdown.ts | sed -n '1,140p'
printf '%s\n' '--- process signal handlers in worker sources ---'
rg -n -C12 "process\.on\(['\"](SIGINT|SIGTERM)|requestStop\(" apps/worker/src \
-g '*.ts' -g '!*.spec.ts'
printf '%s\n' '--- worker entrypoint imports and startup/shutdown section ---'
rg -n -C10 'from "./shutdown\.js"|from "`@/shutdown`\.js"|process\.on|startWorker|main\(' apps/worker/src/worker.ts apps/worker/src/index.ts apps/worker/src/main.ts 2>/dev/null || trueRepository: theopenco/llmgateway Length of output: 11347 🏁 Script executed: #!/bin/bash
set -eu
python3 - <<'PY'
from pathlib import Path
import re
script = Path("apps/worker/src/scripts/rebuild-global-stats.ts").read_text()
aggregator = Path("apps/worker/src/services/global-stats-aggregator.ts").read_text()
shutdown = Path("apps/worker/src/shutdown.ts").read_text()
worker_entry = Path("apps/worker/src/index.ts").read_text()
checks = {
"script imports recomputeDayFully": "recomputeDayFully" in script,
"aggregator imports isStopRequested": 'import { isStopRequested } from "`@/shutdown.js`";' in aggregator,
"recomputeDayFully checks isStopRequested": bool(
re.search(r"export async function recomputeDayFully[\s\S]*?isStopRequested\(\)", aggregator)
),
"shutdown initializes shouldStop=false": "let shouldStop = false;" in shutdown,
"standalone script has SIGINT handler": bool(
re.search(r"process\.on\(\s*[\"']SIGINT[\"']", script)
),
"standalone script has SIGTERM handler": bool(
re.search(r"process\.on\(\s*[\"']SIGTERM[\"']", script)
),
"worker entrypoint has SIGINT handler": bool(
re.search(r"process\.on\(\s*[\"']SIGINT[\"']", worker_entry)
),
"worker entrypoint has SIGTERM handler": bool(
re.search(r"process\.on\(\s*[\"']SIGTERM[\"']", worker_entry)
),
}
for name, result in checks.items():
print(f"{name}: {result}")
PYRepository: theopenco/llmgateway Length of output: 505 Register signal handlers for the standalone script.
🤖 Prompt for AI Agents |
||
| done++; | ||
| if (done % 10 === 0 || done === planned.length) { | ||
| const elapsed = (Date.now() - startedAt) / 1000; | ||
| const rate = done / elapsed; | ||
| const remaining = planned.length - done; | ||
| console.log( | ||
| `${done}/${planned.length} days (${toDateString(day)}) — ${elapsed.toFixed(0)}s elapsed, ~${rate > 0 ? Math.round(remaining / rate) : 0}s remaining`, | ||
| ); | ||
| } | ||
| } | ||
|
|
||
| console.log( | ||
| `Rebuilt ${done} day(s) in ${((Date.now() - startedAt) / 1000).toFixed(0)}s.`, | ||
| ); | ||
| } | ||
|
|
||
| main() | ||
| .then(() => process.exit(0)) | ||
| .catch((error) => { | ||
| console.error(error); | ||
| process.exit(1); | ||
| }); | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
Repository: theopenco/llmgateway
Length of output: 4960
🏁 Script executed:
Repository: theopenco/llmgateway
Length of output: 50376
🏁 Script executed:
Repository: theopenco/llmgateway
Length of output: 50376
🏁 Script executed:
Repository: theopenco/llmgateway
Length of output: 50376
🏁 Script executed:
Repository: theopenco/llmgateway
Length of output: 45016
🌐 Web query:
node-postgres timestamp without time zone parses local timezone Date parser💡 Result:
By default, node-postgres parses PostgreSQL TIMESTAMP WITHOUT TIME ZONE columns as JavaScript Date objects based on the local time zone of the Node.js process [1][2]. The driver uses the postgres-date library to handle this conversion [3]. Because TIMESTAMP WITHOUT TIME ZONE lacks timezone information, the parser treats the input string as a local time when constructing the Date object [4][3]. This can lead to unexpected behavior if your application expects UTC or if your Node.js server and database server operate in different time zones [1][5]. To change this behavior, you can override the default type parser provided by the pg library [1][6]. Common approaches include: 1. Returning the raw string: If you want to avoid automatic Date conversion entirely, you can disable the parser [7][5]. const { types } = require('pg'); const TIMESTAMP_OID = 1114; types.setTypeParser(TIMESTAMP_OID, (timestamp) => timestamp); 2. Forcing UTC interpretation: If you need a Date object but want it interpreted as UTC, you can append a Z (UTC indicator) to the string before passing it to the Date constructor [8]. const { types } = require('pg'); const TIMESTAMP_OID = 1114; types.setTypeParser(TIMESTAMP_OID, (stringValue) => { return new Date(
${stringValue.replace(' ', 'T')}Z); }); Note: If you are experiencing issues with time offsets, it is often recommended to use the TIMESTAMPTZ (timestamp with time zone) data type in PostgreSQL, which stores the value in UTC and avoids these ambiguity issues [1][2].Citations:
🏁 Script executed:
Repository: theopenco/llmgateway
Length of output: 439
Interpret the
timestampboundary as UTC. The three columns aretimestamp without time zonevalues, andnode-postgresparses them as local-timeDateobjects. UnderTZ=Asia/Tokyo,floorToUTCDaycan select the previous day. Convert theleast(...)expression withAT TIME ZONE 'UTC', or parse text with an explicitZ; a bare text cast is insufficient.🤖 Prompt for AI Agents