feat: job queue parallelism & metrics

This commit is contained in:
zk
2026-05-17 22:40:00 +00:00
parent 9e9952e1b6
commit 0cf9e9cf23
8 changed files with 310 additions and 33 deletions
+2 -6
View File
@@ -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
+3
View File
@@ -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();
+78 -26
View File
@@ -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<void> {
// 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();
@@ -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 });
});
});
});
@@ -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<Record<string, string | undefined>>,
): 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 };
}
@@ -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<WorkerMetrics> {
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<void>((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<Record<string, string | undefined>>,
): Promise<WorkerMetrics | null> {
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);
}
@@ -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,
};
+1 -1
View File
@@ -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