diff --git a/.env.example b/.env.example index 4e66a6b6c7..aa9f4d443d 100644 --- a/.env.example +++ b/.env.example @@ -520,6 +520,16 @@ NEXT_PUBLIC_BASE_URL=http://localhost:20128 # cost of more upstream polling; raise to reduce request volume. # OMNIROUTE_CGPT_WEB_PRO_POLL_INTERVAL_MS=4000 +# Timeout for the /api/jobs/:id/run-now endpoint, in milliseconds. +# This bounds the CALL, not the job. runNow() dispatches the handler with +# `void` and returns as soon as it has decided to start, so on the normal +# path it resolves in milliseconds. It only matters when the job is already +# running: runNow() then waits for the in-flight run before starting the +# queued one, and this timeout prevents that wait from hanging forever. +# Used by: src/app/api/jobs/[id]/run-now/route.ts +# Default: 30000 (30 seconds) +# OMNIROUTE_RUNNOW_TIMEOUT_MS=30000 + # Public cloud URL — client-side mirror of CLOUD_URL. NEXT_PUBLIC_CLOUD_URL= diff --git a/config/quality/dependency-allowlist.json b/config/quality/dependency-allowlist.json index 0c36dd1155..842eb696f3 100644 --- a/config/quality/dependency-allowlist.json +++ b/config/quality/dependency-allowlist.json @@ -44,6 +44,7 @@ "commander", "concurrently", "cross-env", + "cron-parser", "csv-stringify", "ctrf", "dompurify", diff --git a/docs/reference/ENVIRONMENT.md b/docs/reference/ENVIRONMENT.md index cd6423d4ea..a013c8526b 100644 --- a/docs/reference/ENVIRONMENT.md +++ b/docs/reference/ENVIRONMENT.md @@ -658,9 +658,10 @@ REQUEST_TIMEOUT_MS (global override) | `OMNIROUTE_DEFAULT_FETCH_TIMEOUT_MS` | `120000` | Fallback used by `src/shared/utils/fetchTimeout.ts` when `FETCH_TIMEOUT_MS` is unset. | | `OMNIROUTE_CHATGPT_TLS_TIMEOUT_MS` | `60000` | Wire-level timeout for the bogdanfinn/tls-client koffi binding (`chatgptTlsClient.ts`). | | `OMNIROUTE_CHATGPT_TLS_GRACE_MS` | `10000` | JS-side grace added on top of the wire timeout when the native binding is wedged. | -| `OMNIROUTE_CHATGPT_STREAM_FIRST_BYTE_TIMEOUT_MS` | `30000` (30s) | Max wait for the first streamed byte from the ChatGPT TLS sidecar (`chatgptTlsClient.ts`) before aborting a dead stream. Raise if upstream cold-starts exceed the window. | -| `OMNIROUTE_CLAUDE_TLS_TIMEOUT_MS` | `60000` | Wire-level timeout for the bogdanfinn/tls-client koffi binding (`claudeTlsClient.ts`). | -| `OMNIROUTE_CLAUDE_TLS_GRACE_MS` | `10000` | JS-side grace added on top of the wire timeout when the native binding is wedged. | +| `OMNIROUTE_CHATGPT_STREAM_FIRST_BYTE_TIMEOUT_MS` | `30000` | Max wait for the first streamed byte from the ChatGPT TLS sidecar. | +| `OMNIROUTE_CLAUDE_TLS_TIMEOUT_MS` | `60000` | Wire-level timeout for the bogdanfinn/tls-client koffi binding. | +| `OMNIROUTE_CLAUDE_TLS_GRACE_MS` | `10000` | JS-side grace added on top of the wire timeout when the native binding is wedged. | +| `OMNIROUTE_RUNNOW_TIMEOUT_MS` | `30000` | Timeout for the `/api/jobs/:id/run-now` endpoint. Bounds how long a run-now call waits for an in-flight job to finish before starting the queued run. See `src/app/api/jobs/[id]/run-now/route.ts`. | | `OMNIROUTE_PPLX_TLS_TIMEOUT_MS` | `30000` | Wire-level timeout for the bogdanfinn/tls-client koffi binding (`perplexityTlsClient.ts`). | | `OMNIROUTE_PPLX_TLS_GRACE_MS` | `10000` | JS-side grace added on top of the wire timeout when the native binding is wedged. | | `OMNIROUTE_GROK_TLS_TIMEOUT_MS` | `60000` | Wire-level timeout for the bogdanfinn/tls-client koffi binding (`grokTlsClient.ts`). | diff --git a/package-lock.json b/package-lock.json index f83e6a5940..95985d2f8a 100644 --- a/package-lock.json +++ b/package-lock.json @@ -30,6 +30,7 @@ "bottleneck": "^2.19.5", "clsx": "^2.1.1", "commander": "^15.0.0", + "cron-parser": "^5.6.2", "csv-stringify": "^6.7.0", "dompurify": "^3.4.13", "express": "^5.2.1", @@ -15879,6 +15880,18 @@ "node": ">= 6" } }, + "node_modules/cron-parser": { + "version": "5.7.0", + "resolved": "https://registry.npmjs.org/cron-parser/-/cron-parser-5.7.0.tgz", + "integrity": "sha512-iSpDHpwwW/GhIg4JVODYlWUEpMNSimaHvqOhHpOz1W+Y97z1lL1nf+dpcF17cNwFRpTtKN9devgi1fxflp3Phw==", + "license": "MIT", + "dependencies": { + "luxon": "^3.7.2" + }, + "engines": { + "node": ">=18" + } + }, "node_modules/cross-env": { "version": "10.1.0", "resolved": "https://registry.npmjs.org/cross-env/-/cross-env-10.1.0.tgz", @@ -25499,6 +25512,15 @@ "react": "^16.5.1 || ^17.0.0 || ^18.0.0 || ^19.0.0" } }, + "node_modules/luxon": { + "version": "3.7.2", + "resolved": "https://registry.npmjs.org/luxon/-/luxon-3.7.2.tgz", + "integrity": "sha512-vtEhXh/gNjI9Yg1u4jX/0YVPMvxzHuGgCm6tC5kZyb08yjGWGnqAjGJvcXbqQR2P3MyMEFnRbpcdFS6PBcLqew==", + "license": "MIT", + "engines": { + "node": ">=12" + } + }, "node_modules/magic-string": { "version": "0.30.21", "resolved": "https://registry.npmjs.org/magic-string/-/magic-string-0.30.21.tgz", diff --git a/package.json b/package.json index f90b56d96e..7aaed13f0d 100644 --- a/package.json +++ b/package.json @@ -255,6 +255,7 @@ "bottleneck": "^2.19.5", "clsx": "^2.1.1", "commander": "^15.0.0", + "cron-parser": "^5.6.2", "csv-stringify": "^6.7.0", "dompurify": "^3.4.13", "express": "^5.2.1", diff --git a/src/app/api/jobs/[id]/disable/route.ts b/src/app/api/jobs/[id]/disable/route.ts new file mode 100644 index 0000000000..cc03469ecb --- /dev/null +++ b/src/app/api/jobs/[id]/disable/route.ts @@ -0,0 +1,26 @@ +/** + * POST /api/jobs/:id/disable + * + * Disable a job and stop its timer. LOCAL_ONLY. + */ +import { NextResponse } from "next/server"; + +import { buildErrorBody } from "@omniroute/open-sse/utils/error.ts"; +import { getJobRegistry } from "@/lib/jobRegistry"; + +export const dynamic = "force-dynamic"; + +export async function POST(_request: Request, { params }: { params: Promise<{ id: string }> }) { + try { + const { id } = await params; + const registry = getJobRegistry(); + if (!registry.listJobs().some((j) => j.id === id)) { + return NextResponse.json(buildErrorBody(404, "Job not found"), { status: 404 }); + } + registry.setEnabled(id, false); + return NextResponse.json({ data: { id, enabled: false } }); + } catch (err) { + console.error("[API] POST /api/jobs/:id/disable error:", err); + return NextResponse.json(buildErrorBody(500, "Failed to disable job"), { status: 500 }); + } +} diff --git a/src/app/api/jobs/[id]/enable/route.ts b/src/app/api/jobs/[id]/enable/route.ts new file mode 100644 index 0000000000..38e5448c41 --- /dev/null +++ b/src/app/api/jobs/[id]/enable/route.ts @@ -0,0 +1,26 @@ +/** + * POST /api/jobs/:id/enable + * + * Enable a disabled job and (re)start its timer. LOCAL_ONLY. + */ +import { NextResponse } from "next/server"; + +import { buildErrorBody } from "@omniroute/open-sse/utils/error.ts"; +import { getJobRegistry } from "@/lib/jobRegistry"; + +export const dynamic = "force-dynamic"; + +export async function POST(_request: Request, { params }: { params: Promise<{ id: string }> }) { + try { + const { id } = await params; + const registry = getJobRegistry(); + if (!registry.listJobs().some((j) => j.id === id)) { + return NextResponse.json(buildErrorBody(404, "Job not found"), { status: 404 }); + } + registry.setEnabled(id, true); + return NextResponse.json({ data: { id, enabled: true } }); + } catch (err) { + console.error("[API] POST /api/jobs/:id/enable error:", err); + return NextResponse.json(buildErrorBody(500, "Failed to enable job"), { status: 500 }); + } +} diff --git a/src/app/api/jobs/[id]/run-now/route.ts b/src/app/api/jobs/[id]/run-now/route.ts new file mode 100644 index 0000000000..bc4cf429ab --- /dev/null +++ b/src/app/api/jobs/[id]/run-now/route.ts @@ -0,0 +1,51 @@ +/** + * POST /api/jobs/:id/run-now -- manually trigger a job run. + * LOCAL_ONLY (enforced by routeGuard). + * + * The timeout bounds the CALL, not the job. runNow() dispatches the handler + * with `void` and returns as soon as it has decided to start, so on the normal + * path this resolves in milliseconds and the timer never fires. It only has + * something to bound when the job is already running: runNow() then returns a + * promise that waits for the in-flight run to finish before starting the queued + * one. Cancelling here does not cancel the job -- the handler keeps running. + */ +import { NextResponse } from "next/server"; + +import { buildErrorBody } from "@omniroute/open-sse/utils/error.ts"; +import { getJobRegistry } from "@/lib/jobRegistry"; + +export const dynamic = "force-dynamic"; + +const DEFAULT_TIMEOUT_MS = 30_000; + +export async function POST(_request: Request, { params }: { params: Promise<{ id: string }> }) { + try { + const { id } = await params; + const registry = getJobRegistry(); + if (!registry.listJobs().some((j) => j.id === id)) { + return NextResponse.json(buildErrorBody(404, "Job not found"), { status: 404 }); + } + const timeoutMs = Number(process.env.OMNIROUTE_RUNNOW_TIMEOUT_MS) || DEFAULT_TIMEOUT_MS; + // Clear the loser: Promise.race settles on the first result but leaves the + // other timer armed, so without this every call keeps a live timeout for + // the full window even though it resolved in milliseconds. + let timer: ReturnType | undefined; + try { + const result = await Promise.race([ + registry.runNow(id), + new Promise((_, reject) => { + timer = setTimeout( + () => reject(new Error(`runNow timed out after ${timeoutMs}ms`)), + timeoutMs + ); + }), + ]); + return NextResponse.json({ data: result }); + } finally { + if (timer) clearTimeout(timer); + } + } catch (err) { + console.error("[API] POST /api/jobs/:id/run-now error:", err); + return NextResponse.json(buildErrorBody(500, "Failed to run job"), { status: 500 }); + } +} diff --git a/src/app/api/jobs/[id]/runs/route.ts b/src/app/api/jobs/[id]/runs/route.ts new file mode 100644 index 0000000000..8e17b5c9c0 --- /dev/null +++ b/src/app/api/jobs/[id]/runs/route.ts @@ -0,0 +1,26 @@ +/** + * GET /api/jobs/:id/runs + * + * Return run history for a single job (newest-first). LOCAL_ONLY. + * Next 16 async params: `const { id } = await params`. + */ +import { NextResponse } from "next/server"; + +import { buildErrorBody } from "@omniroute/open-sse/utils/error.ts"; +import { getJobRegistry } from "@/lib/jobRegistry"; + +export const dynamic = "force-dynamic"; + +export async function GET(_request: Request, { params }: { params: Promise<{ id: string }> }) { + try { + const { id } = await params; + const registry = getJobRegistry(); + if (!registry.listJobs().some((j) => j.id === id)) { + return NextResponse.json(buildErrorBody(404, "Job not found"), { status: 404 }); + } + return NextResponse.json({ data: registry.getRuns(id) }); + } catch (err) { + console.error("[API] GET /api/jobs/:id/runs error:", err); + return NextResponse.json(buildErrorBody(500, "Failed to load runs"), { status: 500 }); + } +} diff --git a/src/app/api/jobs/route.ts b/src/app/api/jobs/route.ts new file mode 100644 index 0000000000..e2505fb446 --- /dev/null +++ b/src/app/api/jobs/route.ts @@ -0,0 +1,37 @@ +/** + * GET /api/jobs + * + * List all registered jobs with their last run. LOCAL_ONLY - loopback enforced by + * routeGuard's isLocalOnlyPath() before this handler runs. + * + * Response: { data: JobDto[] } - DTO whitelist (no handler/timer, which are + * non-serializable live objects). + */ +import { NextResponse } from "next/server"; + +import { buildErrorBody } from "@omniroute/open-sse/utils/error.ts"; +import { getJobRegistry } from "@/lib/jobRegistry"; + +export const dynamic = "force-dynamic"; + +export async function GET() { + try { + const registry = getJobRegistry(); + const jobs = registry.listJobs().map((job) => ({ + id: job.id, + type: job.type, + cron: job.cron, + intervalMs: job.intervalMs, + enabled: job.enabled, + envFlag: job.envFlag, + config: job.config, + createdAt: job.createdAt, + updatedAt: job.updatedAt, + lastRun: registry.getRuns(job.id, 1)[0] ?? null, + })); + return NextResponse.json({ data: jobs }); + } catch (err) { + console.error("[API] GET /api/jobs error:", err); + return NextResponse.json(buildErrorBody(500, "Failed to list jobs"), { status: 500 }); + } +} diff --git a/src/lib/db/jobRegistryDb.ts b/src/lib/db/jobRegistryDb.ts new file mode 100644 index 0000000000..da4dbfa1a3 --- /dev/null +++ b/src/lib/db/jobRegistryDb.ts @@ -0,0 +1,175 @@ +/** + * JobRegistry persistence layer. + * + * CRUD for the `jobs` and `job_runs` tables (migration 136). All timestamps are + * ISO-8601 strings; the registry computes thresholds in JS (never SQL datetime + * arithmetic) so comparisons are plain string compares. + * + * Column naming follows the rest of src/lib/db: snake_case in SQLite, camelCase in + * the returned objects (mapRow). `enabled` is stored INTEGER (0/1), `config` is a + * JSON string parsed/stringified at the boundary. + */ + +import { getDbInstance } from "./core"; +import type { JobRecord, JobRun } from "../jobRegistry/core"; + +function mapJob(row: any): JobRecord { + return { + id: row.id, + type: row.type, + cron: row.cron, + intervalMs: row.interval_ms, + enabled: row.enabled === 1, + envFlag: row.env_flag, + config: row.config ? JSON.parse(row.config) : {}, + createdAt: row.created_at, + updatedAt: row.updated_at, + }; +} + +function mapRun(row: any): JobRun { + return { + id: row.id, + jobId: row.job_id, + startedAt: row.started_at, + finishedAt: row.finished_at, + status: row.status, + errorMessage: row.error_message, + recordsAffected: row.records_affected ?? 0, + durationMs: row.duration_ms, + }; +} + +export function getAllJobs(): JobRecord[] { + const db = getDbInstance(); + const rows = db.prepare("SELECT * FROM jobs ORDER BY id").all(); + return rows.map(mapJob); +} + +export function getJob(id: string): JobRecord | null { + const db = getDbInstance(); + const row = db.prepare("SELECT * FROM jobs WHERE id = ?").get(id); + return row ? mapJob(row) : null; +} + +/** + * Idempotent register: INSERT OR IGNORE on first sight, then a column-level UPDATE + * that refreshes scheduling fields (type/cron/interval/env_flag/config) and bumps + * updated_at - but NEVER overwrites `enabled` (the user's API-driven toggle) nor + * `created_at`. Pass enabled=true for new jobs; the UPDATE simply skips the column. + */ +export function upsertJob(job: JobRecord): void { + const db = getDbInstance(); + db.prepare( + `INSERT INTO jobs (id, type, cron, interval_ms, enabled, env_flag, config, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, datetime('now')) + ON CONFLICT(id) DO UPDATE SET + type = excluded.type, + cron = excluded.cron, + interval_ms = excluded.interval_ms, + env_flag = excluded.env_flag, + config = excluded.config, + updated_at = datetime('now')` + ).run( + job.id, + job.type, + job.cron, + job.intervalMs, + job.enabled ? 1 : 0, + job.envFlag, + JSON.stringify(job.config ?? {}), + job.createdAt + ); +} + +export function updateJobEnabled(id: string, enabled: boolean): void { + const db = getDbInstance(); + db.prepare("UPDATE jobs SET enabled = ?, updated_at = datetime('now') WHERE id = ?").run( + enabled ? 1 : 0, + id + ); +} + +export interface RecordRunOptions { + startedAt?: string; + durationMs?: number; + errorMessage?: string; + recordsAffected?: number; +} + +/** + * Insert a completed run. `startedAt` is captured by the registry before the handler + * runs; `finishedAt` is derived from startedAt + durationMs so the two never drift. + * A status='running' insert leaves finishedAt NULL (used to mark in-flight work). + */ +export function recordRun( + jobId: string, + status: JobRun["status"], + opts: RecordRunOptions = {} +): void { + const db = getDbInstance(); + const startedAt = opts.startedAt ?? new Date().toISOString(); + const finishedAt = + status === "running" + ? null + : opts.startedAt && opts.durationMs != null + ? new Date(new Date(opts.startedAt).getTime() + opts.durationMs).toISOString() + : new Date().toISOString(); + db.prepare( + `INSERT INTO job_runs (job_id, started_at, finished_at, status, error_message, records_affected, duration_ms) + VALUES (?, ?, ?, ?, ?, ?, ?)` + ).run( + jobId, + startedAt, + finishedAt, + status, + opts.errorMessage ?? null, + opts.recordsAffected ?? 0, + opts.durationMs ?? null + ); +} + +/** + * Dual-dimension pruning: keep the most recent `maxRuns` OR anything younger than + * `maxDays`. A run is deleted only when it is BOTH outside the recent-N window AND + * older than maxDays - so a job that runs rarely keeps its full history. + * All thresholds are ISO-8601 string compares (no SQL datetime math). + */ +export function pruneRuns(jobId: string, maxRuns = 100, maxDays = 30): void { + const db = getDbInstance(); + const threshold = new Date(Date.now() - maxDays * 86_400_000).toISOString(); + db.prepare( + `DELETE FROM job_runs + WHERE job_id = ? + AND id NOT IN ( + SELECT id FROM job_runs WHERE job_id = ? ORDER BY started_at DESC LIMIT ? + ) + AND started_at < ?` + ).run(jobId, jobId, maxRuns, threshold); +} + +export function getRuns(jobId: string, limit = 20): JobRun[] { + const db = getDbInstance(); + const rows = db + .prepare("SELECT * FROM job_runs WHERE job_id = ? ORDER BY started_at DESC LIMIT ?") + .all(jobId, limit); + return rows.map(mapRun); +} + +/** + * Startup repair: any run still marked `running` whose started_at is older than + * `timeoutMinutes` is a leftover from a crashed/hung process. Mark it `failure` so + * the history is accurate and the slot frees up. The 5-minute default avoids + * clobbering a genuinely in-flight run on a slow box. + */ +export function cleanupOrphanedRuns(timeoutMinutes = 5): void { + const db = getDbInstance(); + const threshold = new Date(Date.now() - timeoutMinutes * 60_000).toISOString(); + db.prepare( + `UPDATE job_runs + SET status = 'failure', + error_message = 'orphaned: exceeded timeout', + finished_at = datetime('now') + WHERE status = 'running' AND started_at < ?` + ).run(threshold); +} diff --git a/src/lib/db/migrations/139_job_registry.sql b/src/lib/db/migrations/139_job_registry.sql new file mode 100644 index 0000000000..3bf29ac3d2 --- /dev/null +++ b/src/lib/db/migrations/139_job_registry.sql @@ -0,0 +1,48 @@ +-- Migration 139: generic job registry (jobs + job_runs tables) +-- Job registry (#8848): centralized periodic-job scheduling + run history. +-- +-- jobs: one row per registered job (interval or cron), with env-flag gating +-- job_runs: one row per execution (running/success/failure), pruned by count + age + +CREATE TABLE IF NOT EXISTS jobs ( + id TEXT PRIMARY KEY, + type TEXT NOT NULL DEFAULT 'interval' CHECK(type IN ('interval', 'cron')), + cron TEXT, -- cron expression (type='cron'); NULL for interval jobs + interval_ms INTEGER, -- interval in ms (type='interval'); NULL for cron jobs + enabled INTEGER NOT NULL DEFAULT 1, -- 0=disabled, 1=enabled + env_flag TEXT, -- env var name (boolean gate), e.g. 'OMNIROUTE_WARMUP_ENABLED'; NULL = no gate + config TEXT NOT NULL DEFAULT '{}', -- JSON config (concurrency, timezone, envDefault, ...) + created_at TEXT NOT NULL DEFAULT (datetime('now')), + updated_at TEXT NOT NULL DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS job_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + job_id TEXT NOT NULL REFERENCES jobs(id) ON DELETE CASCADE, + started_at TEXT NOT NULL, -- ISO-8601, written explicitly by recordRun (no DB default) + finished_at TEXT, -- ISO-8601; NULL while running + status TEXT NOT NULL DEFAULT 'running', -- 'running' | 'success' | 'failure' + error_message TEXT, -- sanitized error message (no stack trace) + records_affected INTEGER DEFAULT 0, -- job-specific meaning (see below) + duration_ms INTEGER, -- execution duration in ms + created_at TEXT NOT NULL DEFAULT (datetime('now')) +); +CREATE INDEX IF NOT EXISTS idx_jr_job_id ON job_runs(job_id, started_at DESC); +CREATE INDEX IF NOT EXISTS idx_jr_started_at ON job_runs(started_at); + +-- Built-in job registration (idempotent - INSERT OR IGNORE). +-- env_flag = NULL means "no registry-level boolean gate"; per-job disable semantics +-- live inside the handler itself (see token_health_check wrapper). +-- 'warmup' is seeded disabled because its handler is not part of this change. +-- startAll() filters on enabled before it looks for a handler, so a disabled row +-- stays quiet instead of warning on every boot; the change that brings the warmup +-- handler flips it on. +INSERT OR IGNORE INTO jobs (id, type, cron, interval_ms, enabled, env_flag, config) VALUES + ('budget_reset', 'interval', NULL, 600000, 1, NULL, '{}'), + ('warmup', 'cron', '0 7 * * *', NULL, 0, 'OMNIROUTE_WARMUP_ENABLED', '{"timezone":"America/Los_Angeles","envDefault":false}'), + ('token_health_check', 'interval', NULL, 60000, 1, NULL, '{}'); + +-- records_affected semantics: +-- budget_reset = number of budget records reset (UPDATE ... SET budget_used=0 row count) +-- warmup = number of connections attempted for warmup +-- token_health_check = number of connections swept by the health check diff --git a/src/lib/initCloudSync.ts b/src/lib/initCloudSync.ts index 9032452cd0..3464ea0d29 100644 --- a/src/lib/initCloudSync.ts +++ b/src/lib/initCloudSync.ts @@ -1,8 +1,9 @@ import initializeCloudSync from "@/shared/services/initializeCloudSync"; -import { startBudgetResetJob } from "@/lib/jobs/budgetResetJob"; import { startModelSyncScheduler } from "@/shared/services/modelSyncScheduler"; -import { startWarmupScheduler } from "@/lib/warmupScheduler"; import { isAutomatedTestProcess } from "@/shared/utils/testProcess"; +import { getJobRegistry } from "@/lib/jobRegistry"; +import { registerBudgetResetJob } from "@/lib/jobs/budgetResetJob"; +import { registerTokenHealthCheck } from "@/lib/jobs/tokenHealthCheckJob"; // Initialize runtime background sync services once per server process. let initialized = false; @@ -29,12 +30,19 @@ export async function ensureCloudSyncInitialized() { } if (!initialized) { try { - const { initTokenHealthCheck } = await import("@/lib/tokenHealthCheck"); - initTokenHealthCheck(); await initializeCloudSync(); startModelSyncScheduler(); - startBudgetResetJob(); - startWarmupScheduler(); + + // startAll() runs each interval job's first tick synchronously, so it has to + // come after initializeCloudSync(). The old wiring got that ordering two + // different ways: the budget reset was started right here, and the health + // check's first sweep sat behind a 10s timer. Awaiting the init is a firmer + // guarantee than the timer was. + const registry = getJobRegistry(); + registerBudgetResetJob(registry); + registerTokenHealthCheck(registry); + await registry.startAll(); + initialized = true; } catch (error) { console.error("[ServerInit] Error initializing background sync services:", error); diff --git a/src/lib/jobRegistry/core.ts b/src/lib/jobRegistry/core.ts new file mode 100644 index 0000000000..ba8ecd770b --- /dev/null +++ b/src/lib/jobRegistry/core.ts @@ -0,0 +1,51 @@ +/** JobRegistry shared types + env-gate helper . */ + +export interface JobRecord { + id: string; + type: "interval" | "cron"; + cron: string | null; + intervalMs: number | null; + enabled: boolean; + envFlag: string | null; + config: Record; + createdAt: string; + updatedAt: string; +} + +export interface JobDefinition extends JobRecord { + handler: () => Promise; + cronGetter?: () => string; +} + +export interface HandlerResult { + success: boolean; + recordsAffected?: number; + error?: string; +} + +export interface JobRun { + id: number; + jobId: string; + startedAt: string; + finishedAt: string | null; + status: "running" | "success" | "failure"; + errorMessage: string | null; + recordsAffected: number; + durationMs: number | null; +} + +const TRUE_ENV_VALUES = new Set(["1", "true", "yes", "on"]); + +/** + * Boolean env-var gate used by cron jobs (e.g. warmup). + * + * - No envFlag -> always enabled (no gate). + * - envFlag, env unset -> fall back to `defaultWhenUnset` (job-level default; warmup=false). + * - envFlag, env set -> truthy check against TRUE_ENV_VALUES. + */ +export function isEnvEnabled(envName: string | undefined, defaultWhenUnset = true): boolean { + if (!envName) return true; + const v = process.env[envName]; + if (v === undefined) return defaultWhenUnset; + return TRUE_ENV_VALUES.has(v.trim().toLowerCase()); +} diff --git a/src/lib/jobRegistry/index.ts b/src/lib/jobRegistry/index.ts new file mode 100644 index 0000000000..fc9b56e0a2 --- /dev/null +++ b/src/lib/jobRegistry/index.ts @@ -0,0 +1,22 @@ +/** JobRegistry singleton - survives Next.js HMR via globalThis. */ + +import { JobRegistry } from "./registry"; +import type { JobDefinition } from "./core"; + +declare global { + var __omnirouteJobRegistry: JobRegistry | undefined; +} + +export function getJobRegistry(): JobRegistry { + if (!globalThis.__omnirouteJobRegistry) { + globalThis.__omnirouteJobRegistry = new JobRegistry(); + } + return globalThis.__omnirouteJobRegistry; +} + +/** Test-only: drop the singleton so each test starts fresh. */ +export function __resetJobRegistry(): void { + globalThis.__omnirouteJobRegistry = undefined; +} + +export type { JobDefinition, JobRecord, HandlerResult, JobRun } from "./core"; diff --git a/src/lib/jobRegistry/registry.ts b/src/lib/jobRegistry/registry.ts new file mode 100644 index 0000000000..3c7c227c5d --- /dev/null +++ b/src/lib/jobRegistry/registry.ts @@ -0,0 +1,271 @@ +/** JobRegistry - unified scheduler for all periodic background jobs. + * + * Two strategies: interval (setInterval) and cron (nextTick via cron-parser, DST-safe). + * Cron re-arms the next fire BEFORE running the current handler, so a slow handler + * cannot kill the recursion chain. + */ + +import { CronExpressionParser } from "cron-parser"; +import { + getAllJobs, + getJob, + recordRun, + pruneRuns, + cleanupOrphanedRuns, + upsertJob, + updateJobEnabled, + getRuns as dbGetRuns, +} from "../db/jobRegistryDb"; +import { isEnvEnabled } from "./core"; +import { sanitizeErrorMessage } from "@omniroute/open-sse/utils/error.ts"; +import type { JobDefinition, HandlerResult, JobRun } from "./core"; + +function errMessage(err: unknown): string { + return err instanceof Error ? err.message : String(err); +} + +export class JobRegistry { + private timers = new Map>(); + private running = new Set(); + private handlers = new Map Promise>(); + private queued = new Set(); + private waiters = new Map void>>(); + private cronFailCount = new Map(); + private static readonly MAX_CRON_PARSE_FAILURES = 5; + private cronGetters = new Map string>(); + register(def: JobDefinition): void { + // Never clobber enabled toggle or created_at on re-register. + // the original created_at. Pass enabled=true for new jobs (seed default). + upsertJob({ + id: def.id, + type: def.type, + cron: def.cron, + intervalMs: def.intervalMs, + enabled: getJob(def.id)?.enabled ?? def.enabled ?? true, + envFlag: def.envFlag, + config: def.config, + createdAt: getJob(def.id)?.createdAt ?? new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + this.handlers.set(def.id, def.handler); + if (def.cronGetter) this.cronGetters.set(def.id, def.cronGetter); + } + + start(jobId: string): void { + if (this.timers.has(jobId)) return; + const job = getJob(jobId); + const handler = this.handlers.get(jobId); + if (!job || !handler) { + console.warn(`[JobRegistry] Skip ${jobId}: no ${job ? "handler" : "job"} registered`); + return; + } + if (job.type === "interval") { + this.startInterval(jobId, handler, job.intervalMs ?? 60_000); + } else { + const tz = (job.config?.timezone as string) || "UTC"; + let cronExpr = job.cron ?? "* * * * *"; + const getter = this.cronGetters.get(jobId); + if (getter) { + try { + cronExpr = getter(); + } catch (err) { + console.error(`[JobRegistry] cronGetter for ${jobId} threw:`, errMessage(err)); + cronExpr = job.cron ?? "* * * * *"; + } + } + this.scheduleNextCron(jobId, cronExpr, tz, job.envFlag); + } + } + + stop(jobId: string): void { + const timer = this.timers.get(jobId); + if (timer) { + clearTimeout(timer); + clearInterval(timer); + this.timers.delete(jobId); + } + } + + async startAll(): Promise { + if (this.handlers.size === 0) { + throw new Error("[JobRegistry] No handlers registered. Call register() before startAll()."); + } + await cleanupOrphanedRuns(); + const jobs = getAllJobs().filter((j) => j.enabled); + // start() is synchronous and can throw (getJob reads the DB), so each call + // needs its own try/catch. Collecting them with Promise.allSettled would not + // work: the callback returns void, every entry settles as fulfilled, and a + // synchronous throw escapes the map before allSettled is ever reached. + for (const job of jobs) { + if (!this.handlers.has(job.id)) { + console.warn(`[JobRegistry] Skip ${job.id}: no handler registered`); + continue; + } + try { + this.start(job.id); + } catch (err) { + console.error(`[JobRegistry] Failed to start ${job.id}:`, errMessage(err)); + } + } + } + + stopAll(): void { + for (const jobId of [...this.timers.keys()]) this.stop(jobId); + } + + dispose(): void { + this.stopAll(); + this.cronFailCount.clear(); + } + + async runNow(jobId: string): Promise<{ started: boolean; reason?: string }> { + const job = getJob(jobId); + if (!job) return { started: false, reason: "not_found" }; + if (!job.enabled) return { started: false, reason: "disabled" }; + const handler = this.handlers.get(jobId); + if (!handler) return { started: false, reason: "no_handler" }; + + if (this.running.has(jobId)) { + if (this.queued.has(jobId)) return { started: false, reason: "already_queued" }; + this.queued.add(jobId); + return new Promise<{ started: boolean; reason?: string }>((resolve) => { + const waiters = this.waiters.get(jobId) ?? []; + waiters.push(() => { + this.queued.delete(jobId); + resolve(this.runNow(jobId)); + }); + this.waiters.set(jobId, waiters); + }); + } + + if ( + job.envFlag && + !isEnvEnabled(job.envFlag, job.config?.envDefault === false ? false : true) + ) { + return { started: false, reason: "env_disabled" }; + } + void this.safeRun(jobId, handler); + return { started: true }; + } + + setEnabled(jobId: string, enabled: boolean): void { + updateJobEnabled(jobId, enabled); + if (enabled) { + this.stop(jobId); + this.start(jobId); + } else { + this.stop(jobId); + } + } + + getRuns(jobId: string, limit = 20): JobRun[] { + return dbGetRuns(jobId, limit); + } + + listJobs(): JobDefinition[] { + return getAllJobs().map((job) => ({ + ...job, + handler: this.handlers.get(job.id) ?? (() => Promise.resolve({ success: true })), + })); + } + + private getJobConfig(jobId: string): Record | null { + return getJob(jobId)?.config ?? null; + } + + /** Core execution wrapper: re-entrancy guard + timing + recording + prune. */ + private async safeRun(jobId: string, handler: () => Promise): Promise { + if (this.running.has(jobId)) return; + this.running.add(jobId); + const startedAt = new Date().toISOString(); + const start = Date.now(); + try { + const result = await handler(); + await recordRun(jobId, result.success ? "success" : "failure", { + startedAt, + durationMs: Date.now() - start, + errorMessage: result.error ? sanitizeErrorMessage(result.error) : undefined, + recordsAffected: result.recordsAffected, + }); + } catch (err) { + await recordRun(jobId, "failure", { + startedAt, + durationMs: Date.now() - start, + errorMessage: sanitizeErrorMessage(errMessage(err)), + }); + } finally { + this.running.delete(jobId); + const ws = this.waiters.get(jobId); + if (ws) { + this.waiters.delete(jobId); + for (const resolve of ws) resolve(); + } + await pruneRuns(jobId); + } + } + + private startInterval( + jobId: string, + handler: () => Promise, + intervalMs: number + ): void { + void this.safeRun(jobId, handler); + const timer = setInterval(() => void this.safeRun(jobId, handler), intervalMs); + timer.unref?.(); + this.timers.set(jobId, timer); + } + + /** + * Cron nextTick: schedule the next fire, and on fire run the handler then + * re-arm. The next timer is armed BEFORE the handler runs (recursion-safe - + * see file header). DST-safe because cron-parser computes the next UTC instant + * from the IANA timezone; we never compare wall-clock fields. + */ + private scheduleNextCron( + jobId: string, + cronExpr: string, + tz: string, + envFlag?: string | null + ): void { + let delay: number; + try { + const interval = CronExpressionParser.parse(cronExpr, { tz }); + const next = interval.next().getTime(); + delay = Math.max(next - Date.now(), 0); + this.cronFailCount.delete(jobId); + } catch (err) { + const fails = (this.cronFailCount.get(jobId) ?? 0) + 1; + this.cronFailCount.set(jobId, fails); + console.error( + `[JobRegistry] Invalid cron "${cronExpr}" for ${jobId} (fail #${fails}):`, + errMessage(err) + ); + if (fails >= JobRegistry.MAX_CRON_PARSE_FAILURES) { + console.error( + `[JobRegistry] Cron "${cronExpr}" for ${jobId} failed ${fails} times - stopping job.` + ); + return; + } + delay = 60_000; + } + const timer = setTimeout(() => { + let freshCron = cronExpr; + const getter = this.cronGetters.get(jobId); + if (getter) { + try { + freshCron = getter(); + } catch (err) { + console.error(`[JobRegistry] cronGetter for ${jobId} threw in timer:`, errMessage(err)); + } + } + this.scheduleNextCron(jobId, freshCron, tz, envFlag); + if (this.running.has(jobId)) return; + const jobConfig = this.getJobConfig(jobId); + if (envFlag && !isEnvEnabled(envFlag, jobConfig?.envDefault === false ? false : true)) return; + const handler = this.handlers.get(jobId); + if (handler) void this.safeRun(jobId, handler); + }, delay); + timer.unref?.(); + this.timers.set(jobId, timer); + } +} diff --git a/src/lib/jobRegistry/timeUtils.ts b/src/lib/jobRegistry/timeUtils.ts new file mode 100644 index 0000000000..cec530a65c --- /dev/null +++ b/src/lib/jobRegistry/timeUtils.ts @@ -0,0 +1,36 @@ +/** + * Timezone conversion helpers for JobRegistry. + * + * Generalized from warmupScheduler.toPacificTime(): converts a UTC Date into the + * wall-clock Date of any IANA timezone. Used by cron scheduling to evaluate + * "what time is it in the job's timezone" without pulling in a heavy tz library. + */ + +/** + * Convert a Date to the same wall-clock time in the given IANA timezone. + * Returns a Date whose local fields (getHours() etc.) read as the target zone's + * wall clock. Pure Intl - no DST logic here; cron-parser handles DST on the + * scheduling side. + */ +export function convertToTimeZone(date: Date, timeZone: string): Date { + const parts = new Intl.DateTimeFormat("en-US", { + timeZone, + hour12: false, + hourCycle: "h23", + year: "numeric", + month: "numeric", + day: "numeric", + hour: "numeric", + minute: "numeric", + second: "numeric", + }).formatToParts(date); + const get = (t: string) => parseInt(parts.find((p) => p.type === t)?.value || "0", 10); + return new Date( + get("year"), + get("month") - 1, + get("day"), + get("hour"), + get("minute"), + get("second") + ); +} diff --git a/src/lib/jobs/budgetResetJob.ts b/src/lib/jobs/budgetResetJob.ts index bd640569e8..f29ae383a1 100644 --- a/src/lib/jobs/budgetResetJob.ts +++ b/src/lib/jobs/budgetResetJob.ts @@ -1,40 +1,42 @@ import { syncAllBudgetSchedules } from "@/domain/costRules"; +import type { JobRegistry } from "../jobRegistry/registry"; const DEFAULT_INTERVAL_MS = 10 * 60 * 1000; -let timer: NodeJS.Timeout | null = null; - function getIntervalMs() { const raw = process.env.OMNIROUTE_BUDGET_RESET_JOB_INTERVAL_MS; const parsed = raw ? Number(raw) : Number.NaN; return Number.isFinite(parsed) && parsed >= 10_000 ? parsed : DEFAULT_INTERVAL_MS; } -export function startBudgetResetJob() { - if (timer) { - return timer; +/** + * Budget-reset handler - resets all budget schedules once per interval. + * + * Migrated to the JobRegistry: the old start/stop timer logic is + * gone; the registry owns scheduling + run recording. This module exports just the + * `run` handler and a `registerBudgetResetJob` wiring helper. Errors are no longer + * caught here - safeRun records them as a failure run. + */ +export async function run(): Promise<{ success: boolean; recordsAffected: number }> { + const result = syncAllBudgetSchedules(Date.now()); + if (result.resetCount > 0) { + console.log(`[BudgetReset] processed=${result.processed} reset=${result.resetCount}`); } - - const run = () => { - try { - const result = syncAllBudgetSchedules(Date.now()); - if (result.resetCount > 0) { - console.log(`[BudgetReset] processed=${result.processed} reset=${result.resetCount}`); - } - } catch (error) { - console.error("[BudgetReset] Job failed:", error); - } - }; - - run(); - timer = setInterval(run, getIntervalMs()); - timer.unref?.(); - return timer; + return { success: true, recordsAffected: result.resetCount }; } -export function stopBudgetResetJob() { - if (timer) { - clearInterval(timer); - timer = null; - } +/** Wire the budget_reset job into a JobRegistry (idempotent - call at boot). */ +export function registerBudgetResetJob(registry: JobRegistry): void { + registry.register({ + id: "budget_reset", + type: "interval", + cron: null, + intervalMs: getIntervalMs(), + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: run, + }); } diff --git a/src/lib/jobs/tokenHealthCheckJob.ts b/src/lib/jobs/tokenHealthCheckJob.ts new file mode 100644 index 0000000000..73e15d784c --- /dev/null +++ b/src/lib/jobs/tokenHealthCheckJob.ts @@ -0,0 +1,40 @@ +/** + * Registry adapter for the token health check. + * + * The sweep used to start itself on import. The registry owns the schedule now, + * so the adapter lives here next to the other job registrations rather than in + * tokenHealthCheck.ts, which is already at its size ceiling. + * + * Disable semantics are unchanged: isHealthCheckDisabled() still honours + * OMNIROUTE_DISABLE_TOKEN_HEALTHCHECK, NEXT_PHASE=phase-production-build and + * isAutomatedTestProcess(). The registry only fires the interval; the handler + * decides whether there is anything to do. + */ +import { isHealthCheckDisabled, sweep } from "@/lib/tokenHealthCheck"; +import type { JobRegistry } from "@/lib/jobRegistry/registry"; + +const TOKEN_HEALTH_CHECK_INTERVAL_MS = 60_000; + +export function registerTokenHealthCheck(registry: JobRegistry): void { + const now = new Date().toISOString(); + registry.register({ + id: "token_health_check", + type: "interval", + cron: null, + intervalMs: TOKEN_HEALTH_CHECK_INTERVAL_MS, + enabled: true, + envFlag: null, + config: {}, + createdAt: now, + updatedAt: now, + handler: async () => { + if (isHealthCheckDisabled()) { + return { success: true, recordsAffected: 0 }; + } + // Errors are not caught here: safeRun records a thrown error as a failure run + // with its message, the same as the budget reset job. + const swept = await sweep(); + return { success: true, recordsAffected: swept }; + }, + }); +} diff --git a/src/lib/localDb.ts b/src/lib/localDb.ts index c2a6b834df..99993f2538 100755 --- a/src/lib/localDb.ts +++ b/src/lib/localDb.ts @@ -94,7 +94,7 @@ export * from "./db/compressionCacheStats"; export * from "./db/compressionCombos"; export * from "./db/compressionContextBudget"; export * from "./db/compressionRunTelemetry"; -export * from "./db/connectionRuntimeState"; +export * from "./db/jobRegistryDb"; export * from "./db/modelContextOverrides"; export { @@ -818,7 +818,5 @@ export { getRadarSettings, setRadarOptIn, setRadarKey, - getRadarReferralsCache, - setRadarReferralsCache, } from "./db/radar"; -export type { RadarCache, RadarSettings, RadarReferralsCache } from "./db/radar"; +export type { RadarCache, RadarSettings } from "./db/radar"; diff --git a/src/lib/tokenHealthCheck.ts b/src/lib/tokenHealthCheck.ts index 5039e26419..73dd639d91 100644 --- a/src/lib/tokenHealthCheck.ts +++ b/src/lib/tokenHealthCheck.ts @@ -267,7 +267,7 @@ function isEnvFlagEnabled(name: string): boolean { return TRUE_ENV_VALUES.has(value.trim().toLowerCase()); } -function isHealthCheckDisabled(): boolean { +export function isHealthCheckDisabled(): boolean { return ( isEnvFlagEnabled("OMNIROUTE_DISABLE_TOKEN_HEALTHCHECK") || isBuildProcess() || @@ -421,16 +421,18 @@ export function stopTokenHealthCheck() { } // ── Core sweep (batch concurrent) ────────────────────────────────────────── -export async function sweep() { +/** Returns the number of connections swept, which the job registry records. */ +export async function sweep(): Promise { const state = getHCState(); if (state.sweeping) { - return log(`${LOG_PREFIX} Sweep skipped — previous sweep still in progress`); + log(`${LOG_PREFIX} Sweep skipped — previous sweep still in progress`); + return 0; } state.sweeping = true; try { const connections = await getProviderConnections({ authType: "oauth" }); - if (!connections || connections.length === 0) return; + if (!connections || connections.length === 0) return 0; const staggerMs = parseInt(process.env.HEALTHCHECK_STAGGER_MS || "3000", 10); const total = connections.length; @@ -471,8 +473,10 @@ export async function sweep() { await new Promise((resolve) => setTimeout(resolve, 0)); } } + return total; } catch (err) { logError(`${LOG_PREFIX} Sweep error:`, err.message); + return 0; } finally { state.sweeping = false; } @@ -1046,8 +1050,3 @@ export async function checkConnection(conn) { ); } } - -// Auto-start when imported -initTokenHealthCheck(); - -export default initTokenHealthCheck; diff --git a/src/server-init.ts b/src/server-init.ts index bdbcb64dde..04e6190d5a 100644 --- a/src/server-init.ts +++ b/src/server-init.ts @@ -4,7 +4,8 @@ import { enforceWebRuntimeEnv } from "./lib/env/runtimeEnv"; import { enforceSecrets } from "./shared/utils/secretsValidator"; import { initAuditLog, cleanupExpiredLogs, logAuditEvent } from "./lib/compliance/index"; import { initConsoleInterceptor } from "./lib/consoleInterceptor"; -import { startBudgetResetJob } from "./lib/jobs/budgetResetJob"; +import { registerBudgetResetJob } from "./lib/jobs/budgetResetJob"; +import { registerTokenHealthCheck } from "./lib/jobs/tokenHealthCheckJob"; import { startReasoningCacheCleanupJob } from "./lib/jobs/reasoningCacheCleanupJob"; import { startCleanupScheduler } from "./lib/db/cleanup"; import { getSettings } from "./lib/db/settings"; @@ -113,7 +114,16 @@ async function startServer() { } await initializeCloudSync(); - startBudgetResetJob(); + // register() only persists the definition; startAll() is what arms the timers. + // This path does not call ensureCloudSyncInitialized(), so nothing else here + // would start the jobs on our behalf. It registers the same set as that path: + // starting one job and not the other is how a background job goes missing + // without anything failing. + const { getJobRegistry } = await import("./lib/jobRegistry"); + const jobRegistry = getJobRegistry(); + registerBudgetResetJob(jobRegistry); + registerTokenHealthCheck(jobRegistry); + await jobRegistry.startAll(); startReasoningCacheCleanupJob(); startCleanupScheduler(); startRuntimeConfigHotReload(); diff --git a/src/server/authz/routeGuard.ts b/src/server/authz/routeGuard.ts index be7220b00a..e61dcf0bf1 100644 --- a/src/server/authz/routeGuard.ts +++ b/src/server/authz/routeGuard.ts @@ -50,6 +50,8 @@ export const LOCAL_ONLY_API_PREFIXES: ReadonlyArray = [ "/api/local/", // T-12: 1-click local service launchers (Redis today; spawns podman/docker) — loopback-enforced by isLocalRequestAllowed() in src/lib/security/localEndpoints.ts (Hard Rules #15 + #17) "/api/headroom/start", // Headroom token-saver proxy lifecycle: spawns headroom-ai python CLI (Hard Rules #15 + #17) "/api/headroom/stop", // Headroom token-saver proxy lifecycle: sends SIGTERM/SIGKILL to managed PID (Hard Rules #15 + #17) + "/api/jobs", // JobRegistry control (enable/disable/run-now) + run history - runtime job administration, loopback-only (Hard Rules #15 + #17) + "/api/jobs/", // sub-paths: /api/jobs/:id/{runs,enable,disable,run-now} (the bare `/api/jobs` above matches the list route; this matches children) "/api/oauth/cursor/auto-import", // spawns `execFile("which", ["cursor"])` to verify a local Cursor install before importing creds — RCE-via-tunnel surface (Hard Rules #15 + #17, found by 6A.8 route-guard gate). Specific path only: the rest of /api/oauth/ (browser redirect/callback flows) must stay remote-reachable. "/api/skills/collect/", // Skill Collector CLI detection: GET .../detect probes getCliRuntimeStatus() per CLI_TOOL_IDS entry, which spawns a child process to check each tool — RCE-via-tunnel surface (Hard Rules #15 + #17, PR #6294 review). "/api/discovery/", // Discovery tool (opt-in provider scanner): the scan route makes outbound probes to provider endpoints (SSRF-adjacent) and the whole surface is an admin research tool — strict-loopback only, no manage-scope bypass (NOT in LOCAL_ONLY_MANAGE_SCOPE_BYPASS_PREFIXES). See _tasks/features-v3.8.42/gaps/DISCOVERY_TOOL_DESIGN.md. @@ -174,9 +176,7 @@ export function isPrivateLanHost(hostHeader: string | null): boolean { * triggers the auto-update flow (spawns git checkout + npm install + pm2). * Hard Rules #15/#17 still apply to POST. */ -export const LOCAL_ONLY_API_GET_EXEMPTIONS: ReadonlySet = new Set([ - "/api/system/version", -]); +export const LOCAL_ONLY_API_GET_EXEMPTIONS: ReadonlySet = new Set(["/api/system/version"]); /** Safe HTTP methods that can be exempted for read-only paths. */ const SAFE_METHODS = new Set(["GET", "HEAD", "OPTIONS"]); diff --git a/tests/unit/api/jobs.test.ts b/tests/unit/api/jobs.test.ts new file mode 100644 index 0000000000..53ff2b1474 --- /dev/null +++ b/tests/unit/api/jobs.test.ts @@ -0,0 +1,210 @@ +/** /api/jobs route tests . */ + +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-jobs-api-")); +process.env.DATA_DIR = TEST_DATA_DIR; +process.env.DISABLE_SQLITE_AUTO_BACKUP = "true"; +process.env.OMNIROUTE_WARMUP_ENABLED = "1"; + +const core = await import("../../../src/lib/db/core.ts"); +const { getJobRegistry, __resetJobRegistry } = + await import("../../../src/lib/jobRegistry/index.ts"); +const route = await import("../../../src/app/api/jobs/route.ts"); +const runsRoute = await import("../../../src/app/api/jobs/[id]/runs/route.ts"); +const enableRoute = await import("../../../src/app/api/jobs/[id]/enable/route.ts"); +const disableRoute = await import("../../../src/app/api/jobs/[id]/disable/route.ts"); +const runNowRoute = await import("../../../src/app/api/jobs/[id]/run-now/route.ts"); +const { isLocalOnlyPath } = await import("../../../src/server/authz/routeGuard.ts"); + +function resetAll() { + try { + getJobRegistry().stopAll(); + } catch { + // no singleton yet + } + __resetJobRegistry(); + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); + fs.mkdirSync(TEST_DATA_DIR, { recursive: true }); +} + +test.beforeEach(() => { + resetAll(); +}); + +test.after(() => { + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); +}); + +function params(id: string) { + return Promise.resolve({ id }); +} + +async function json(res: Response) { + return res.json(); +} + +test("GET /api/jobs -> 200 + list with lastRun, seeds present", async () => { + const reg = getJobRegistry(); + // Register a custom job so listJobs reflects runtime registrations. + reg.register({ + id: "custom", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => ({ success: true }), + }); + const res = await route.GET(); + assert.equal(res.status, 200); + const body = (await json(res)) as { data: Array<{ id: string; lastRun: unknown }> }; + assert.ok(Array.isArray(body.data)); + const ids = body.data.map((j) => j.id); + assert.ok(ids.includes("budget_reset"), "seeded budget_reset present"); + assert.ok(ids.includes("warmup"), "seeded warmup present"); + assert.ok(ids.includes("custom"), "runtime-registered custom present"); + // DTO whitelist: no handler/timer leaked. + for (const job of body.data) { + assert.ok(!("handler" in job), "DTO must not expose handler"); + assert.ok(Object.prototype.hasOwnProperty.call(job, "lastRun"), "lastRun present"); + } +}); + +test("GET /api/jobs/:id/runs -> 200 + history", async () => { + const reg = getJobRegistry(); + reg.register({ + id: "custom", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => ({ success: true, recordsAffected: 2 }), + }); + await reg.runNow("custom"); + await new Promise((r) => setTimeout(r, 30)); + const res = await runsRoute.GET(new Request("http://localhost/api/jobs/custom/runs"), { + params: params("custom"), + }); + assert.equal(res.status, 200); + const body = (await json(res)) as { data: Array<{ status: string }> }; + assert.ok(body.data.length >= 1); + assert.equal(body.data[0].status, "success"); +}); + +test("GET /api/jobs/unknown/runs -> 404", async () => { + const res = await runsRoute.GET(new Request("http://localhost/api/jobs/nope/runs"), { + params: params("nope"), + }); + assert.equal(res.status, 404); +}); + +test("POST /api/jobs/:id/enable -> 200 + enabled:true", async () => { + const reg = getJobRegistry(); + reg.register({ + id: "custom", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: false, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => ({ success: true }), + }); + const res = await enableRoute.POST( + new Request("http://localhost/api/jobs/custom/enable", { method: "POST" }), + { params: params("custom") } + ); + assert.equal(res.status, 200); + const body = (await json(res)) as { data: { id: string; enabled: boolean } }; + assert.deepEqual(body.data, { id: "custom", enabled: true }); + assert.equal(reg.listJobs().find((j) => j.id === "custom")!.enabled, true); +}); + +test("POST /api/jobs/:id/disable -> 200 + enabled:false", async () => { + const reg = getJobRegistry(); + reg.register({ + id: "custom", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => ({ success: true }), + }); + const res = await disableRoute.POST( + new Request("http://localhost/api/jobs/custom/disable", { method: "POST" }), + { params: params("custom") } + ); + assert.equal(res.status, 200); + const body = (await json(res)) as { data: { id: string; enabled: boolean } }; + assert.deepEqual(body.data, { id: "custom", enabled: false }); + assert.equal(reg.listJobs().find((j) => j.id === "custom")!.enabled, false); +}); + +test("POST /api/jobs/:id/run-now -> 200 + started:true", async () => { + const reg = getJobRegistry(); + reg.register({ + id: "custom", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => ({ success: true }), + }); + const res = await runNowRoute.POST( + new Request("http://localhost/api/jobs/custom/run-now", { method: "POST" }), + { params: params("custom") } + ); + assert.equal(res.status, 200); + const body = (await json(res)) as { data: { started: boolean } }; + assert.equal(body.data.started, true); +}); + +test("POST /api/jobs/unknown/run-now -> 404", async () => { + const res = await runNowRoute.POST( + new Request("http://localhost/api/jobs/nope/run-now", { method: "POST" }), + { params: params("nope") } + ); + assert.equal(res.status, 404); +}); + +test("LOCAL_ONLY guard: /api/jobs and children are loopback-only", () => { + assert.equal(isLocalOnlyPath("/api/jobs"), true, "bare /api/jobs"); + assert.equal(isLocalOnlyPath("/api/jobs/"), true, "/api/jobs/"); + assert.equal(isLocalOnlyPath("/api/jobs/budget_reset/runs"), true, "runs sub-path"); + assert.equal(isLocalOnlyPath("/api/jobs/warmup/enable"), true, "enable sub-path"); + assert.equal(isLocalOnlyPath("/api/jobs/warmup/disable"), true, "disable sub-path"); + assert.equal(isLocalOnlyPath("/api/jobs/warmup/run-now"), true, "run-now sub-path"); +}); + +test("error responses do not leak stack traces", async () => { + // 404 path - assert the error message is sanitized (no "at /" frame). + const res = await runsRoute.GET(new Request("http://localhost/api/jobs/nope/runs"), { + params: params("nope"), + }); + const body = (await json(res)) as { error: { message: string } }; + assert.ok(!body.error.message.includes("at /"), "error message must not leak stack path"); +}); diff --git a/tests/unit/db/jobRegistryDb.test.ts b/tests/unit/db/jobRegistryDb.test.ts new file mode 100644 index 0000000000..08aa22a9c7 --- /dev/null +++ b/tests/unit/db/jobRegistryDb.test.ts @@ -0,0 +1,290 @@ +/** + * Tests for jobRegistry persistence (migration 139 + jobRegistryDb.ts). + * + * Verifies: + * - 3 built-in jobs seeded by the migration + * - upsertJob insert + column-level update (preserves enabled + created_at) + * - recordRun writes ISO-8601 timestamps + * - pruneRuns count dimension (150 rows -> keep 100) + * - pruneRuns time dimension (rows older than 30 days deleted) + * - pruneRuns dual dimension (150 rows, 50 older than 30 days -> keep 100) + * - cleanupOrphanedRuns fixes stale 'running' records past the timeout + * + * Runs against an isolated temp DATA_DIR so the real ~/.omniroute DB is never touched. + */ + +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-jr-")); +process.env.DATA_DIR = TEST_DATA_DIR; +process.env.DISABLE_SQLITE_AUTO_BACKUP = "true"; + +const core = await import("../../../src/lib/db/core.ts"); +const db = await import("../../../src/lib/db/jobRegistryDb.ts"); + +function resetDb() { + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); + fs.mkdirSync(TEST_DATA_DIR, { recursive: true }); +} + +test.beforeEach(() => { + resetDb(); +}); + +test.after(() => { + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); +}); + +test("seed: migration registers 3 built-in jobs", () => { + const jobs = db.getAllJobs(); + assert.equal(jobs.length, 3); + const ids = jobs.map((j) => j.id).sort(); + assert.deepEqual(ids, ["budget_reset", "token_health_check", "warmup"]); +}); + +test("seed: warmup is cron type with env gate + envDefault=false", () => { + const warmup = db.getAllJobs().find((j) => j.id === "warmup"); + assert.ok(warmup); + assert.equal(warmup.type, "cron"); + assert.equal(warmup.cron, "0 7 * * *"); + assert.equal(warmup.envFlag, "OMNIROUTE_WARMUP_ENABLED"); + // Seeded disabled: the warmup handler is not registered by this change, and + // startAll() would otherwise warn about the missing handler on every boot. + assert.equal(warmup.enabled, false); + assert.equal(warmup.config.envDefault, false); + assert.equal(warmup.config.timezone, "America/Los_Angeles"); +}); + +test("seed: budget_reset is interval type with no env gate", () => { + const budget = db.getAllJobs().find((j) => j.id === "budget_reset"); + assert.ok(budget); + assert.equal(budget.type, "interval"); + assert.equal(budget.intervalMs, 600000); + assert.equal(budget.envFlag, null); +}); + +test("upsertJob: insert a new job then update scheduling fields", () => { + const created = new Date().toISOString(); + db.upsertJob({ + id: "custom", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: { foo: "bar" }, + createdAt: created, + updatedAt: created, + }); + let job = db.getJob("custom"); + assert.ok(job); + assert.equal(job.type, "interval"); + assert.equal(job.intervalMs, 1000); + assert.deepEqual(job.config, { foo: "bar" }); + + // Update interval + config; created_at must not change. + db.upsertJob({ + ...job!, + intervalMs: 2000, + config: { foo: "baz" }, + }); + job = db.getJob("custom"); + assert.equal(job!.intervalMs, 2000); + assert.deepEqual(job!.config, { foo: "baz" }); + assert.equal(job!.createdAt, created); +}); + +test("upsertJob: does NOT overwrite enabled (user toggle preserved)", () => { + const ts = new Date().toISOString(); + db.upsertJob({ + id: "toggle", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: ts, + updatedAt: ts, + }); + db.updateJobEnabled("toggle", false); + // Re-register with enabled=true - must not flip the user's disabled state. + db.upsertJob({ + id: "toggle", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: ts, + updatedAt: ts, + }); + assert.equal(db.getJob("toggle")!.enabled, false); +}); + +test("recordRun: writes a completed run with ISO timestamps", () => { + db.upsertJob({ + id: "j", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + db.recordRun("j", "success", { + startedAt: "2026-01-01T00:00:00.000Z", + durationMs: 123, + recordsAffected: 7, + }); + const runs = db.getRuns("j"); + assert.equal(runs.length, 1); + assert.equal(runs[0].status, "success"); + assert.equal(runs[0].recordsAffected, 7); + assert.equal(runs[0].durationMs, 123); + assert.equal(runs[0].startedAt, "2026-01-01T00:00:00.000Z"); + assert.ok(runs[0].finishedAt, "finishedAt must be set for a completed run"); +}); + +test("recordRun: running status leaves finishedAt NULL", () => { + db.upsertJob({ + id: "j", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + db.recordRun("j", "running", { startedAt: "2026-01-01T00:00:00.000Z" }); + const runs = db.getRuns("j"); + assert.equal(runs[0].status, "running"); + assert.equal(runs[0].finishedAt, null); +}); + +test("pruneRuns: count dimension keeps the most recent 100 of 150", () => { + db.upsertJob({ + id: "j", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + // Insert 150 rows, oldest first, 1s apart. + for (let i = 0; i < 150; i++) { + const t = new Date(Date.UTC(2026, 0, 1, 0, 0, i)).toISOString(); + db.recordRun("j", "success", { startedAt: t, durationMs: 1 }); + } + db.pruneRuns("j", 100, 30); + const runs = db.getRuns("j", 200); + assert.equal(runs.length, 100); + // The oldest surviving row is the 50th (i=50); rows i=0..49 are pruned. + const oldestSec = new Date(runs[runs.length - 1].startedAt).getUTCSeconds(); + assert.equal(oldestSec, 50); +}); + +test("pruneRuns: time dimension deletes rows older than 30 days (once outside recent-100)", () => { + // Dual-dimension semantics: a row is deleted only when it is BOTH outside the + // recent-100 window AND older than maxDays. So we insert 100 recent rows to push + // the 40-day-old row out of the recent-100, where age then prunes it. + db.upsertJob({ + id: "j", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + const now = Date.now(); + db.recordRun("j", "success", { + startedAt: new Date(now - 40 * 86_400_000).toISOString(), + durationMs: 1, + }); + // 100 genuinely-recent rows (hourly, within the last ~4 days) fill the recent-100 window. + for (let i = 0; i < 100; i++) { + db.recordRun("j", "success", { + startedAt: new Date(now - i * 3_600_000).toISOString(), + durationMs: 1, + }); + } + db.pruneRuns("j", 100, 30); + const runs = db.getRuns("j", 200); + // The 40-day-old row is outside recent-100 and >30 days -> deleted. + // Of the 100 recent rows, those aged 31..99 days are also pruned (outside recent-100? no - + // they're within recent-100 by rank). Only the single 40-day outlier is deleted. + const hasOld = runs.some((r) => new Date(r.startedAt).getTime() < now - 30 * 86_400_000); + assert.equal(hasOld, false, "no run older than 30 days should survive"); + assert.equal(runs.length, 100); +}); + +test("pruneRuns: dual dimension - old rows inside the recent-100 are kept", () => { + db.upsertJob({ + id: "j", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + const now = Date.now(); + // 150 rows total: the 50 oldest are >30 days ago, the newest 100 are recent. + for (let i = 0; i < 150; i++) { + // Row i=0 oldest. Newest 100 (i=50..149) are recent; oldest 50 (i=0..49) are 40 days old. + const ageDays = i < 50 ? 40 : 1; + const t = new Date(now - ageDays * 86_400_000 - (149 - i) * 1000).toISOString(); + db.recordRun("j", "success", { startedAt: t, durationMs: 1 }); + } + db.pruneRuns("j", 100, 30); + const runs = db.getRuns("j", 200); + // The 50 rows older than 30 days are all outside the recent-100 -> deleted. + // The 100 recent rows survive. + assert.equal(runs.length, 100); +}); + +test("cleanupOrphanedRuns: fixes only running rows past the timeout", () => { + db.upsertJob({ + id: "j", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + }); + const now = Date.now(); + // Orphaned: running, started 10 minutes ago. + db.recordRun("j", "running", { startedAt: new Date(now - 10 * 60_000).toISOString() }); + // Fresh: running, started now - must NOT be touched. + db.recordRun("j", "running", { startedAt: new Date(now).toISOString() }); + + db.cleanupOrphanedRuns(5); + const runs = db.getRuns("j", 200); + const orphaned = runs.find((r) => new Date(r.startedAt).getTime() <= now - 10 * 60_000); + const fresh = runs.find((r) => new Date(r.startedAt).getTime() > now - 60_000); + assert.equal(orphaned.status, "failure"); + assert.equal(orphaned.errorMessage, "orphaned: exceeded timeout"); + assert.equal(fresh.status, "running"); +}); diff --git a/tests/unit/lib/jobRegistry/boot-wiring.test.ts b/tests/unit/lib/jobRegistry/boot-wiring.test.ts new file mode 100644 index 0000000000..4aa09dada4 --- /dev/null +++ b/tests/unit/lib/jobRegistry/boot-wiring.test.ts @@ -0,0 +1,61 @@ +/** + * Boot wiring for the job registry. + * + * startAll() runs each interval job's first tick synchronously (see startInterval + * in registry.ts), so starting a job before initializeCloudSync() has been awaited + * runs its handler against a half-initialised app. The previous wiring avoided that + * two different ways: the budget reset was started after the init call, and the + * token health check's first sweep sat behind a 10s timer. + * + * Neither entry point can be driven from a unit test - both are process-boot + * functions that reach the real cloud sync - so these read the wiring, the same + * way tests/unit/model-sync-scheduler.test.ts already asserts on this file. + */ + +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import path from "node:path"; + +function readSource(relativePath: string): string { + return fs.readFileSync(path.join(process.cwd(), relativePath), "utf8"); +} + +test("initCloudSync: cloud sync is initialised before the registry starts any job", () => { + const source = readSource("src/lib/initCloudSync.ts"); + const initCall = source.indexOf("await initializeCloudSync("); + const startAllCall = source.indexOf("registry.startAll("); + + assert.notEqual(initCall, -1, "initializeCloudSync() call not found"); + assert.notEqual(startAllCall, -1, "registry.startAll() call not found"); + assert.ok( + initCall < startAllCall, + "startAll() fires each interval job's first tick immediately, so it must come after initializeCloudSync()" + ); +}); + +test("server-init: cloud sync is initialised before the registry starts any job", () => { + const source = readSource("src/server-init.ts"); + const initCall = source.indexOf("await initializeCloudSync("); + const startAllCall = source.indexOf("jobRegistry.startAll("); + + assert.notEqual(initCall, -1, "initializeCloudSync() call not found"); + assert.notEqual(startAllCall, -1, "jobRegistry.startAll() call not found"); + assert.ok(initCall < startAllCall, "startAll() must come after initializeCloudSync()"); +}); + +test("both boot paths register the same jobs", () => { + for (const relativePath of ["src/lib/initCloudSync.ts", "src/server-init.ts"]) { + const source = readSource(relativePath); + assert.match( + source, + /registerBudgetResetJob\s*\(/, + `${relativePath} must register the budget reset job` + ); + assert.match( + source, + /registerTokenHealthCheck\s*\(/, + `${relativePath} must register the token health check job` + ); + } +}); diff --git a/tests/unit/lib/jobRegistry/registry.test.ts b/tests/unit/lib/jobRegistry/registry.test.ts new file mode 100644 index 0000000000..64ce9d2050 --- /dev/null +++ b/tests/unit/lib/jobRegistry/registry.test.ts @@ -0,0 +1,594 @@ +/** JobRegistry runtime tests . Uses real timers + isolated temp DB. */ + +// Access private internals (avoids `as any`). +type TestRegistry = JobRegistry & { + timers: Map; + cronFailCount: Map; +}; +function regInternals(reg: JobRegistry): TestRegistry { + return reg as TestRegistry; +} + +import test from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import type { JobDefinition } from "@/lib/jobRegistry/core.ts"; + +const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-jr-rt-")); +process.env.DATA_DIR = TEST_DATA_DIR; +process.env.DISABLE_SQLITE_AUTO_BACKUP = "true"; + +const core = await import("@/lib/db/core.ts"); +const index = await import("@/lib/jobRegistry/index.ts"); +const { getJobRegistry, __resetJobRegistry } = index; + +let nowIso: string; +function def( + id: string, + handler: () => Promise<{ success: boolean; recordsAffected?: number; error?: string }>, + over: Record = {} +): JobDefinition { + return { + id, + type: "interval", + cron: null, + intervalMs: 100000, + enabled: true, + envFlag: null, + config: {}, + createdAt: nowIso, + updatedAt: nowIso, + handler, + ...over, + } as JobDefinition; +} + +function resetAll() { + // Stop any live timers on the current singleton before wiping the DB, otherwise + // an in-flight safeRun from the previous test writes to a closed/wiped DB (FK error). + try { + getJobRegistry().stopAll(); + } catch { + // no singleton yet on first run + } + __resetJobRegistry(); + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); + fs.mkdirSync(TEST_DATA_DIR, { recursive: true }); +} + +test.beforeEach(() => { + nowIso = new Date().toISOString(); + resetAll(); +}); + +test.after(() => { + core.resetDbInstance(); + fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true }); +}); + +test("register + start (interval) fires handler immediately", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def("j", async () => { + calls++; + return { success: true, recordsAffected: 0 }; + }) + ); + reg.start("j"); + await new Promise((r) => setTimeout(r, 30)); + assert.equal(calls, 1, "interval fires once on start"); + reg.stop("j"); +}); + +test("re-entrancy: second tick skipped while handler is running", async () => { + const reg = getJobRegistry(); + let calls = 0; + let release: (v: void) => void = () => {}; + const gate = new Promise((res) => (release = res)); + reg.register( + def( + "j", + async () => { + calls++; + await gate; + return { success: true }; + }, + { intervalMs: 30 } + ) + ); + reg.start("j"); // fires immediately, blocks on gate + await new Promise((r) => setTimeout(r, 10)); + assert.equal(calls, 1); + // ~100ms passes; setInterval would tick again but must be skipped (running guard). + await new Promise((r) => setTimeout(r, 120)); + assert.equal(calls, 1, "must not re-enter while handler runs"); + release(); + await new Promise((r) => setTimeout(r, 30)); + reg.stop("j"); +}); + +test("cron nextTick: handler fires via recursive setTimeout", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def( + "j", + async () => { + calls++; + return { success: true }; + }, + { + type: "cron", + cron: "* * * * * *", // every second (6-field) for fast deterministic test + intervalMs: null, + config: { timezone: "UTC" }, + } + ) + ); + reg.start("j"); + await new Promise((r) => setTimeout(r, 1400)); + reg.stop("j"); + assert.ok(calls >= 1, `cron handler should fire at least once, got ${calls}`); +}); + +test("runNow: manual trigger starts a disabled job's run", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def("j", async () => { + calls++; + return { success: true }; + }) + ); + const res = await reg.runNow("j"); + assert.deepEqual(res, { started: true }); + await new Promise((r) => setTimeout(r, 30)); + assert.equal(calls, 1); +}); + +test("runNow: disabled job returns reason=disabled", async () => { + const reg = getJobRegistry(); + reg.register(def("j", async () => ({ success: true }), { enabled: false })); + const res = await reg.runNow("j"); + assert.deepEqual(res, { started: false, reason: "disabled" }); +}); + +test("runNow: unknown job returns reason=not_found", async () => { + const reg = getJobRegistry(); + const res = await reg.runNow("nope"); + assert.deepEqual(res, { started: false, reason: "not_found" }); +}); + +test("runNow: queue depth=1, coalesces concurrent triggers then re-runs", async () => { + const reg = getJobRegistry(); + let calls = 0; + let release: (v: void) => void = () => {}; + const gate = new Promise((res) => (release = res)); + reg.register( + def("j", async () => { + calls++; + await gate; + return { success: true }; + }) + ); + const first = await reg.runNow("j"); // starts, blocks on gate + assert.deepEqual(first, { started: true }); + // While running, a second runNow should queue (depth=1) and resolve after. + const queued = reg.runNow("j"); + // The queued promise must not pile up: only one extra fire after release. + release(); + const queuedRes = await queued; + await new Promise((r) => setTimeout(r, 30)); + assert.equal(calls, 2, "initial + one queued re-run = 2 fires, no pile-up"); + assert.deepEqual(queuedRes, { started: true }); +}); + +test("runNow: env gate blocks when env explicitly disabled", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def("j", async () => ({ success: true }), { envFlag: "OMNIROUTE_TEST_JOB_ENABLED" }) + ); + process.env.OMNIROUTE_TEST_JOB_ENABLED = "0"; + const res = await reg.runNow("j"); + delete process.env.OMNIROUTE_TEST_JOB_ENABLED; + assert.deepEqual(res, { started: false, reason: "env_disabled" }); + assert.equal(calls, 0); +}); + +test("setEnabled(false) stops timer; handler no longer fires", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def( + "j", + async () => { + calls++; + return { success: true }; + }, + { intervalMs: 40 } + ) + ); + reg.start("j"); + await new Promise((r) => setTimeout(r, 25)); + reg.setEnabled("j", false); + await new Promise((r) => setTimeout(r, 120)); + reg.stop("j"); + assert.equal(calls, 1, "only the immediate fire before disable"); +}); + +test("setEnabled(true) restarts a stopped job", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def( + "j", + async () => { + calls++; + return { success: true }; + }, + { intervalMs: 40 } + ) + ); + reg.setEnabled("j", true); + await new Promise((r) => setTimeout(r, 130)); + reg.stop("j"); + assert.ok(calls >= 2, `expected repeated fires, got ${calls}`); +}); + +test("envFlag gate generic: unset env fires (defaultWhenUnset=true)", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def( + "j", + async () => { + calls++; + return { success: true }; + }, + { + type: "cron", + cron: "* * * * * *", + intervalMs: null, + envFlag: "OMNIROUTE_GENERIC_JOB_ENABLED", // unset -> default true -> fires + config: { timezone: "UTC" }, + } + ) + ); + reg.start("j"); + await new Promise((r) => setTimeout(r, 1400)); + reg.stop("j"); + assert.ok(calls >= 1, "unset env with default=true should fire"); +}); + +test("envFlag gate warmup: unset env does NOT fire (envDefault=false)", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def("j", async () => ({ success: true }), { + type: "cron", + cron: "* * * * * *", + intervalMs: null, + envFlag: "OMNIROUTE_WARMUP_ENABLED", + config: { timezone: "UTC", envDefault: false }, + }) + ); + reg.start("j"); + await new Promise((r) => setTimeout(r, 1400)); + reg.stop("j"); + assert.equal(calls, 0, "unset env with envDefault=false must not fire"); +}); + +test("startAll: registers + starts all enabled jobs; throws if no handlers", async () => { + const reg = getJobRegistry(); + await assert.rejects(() => reg.startAll(), /No handlers registered/); + + let a = 0; + let b = 0; + reg.register( + def("a", async () => { + a++; + return { success: true, recordsAffected: 0 }; + }) + ); + reg.register( + def("b", async () => { + b++; + return { success: true, recordsAffected: 0 }; + }) + ); + await reg.startAll(); + await new Promise((r) => setTimeout(r, 30)); + assert.ok(a >= 1, "job a started"); + assert.ok(b >= 1, "job b started"); + reg.stopAll(); +}); + +test("startAll isolation: a missing handler skips that job, starts the rest", async () => { + const reg = getJobRegistry(); + let a = 0; + reg.register( + def("a", async () => { + a++; + return { success: true, recordsAffected: 0 }; + }) + ); + // Seed a job in DB with no handler registered in the registry. + core.getDbInstance(); + const { upsertJob } = await import("@/lib/db/jobRegistryDb.ts"); + upsertJob({ + id: "orphan", + type: "interval", + cron: null, + intervalMs: 1000, + enabled: true, + envFlag: null, + config: {}, + createdAt: nowIso, + updatedAt: nowIso, + }); + await reg.startAll(); + await new Promise((r) => setTimeout(r, 30)); + assert.ok(a >= 1, "registered job a still starts despite orphan in DB"); + reg.stopAll(); +}); + +test("startAll isolation: one job whose start() throws does not abort the rest", async () => { + const reg = getJobRegistry(); + let a = 0; + let c = 0; + for (const [id, counter] of [ + ["a", () => a++], + ["boom", () => 0], + ["c", () => c++], + ] as Array<[string, () => number]>) { + reg.register( + def(id, async () => { + counter(); + return { success: true, recordsAffected: 0 }; + }) + ); + } + + // start() reads the DB (getJob), so it can throw for one job while the others + // are fine. Nothing in the registry catches per-job, so this asserts the loop + // in startAll() isolates it. + const realStart = reg.start.bind(reg); + reg.start = (jobId: string): void => { + if (jobId === "boom") throw new Error("start blew up for boom"); + realStart(jobId); + }; + + try { + await reg.startAll(); + } finally { + reg.start = realStart; + } + + await new Promise((r) => setTimeout(r, 30)); + assert.ok(a >= 1, "job registered before the throwing one still started"); + assert.ok(c >= 1, "job registered after the throwing one still started"); + reg.stopAll(); +}); + +test("error isolation: handler throw records failure and survives", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def("j", async () => { + calls++; + if (calls === 1) throw new Error("boom at /some/file.ts"); + return { success: true }; + }) + ); + await reg.runNow("j"); + await new Promise((r) => setTimeout(r, 30)); + const runs = reg.getRuns("j"); + assert.equal(runs.length, 1); + assert.equal(runs[0].status, "failure"); + // Error message must be sanitized - no stack path leaked. + assert.ok(!runs[0].errorMessage?.includes("at /"), "error must be sanitized"); +}); + +test("getRuns returns recent records newest-first", async () => { + const reg = getJobRegistry(); + reg.register(def("j", async () => ({ success: true, recordsAffected: 0 }))); + await reg.runNow("j"); + await reg.runNow("j"); + await new Promise((r) => setTimeout(r, 30)); + const runs = reg.getRuns("j"); + assert.equal(runs.length, 2); + assert.ok( + new Date(runs[0].startedAt).getTime() >= new Date(runs[1].startedAt).getTime(), + "newest first" + ); +}); + +test("listJobs returns registered jobs with handlers", () => { + const reg = getJobRegistry(); + reg.register(def("j", async () => ({ success: true }))); + const jobs = reg.listJobs(); + const j = jobs.find((x) => x.id === "j"); + assert.ok(j, "registered job j should appear in listJobs"); + assert.equal(typeof j!.handler, "function"); +}); + +test("cron: invalid expression stops re-scheduling after MAX_CRON_PARSE_FAILURES", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def( + "badcron", + async () => { + calls++; + return { success: true }; + }, + { + type: "cron", + cron: "not-a-valid-cron-expression", + intervalMs: null, + config: { timezone: "UTC" }, + } + ) + ); + reg.start("badcron"); + const failCount = regInternals(reg).cronFailCount.get("badcron") ?? 0; + assert.ok(failCount >= 1, `parse should fail, got failCount=${failCount}`); + assert.equal(calls, 0, "invalid cron handler should never fire"); + reg.stop("badcron"); +}); + +test("runNow: queued trigger waits for current run then re-fires", async () => { + const reg = getJobRegistry(); + let calls = 0; + let release: (v: void) => void = () => {}; + const gate = new Promise((res) => (release = res)); + reg.register( + def("slow", async () => { + calls++; + await gate; + return { success: true }; + }) + ); + const first = await reg.runNow("slow"); + assert.deepEqual(first, { started: true }); + const queued = reg.runNow("slow"); + release(); + const queuedRes = await queued; + await new Promise((r) => setTimeout(r, 30)); + assert.equal(calls, 2, "initial + one queued re-run = 2 fires"); + assert.deepEqual(queuedRes, { started: true }); +}); + +test("dispose clears timers and cron failure counts", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register( + def( + "j", + async () => { + calls++; + return { success: true }; + }, + { intervalMs: 30 } + ) + ); + reg.start("j"); + await new Promise((r) => setTimeout(r, 10)); + assert.ok(regInternals(reg).timers.has("j")); + reg.dispose(); + assert.equal(regInternals(reg).timers.size, 0); + assert.equal(regInternals(reg).cronFailCount.size, 0); + const callsBefore = calls; + await new Promise((r) => setTimeout(r, 100)); + assert.equal(calls, callsBefore, "disposed timers must not fire"); +}); + +test("cronGetter: re-reads cron on each fire", async () => { + const reg = getJobRegistry(); + let calls = 0; + let cronExpr = "*/2 * * * * *"; + reg.register({ + id: "dynamic", + type: "cron", + cron: cronExpr, + intervalMs: null, + enabled: true, + envFlag: null, + config: { timezone: "UTC" }, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => { + calls++; + return { success: true }; + }, + cronGetter: () => cronExpr, + }); + reg.start("dynamic"); + await new Promise((r) => setTimeout(r, 2500)); + const callsAfterInitial = calls; + assert.ok(callsAfterInitial >= 1, `should fire with initial cron, got ${callsAfterInitial}`); + cronExpr = "0 0 1 1 * 2099"; + await new Promise((r) => setTimeout(r, 3000)); + reg.stop("dynamic"); + const extraFires = calls - callsAfterInitial; + assert.ok(extraFires <= 1, `cronGetter change should stop fires, got ${extraFires} extra`); +}); + +test("cronFailCount: resets after successful parse", async () => { + const reg = getJobRegistry(); + let calls = 0; + reg.register({ + id: "recover", + type: "cron", + cron: "not-a-cron", + intervalMs: null, + enabled: true, + envFlag: null, + config: { timezone: "UTC" }, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => { + calls++; + return { success: true }; + }, + }); + reg.start("recover"); + const failCountAfterFail = regInternals(reg).cronFailCount.get("recover") ?? 0; + assert.ok(failCountAfterFail >= 1, `parse should fail, got failCount=${failCountAfterFail}`); + reg.register({ + id: "recover", + type: "cron", + cron: "* * * * * *", + intervalMs: null, + enabled: true, + envFlag: null, + config: { timezone: "UTC" }, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => { + calls++; + return { success: true }; + }, + cronGetter: () => "* * * * * *", + }); + reg.stop("recover"); + reg.start("recover"); + const failCountAfterRecovery = regInternals(reg).cronFailCount.get("recover") ?? 0; + assert.equal(failCountAfterRecovery, 0, "successful parse should reset failure count"); + reg.stop("recover"); +}); + +test("cronGetter: throws falls back to static cron", async () => { + const reg = getJobRegistry(); + let calls = 0; + let shouldThrow = true; + reg.register({ + id: "throwing", + type: "cron", + cron: "* * * * * *", + intervalMs: null, + enabled: true, + envFlag: null, + config: { timezone: "UTC" }, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + handler: async () => { + calls++; + return { success: true }; + }, + cronGetter: () => { + if (shouldThrow) throw new Error("transient failure"); + return "* * * * * *"; + }, + }); + reg.start("throwing"); + await new Promise((r) => setTimeout(r, 1100)); + assert.ok(calls >= 1, "job should fire with static cron fallback when cronGetter throws"); + reg.stop("throwing"); +}); diff --git a/tests/unit/lib/jobRegistry/timeUtils.test.ts b/tests/unit/lib/jobRegistry/timeUtils.test.ts new file mode 100644 index 0000000000..2ed967e870 --- /dev/null +++ b/tests/unit/lib/jobRegistry/timeUtils.test.ts @@ -0,0 +1,48 @@ +/** + * Tests for jobRegistry timeUtils (convertToTimeZone). + * + * Verifies wall-clock conversion across IANA timezones, date-boundary crossing, + * and DST handling (the conversion itself is DST-naive; cron-parser owns DST on the + * scheduling side - we only require the wall clock to read correctly). + */ + +import test from "node:test"; +import assert from "node:assert/strict"; +import { convertToTimeZone } from "@/lib/jobRegistry/timeUtils.ts"; + +test("convertToTimeZone: UTC noon -> Pacific standard time (UTC-8)", () => { + // 2026-01-15T12:00:00Z -> PST is UTC-8 -> 04:00 same day + const utc = new Date("2026-01-15T12:00:00Z"); + const pt = convertToTimeZone(utc, "America/Los_Angeles"); + assert.equal(pt.getHours(), 4); + assert.equal(pt.getDate(), 15); +}); + +test("convertToTimeZone: UTC noon -> Pacific daylight time (UTC-7)", () => { + // 2026-07-15T12:00:00Z -> PDT is UTC-7 -> 05:00 same day + const utc = new Date("2026-07-15T12:00:00Z"); + const pt = convertToTimeZone(utc, "America/Los_Angeles"); + assert.equal(pt.getHours(), 5); +}); + +test("convertToTimeZone: crosses date boundary backward (UTC -> Asia/Tokyo, UTC+9)", () => { + // 2026-01-15T20:00:00Z -> JST +9 -> 05:00 NEXT day + const utc = new Date("2026-01-15T20:00:00Z"); + const jst = convertToTimeZone(utc, "Asia/Tokyo"); + assert.equal(jst.getHours(), 5); + assert.equal(jst.getDate(), 16); +}); + +test("convertToTimeZone: UTC identity zone returns same wall clock", () => { + const utc = new Date("2026-03-01T13:45:00Z"); + const same = convertToTimeZone(utc, "UTC"); + assert.equal(same.getHours(), 13); + assert.equal(same.getMinutes(), 45); +}); + +test("convertToTimeZone: Europe/London winter (UTC+0) vs summer (UTC+1 BST)", () => { + const winter = convertToTimeZone(new Date("2026-01-15T12:00:00Z"), "Europe/London"); + assert.equal(winter.getHours(), 12); + const summer = convertToTimeZone(new Date("2026-07-15T12:00:00Z"), "Europe/London"); + assert.equal(summer.getHours(), 13); +}); diff --git a/tests/unit/token-health-check-sweep.test.ts b/tests/unit/token-health-check-sweep.test.ts index 87e2d548a4..afce89cde8 100644 --- a/tests/unit/token-health-check-sweep.test.ts +++ b/tests/unit/token-health-check-sweep.test.ts @@ -104,9 +104,13 @@ test("sweep() skips re-entrant calls while a previous sweep is still in flight", assert.equal(isSweeping(), true, "first sweep should mark sweeping=true immediately"); const start = Date.now(); - await sweep(); // second, concurrent call — must be skipped by the guard + const skippedCount = await sweep(); // second, concurrent call — must be skipped by the guard const elapsedMs = Date.now() - start; + // The job registry records this as records_affected, so a skipped sweep has + // to report nothing rather than inheriting the in-flight sweep's count. + assert.equal(skippedCount, 0, "a skipped sweep reports 0 connections swept"); + // With 21 connections and a 300ms inter-batch stagger, a REAL second // sweep would take >= 300ms to clear the single batch gap. The guard must // make this call return near- @@ -130,8 +134,9 @@ test("sweep() resets the sweeping flag after a normal completion, allowing the n process.env.HEALTHCHECK_STAGGER_MS = "0"; await createNoOpOauthConnections(2, "sequential"); - await sweep(); + const sweptCount = await sweep(); assert.equal(isSweeping(), false, "flag resets after first sweep completes"); + assert.equal(sweptCount, 2, "sweep reports how many connections it swept"); // A second, fully sequential call must not be skipped by leftover state — // it should run and complete normally rather than hang or short-circuit. @@ -144,7 +149,8 @@ test("sweep() resets the sweeping flag even when there are no connections to pro process.env.HEALTHCHECK_STAGGER_MS = "0"; // No connections created — sweep() takes the early `connections.length === 0` return. - await sweep(); + const emptyCount = await sweep(); assert.equal(isSweeping(), false, "sweeping flag must be false after an empty sweep"); + assert.equal(emptyCount, 0, "an empty sweep reports 0 connections swept"); });