diff --git a/Justfile-apps b/Justfile-apps index 76c05ca07..081487205 100644 --- a/Justfile-apps +++ b/Justfile-apps @@ -21,14 +21,10 @@ server: VERSION_DETAIL="${VERSION_DETAIL:-$(TZ=UTC0 git log -1 --format=%cd --date=format:%Y%m%d HEAD)-$(git rev-parse --short=7 HEAD)}" export VERSION VERSION_DETAIL sidecar_pids=() - # Crons: Postgres-backed scheduler. Jobs: `job_queue` (SKIP LOCKED; run multiple for load). + # Crons: Postgres-backed scheduler. Jobs: `job_queue` (SKIP LOCKED; embedded pool via TACHI_JOB_QUEUE_WORKER_POOL). if [[ -z "${TACHI_SERVER_NO_SIDECARS:-}" ]]; then bun run cron-worker & sidecar_pids+=($!) - n="${TACHI_SERVER_JOB_WORKER_COUNT:-1}" - [[ "$n" =~ ^[0-9]+$ ]] && [[ "$n" -ge 1 ]] || n=1 - for _ in $(seq 1 "$n"); do - bun run job-queue-worker & sidecar_pids+=($!) - done + bun run job-queue-worker & sidecar_pids+=($!) fi _cleanup_sidecars() { trap - INT TERM EXIT diff --git a/typescript/server/src/cron-worker.ts b/typescript/server/src/cron-worker.ts index f3f4ae14d..c715c02ab 100644 --- a/typescript/server/src/cron-worker.ts +++ b/typescript/server/src/cron-worker.ts @@ -4,6 +4,7 @@ loadServerEnvFile(process.env.NODE_ENV === "test" ? ".env.test" : ".env"); import { runCronTickOnce } from "#lib/jobs/cron/cron-service"; import { log } from "#lib/log/log"; +import { maybeStartWorkerMetricsServer } from "#lib/metrics/worker-metrics"; import { Env } from "#lib/setup/config"; import { ClosePgConnection } from "#services/pg/db"; import { Sleep } from "#utils/misc"; @@ -27,6 +28,7 @@ void bootstrap(); */ async function bootstrap() { await applyMigrations(Env.POSTGRES_URL, Env.MIGRATIONS_DIR); + const metrics = await maybeStartWorkerMetricsServer(process.env); log.info({ bootInfo: true }, "tachi cron worker starting."); function touchHeartbeatFile(): void { @@ -59,6 +61,7 @@ async function bootstrap() { } } finally { clearInterval(heartbeatInterval); + metrics?.close(); } log.info("Cron worker stopped."); await ClosePgConnection(); diff --git a/typescript/server/src/job-queue-worker.ts b/typescript/server/src/job-queue-worker.ts index 56db91d82..f2a74d7b5 100644 --- a/typescript/server/src/job-queue-worker.ts +++ b/typescript/server/src/job-queue-worker.ts @@ -3,7 +3,9 @@ import { loadServerEnvFile } from "#lib/setup/load-server-env"; loadServerEnvFile(process.env.NODE_ENV === "test" ? ".env.test" : ".env"); import { JOB_KIND_SCORE_IMPORT } from "#lib/jobs/job-queue/constants"; +import { parseJobQueueWorkerOptions } from "#lib/jobs/job-queue/parse-worker-options"; import { ClaimNextJob, MarkJobDone, MarkJobFailed } from "#lib/jobs/job-queue/queue-ops"; +import { maybeStartWorkerMetricsServer } from "#lib/metrics/worker-metrics"; import { log } from "#lib/log/log"; import { CloseScoreImportQueue } from "#lib/score-import/worker/queue"; import { processScoreImportJobFromPayload } from "#lib/score-import/worker/score-import-job-processor"; @@ -11,6 +13,7 @@ import { Env } from "#lib/setup/config"; import { ClosePgConnection } from "#services/pg/db"; import { CloseRedisConnection } from "#services/redis/redis"; import { Sleep } from "#utils/misc"; +import { Counter, Histogram } from "prom-client"; import { writeFileSync } from "fs"; import { applyMigrations } from "tachi-db-migration-engine"; @@ -18,6 +21,9 @@ const HEARTBEAT_FILE = "/tmp/worker-heartbeat"; const POLL_MS = 250; +/** Seconds — mirrors SCORE_IMPORT_DURATION_BUCKETS in prometheus.ts; jobs can run sub-second to 30 min. */ +const JOB_DURATION_BUCKETS = [0.25, 0.5, 1, 2, 5, 10, 30, 60, 120, 300, 600, 1800]; + process.on("uncaughtException", (err, origin) => { log.fatal({ err, origin }, "Uncaught exception, terminating."); log.flush(() => process.exit(1)); @@ -27,7 +33,33 @@ void bootstrap(); async function bootstrap() { await applyMigrations(Env.POSTGRES_URL, Env.MIGRATIONS_DIR); - log.info({ bootInfo: true }, "tachi job-queue worker starting (Postgres job_queue)."); + + const { workerCount } = parseJobQueueWorkerOptions(process.argv.slice(2), process.env); + const metrics = await maybeStartWorkerMetricsServer(process.env); + + log.info( + { bootInfo: true, workerCount, pgPoolMax: Env.PG_POOL_MAX }, + "tachi job-queue worker starting (Postgres job_queue).", + ); + + let jobsTotal: Counter | null = null; + let jobDurationSeconds: Histogram | null = null; + + if (metrics) { + jobsTotal = new Counter({ + name: "job_queue_jobs_total", + help: "Total number of job_queue jobs completed, by kind and status.", + labelNames: ["job_kind", "status"], + registers: [metrics.registry], + }); + jobDurationSeconds = new Histogram({ + name: "job_queue_job_duration_seconds", + help: "Wall-clock duration of job_queue jobs in seconds (claim through mark-done/failed).", + labelNames: ["job_kind"], + buckets: JOB_DURATION_BUCKETS, + registers: [metrics.registry], + }); + } let stopping = false; const shutdown = () => { @@ -36,38 +68,58 @@ async function bootstrap() { process.on("SIGINT", shutdown); process.on("SIGTERM", shutdown); - // eslint-disable-next-line no-unmodified-loop-condition - while (!stopping) { - writeFileSync(HEARTBEAT_FILE, Date.now().toString()); + // Separate from the worker loops so long-running jobs cannot stall mtime updates. + writeFileSync(HEARTBEAT_FILE, Date.now().toString()); + const heartbeatInterval = setInterval( + () => writeFileSync(HEARTBEAT_FILE, Date.now().toString()), + 5_000, + ); - const job = await ClaimNextJob(); + async function runWorkerLoop(workerId: number): Promise { + // eslint-disable-next-line no-unmodified-loop-condition + while (!stopping) { + const job = await ClaimNextJob(); - if (!job) { - if (stopping) { - break; - } - - await Sleep(POLL_MS); - continue; - } - - try { - switch (job.job_kind) { - case JOB_KIND_SCORE_IMPORT: - await processScoreImportJobFromPayload(job.payload); + if (!job) { + if (stopping) { break; - default: - log.error({ job_kind: job.job_kind, row_id: job.row_id }, "Unknown job_kind."); - throw new Error(`Unknown job_kind ${String(job.job_kind)}`); + } + await Sleep(POLL_MS); + continue; + } + + const startMs = Date.now(); + + try { + switch (job.job_kind) { + case JOB_KIND_SCORE_IMPORT: + await processScoreImportJobFromPayload(job.payload); + break; + default: + log.error( + { job_kind: job.job_kind, row_id: job.row_id, workerId }, + "Unknown job_kind.", + ); + throw new Error(`Unknown job_kind ${String(job.job_kind)}`); + } + await MarkJobDone(job.row_id); + jobDurationSeconds?.observe({ job_kind: job.job_kind }, (Date.now() - startMs) / 1000); + jobsTotal?.inc({ job_kind: job.job_kind, status: "success" }); + } catch (e) { + log.error(e, `Job ${job.row_id} (worker ${workerId}) failed.`); + await MarkJobFailed(job.row_id); + jobDurationSeconds?.observe({ job_kind: job.job_kind }, (Date.now() - startMs) / 1000); + jobsTotal?.inc({ job_kind: job.job_kind, status: "failure" }); } - await MarkJobDone(job.row_id); - } catch (e) { - log.error(e, `Job ${job.row_id} failed.`); - await MarkJobFailed(job.row_id); } } - log.info("Job worker loop stopped, closing resources."); + const workerPromises = Array.from({ length: workerCount }, (_, i) => runWorkerLoop(i)); + await Promise.allSettled(workerPromises); + + clearInterval(heartbeatInterval); + metrics?.close(); + log.info("Job worker loops stopped, closing resources."); await CloseScoreImportQueue(); await CloseRedisConnection(); await ClosePgConnection(); diff --git a/typescript/server/src/lib/jobs/job-queue/parse-worker-options.test.ts b/typescript/server/src/lib/jobs/job-queue/parse-worker-options.test.ts new file mode 100644 index 000000000..ccb0b7e44 --- /dev/null +++ b/typescript/server/src/lib/jobs/job-queue/parse-worker-options.test.ts @@ -0,0 +1,92 @@ +import { describe, expect, it } from "vitest"; +import { parseJobQueueWorkerOptions } from "./parse-worker-options"; + +describe("parseJobQueueWorkerOptions", () => { + describe("defaults", () => { + it("returns workerCount 1 when argv and env are empty", () => { + expect(parseJobQueueWorkerOptions([], {})).toStrictEqual({ workerCount: 1 }); + }); + }); + + describe("env var TACHI_JOB_QUEUE_WORKER_POOL", () => { + it("parses a valid integer", () => { + expect( + parseJobQueueWorkerOptions([], { TACHI_JOB_QUEUE_WORKER_POOL: "4" }), + ).toStrictEqual({ workerCount: 4 }); + }); + + it("clamps values below 1 to default", () => { + expect( + parseJobQueueWorkerOptions([], { TACHI_JOB_QUEUE_WORKER_POOL: "0" }), + ).toStrictEqual({ workerCount: 1 }); + }); + + it("ignores non-integer values and falls back to default", () => { + expect( + parseJobQueueWorkerOptions([], { TACHI_JOB_QUEUE_WORKER_POOL: "banana" }), + ).toStrictEqual({ workerCount: 1 }); + }); + + it("ignores empty string and falls back to default", () => { + expect( + parseJobQueueWorkerOptions([], { TACHI_JOB_QUEUE_WORKER_POOL: "" }), + ).toStrictEqual({ workerCount: 1 }); + }); + + it("ignores negative values and falls back to default", () => { + expect( + parseJobQueueWorkerOptions([], { TACHI_JOB_QUEUE_WORKER_POOL: "-3" }), + ).toStrictEqual({ workerCount: 1 }); + }); + }); + + describe("CLI --workers argument", () => { + it("parses --workers N (space-separated)", () => { + expect(parseJobQueueWorkerOptions(["--workers", "8"], {})).toStrictEqual({ + workerCount: 8, + }); + }); + + it("parses --workers=N (equals form)", () => { + expect(parseJobQueueWorkerOptions(["--workers=3"], {})).toStrictEqual({ + workerCount: 3, + }); + }); + + it("ignores non-integer --workers N and falls back to env", () => { + expect( + parseJobQueueWorkerOptions(["--workers", "abc"], { + TACHI_JOB_QUEUE_WORKER_POOL: "2", + }), + ).toStrictEqual({ workerCount: 2 }); + }); + + it("ignores --workers=0 and falls back to env", () => { + expect( + parseJobQueueWorkerOptions(["--workers=0"], { TACHI_JOB_QUEUE_WORKER_POOL: "5" }), + ).toStrictEqual({ workerCount: 5 }); + }); + + it("ignores --workers N when N is missing (end of argv) and falls back to env", () => { + expect( + parseJobQueueWorkerOptions(["--workers"], { TACHI_JOB_QUEUE_WORKER_POOL: "3" }), + ).toStrictEqual({ workerCount: 3 }); + }); + + it("--workers takes precedence over env when both are valid", () => { + expect( + parseJobQueueWorkerOptions(["--workers=6"], { TACHI_JOB_QUEUE_WORKER_POOL: "2" }), + ).toStrictEqual({ workerCount: 6 }); + }); + + it("tolerates other CLI args before --workers", () => { + expect( + parseJobQueueWorkerOptions(["--some-flag", "--workers", "4"], {}), + ).toStrictEqual({ workerCount: 4 }); + }); + + it("workerCount 1 is valid (does not clamp up)", () => { + expect(parseJobQueueWorkerOptions(["--workers=1"], {})).toStrictEqual({ workerCount: 1 }); + }); + }); +}); diff --git a/typescript/server/src/lib/jobs/job-queue/parse-worker-options.ts b/typescript/server/src/lib/jobs/job-queue/parse-worker-options.ts new file mode 100644 index 000000000..01c6131de --- /dev/null +++ b/typescript/server/src/lib/jobs/job-queue/parse-worker-options.ts @@ -0,0 +1,44 @@ +export interface JobQueueWorkerOptions { + workerCount: number; +} + +/** + * Resolve how many parallel worker loops to run. + * + * Priority (highest first): + * 1. `--workers=N` / `--workers N` CLI argument + * 2. `TACHI_JOB_QUEUE_WORKER_POOL` env var + * 3. Default: 1 + * + * Any non-integer or value < 1 is silently clamped to 1. + */ +export function parseJobQueueWorkerOptions( + argv: readonly string[], + env: Readonly>, +): JobQueueWorkerOptions { + // CLI takes precedence over env. + for (let i = 0; i < argv.length; i++) { + const arg = argv[i]; + if (arg === "--workers") { + const n = Number.parseInt(argv[i + 1] ?? "", 10); + if (!Number.isNaN(n) && n >= 1) { + return { workerCount: n }; + } + } else if (arg?.startsWith("--workers=")) { + const n = Number.parseInt(arg.slice("--workers=".length), 10); + if (!Number.isNaN(n) && n >= 1) { + return { workerCount: n }; + } + } + } + + const envVal = env["TACHI_JOB_QUEUE_WORKER_POOL"]; + if (envVal !== undefined && envVal !== "") { + const n = Number.parseInt(envVal, 10); + if (!Number.isNaN(n) && n >= 1) { + return { workerCount: n }; + } + } + + return { workerCount: 1 }; +} diff --git a/typescript/server/src/lib/metrics/worker-metrics.ts b/typescript/server/src/lib/metrics/worker-metrics.ts new file mode 100644 index 000000000..75f3798a9 --- /dev/null +++ b/typescript/server/src/lib/metrics/worker-metrics.ts @@ -0,0 +1,82 @@ +import http from "http"; + +import { log } from "#lib/log/log"; +import { Registry, collectDefaultMetrics } from "prom-client"; + +export interface WorkerMetrics { + registry: Registry; + /** Stops the HTTP listener. Safe to call more than once. */ + close: () => void; +} + +/** + * Start a bare HTTP server serving `GET /metrics` for Prometheus scraping. + * + * Only call this when `WORKER_METRICS_PORT` is explicitly set in env — workers + * skip metrics in environments (e.g. local dev) where all processes share a host + * and would conflict on the same port. + * + * Node.js default process metrics (`process_cpu_*`, `nodejs_heap_*`, + * `nodejs_eventloop_lag_*`, etc.) are registered automatically. Custom metrics + * can be added by the caller against the returned `registry`. + */ +export async function startWorkerMetricsServer(port: number): Promise { + const registry = new Registry(); + collectDefaultMetrics({ register: registry }); + + const server = http.createServer((req, res) => { + if (req.method === "GET" && req.url === "/metrics") { + registry + .metrics() + .then((body) => { + res.writeHead(200, { "Content-Type": registry.contentType }); + res.end(body); + }) + .catch((err: unknown) => { + log.error(err, "Failed to collect metrics."); + res.writeHead(500); + res.end(); + }); + } else { + res.writeHead(404); + res.end(); + } + }); + + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(port, resolve); + }); + + log.info({ bootInfo: true }, `Worker metrics listening on port ${port} (/metrics).`); + + let closed = false; + const close = () => { + if (closed) { + return; + } + closed = true; + server.close(); + }; + + return { registry, close }; +} + +/** + * Parse `WORKER_METRICS_PORT` from env and start the metrics server if present. + * Returns `null` when the env var is absent or not a valid integer ≥ 1. + */ +export async function maybeStartWorkerMetricsServer( + env: Readonly>, +): Promise { + const raw = env["WORKER_METRICS_PORT"]; + if (!raw) { + return null; + } + const port = Number.parseInt(raw, 10); + if (Number.isNaN(port) || port < 1) { + log.warn(`WORKER_METRICS_PORT "${raw}" is not a valid port, metrics disabled.`); + return null; + } + return startWorkerMetricsServer(port); +} diff --git a/typescript/server/src/lib/setup/config.ts b/typescript/server/src/lib/setup/config.ts index c26d5f616..61ab8b6b9 100644 --- a/typescript/server/src/lib/setup/config.ts +++ b/typescript/server/src/lib/setup/config.ts @@ -579,6 +579,13 @@ if (!versionDetail) { versionDetail = "unknown"; } +/** + * Maximum number of connections in the Postgres pool for this process. + * For job-queue workers, set this to at least `TACHI_JOB_QUEUE_WORKER_POOL + 2` + * so every worker loop can hold a connection simultaneously without queuing. + */ +const PG_POOL_MAX = parseIntEnv("PG_POOL_MAX", 10); + export const Env = { PORT, REDIS_URL, @@ -588,4 +595,5 @@ export const Env = { VERSION_DETAIL: versionDetail, NODE_ENV: NODE_ENV as "dev" | "production" | "staging" | "test", LOG_LEVEL: logLevel as "crit" | "debug" | "error" | "info" | "severe" | "verbose" | "warn", + PG_POOL_MAX, }; diff --git a/typescript/server/src/services/pg/db.ts b/typescript/server/src/services/pg/db.ts index 5059911d1..d320162c8 100644 --- a/typescript/server/src/services/pg/db.ts +++ b/typescript/server/src/services/pg/db.ts @@ -12,7 +12,7 @@ pg.types.setTypeParser(pg.types.builtins.INT4, (val) => Number(val)); pg.types.setTypeParser(pg.types.builtins.INT2, (val) => Number(val)); pg.types.setTypeParser(pg.types.builtins.INT8, (val) => Number(val)); -const pool = new Pool({ connectionString: Env.POSTGRES_URL }); +const pool = new Pool({ connectionString: Env.POSTGRES_URL, max: Env.PG_POOL_MAX }); if (process.env.NODE_ENV === "test") { // Swallow 57P01 (admin_shutdown) errors that arrive on idle pool connections