mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-11 09:42:15 +03:00
* feat(db): add a job registry for scheduled background work Background jobs each ship their own timer today, so there is no list of what is scheduled, no history of what ran, and no way to pause one without an environment variable and a restart. The registry gives them one home: a jobs table holding the schedule, a job_runs table holding the outcomes, and a loopback-only API to inspect and control both. Cron jobs read their expression through an optional cronGetter rather than the stored column, so an operator changing OMNIROUTE_WARMUP_CRON does not need the row rewritten. register() is an idempotent upsert that refreshes the schedule but never overwrites `enabled` or `created_at`, which is what lets a job be re-registered on every boot without discarding the operator's toggle. Run history is pruned per job rather than globally, and safeRun records a failure for a handler that throws as well as one that returns success:false, so a crashing job leaves a trail instead of a gap. The API is under /api/jobs and gated to loopback in the route guard. It can trigger a run and flip a job off, which is runtime administration and does not belong on a remotely reachable surface. Signed-off-by: Minxi Hou <houminxi@gmail.com> * feat(jobs): move the budget reset and token health check onto the registry Both jobs owned their own timer and started themselves as an import side effect, so nothing could report whether they were running, when they last ran, or why a run failed. They now register with the job registry and are started from it, which also means their schedule and run history are visible through /api/jobs. startAll() runs each interval job's first tick synchronously, so both entry points start the registry only after initializeCloudSync() has been awaited. The old wiring reached that ordering two different ways: the budget reset was started after the init call, and the health check's first sweep sat behind a 10s timer. Replacing both with one startAll() would otherwise have moved the two handlers in front of the initialisation they run against. Both entry points also register the same pair of jobs. Registering one and not the other is how a background job goes missing without anything failing. sweep() now returns how many connections it swept, so the health check can record a real records_affected the way the budget reset does. The migration documents that column as a per-job count, and hardcoding zero would have left one of the two jobs reporting a number the schema promises but the code never produces. A skipped or empty sweep reports zero. Every existing caller ignores the return value. The token health check keeps its own disable semantics: the handler still calls isHealthCheckDisabled() before sweeping, so OMNIROUTE_DISABLE_TOKEN_HEALTHCHECK, the production-build phase and the automated-test guard behave as before. Its registry adapter lives in src/lib/jobs/ next to the budget reset rather than in tokenHealthCheck.ts, which is already above its frozen size ceiling on the base branch and should not grow further. The adapter lets a failing sweep throw rather than reporting it itself, matching the budget reset: safeRun records a thrown error as a failure run with its message. The warmup job is seeded disabled. Its handler arrives with the warmup scheduler, and startAll() filters on enabled before it looks for a handler, so seeding it enabled here would warn about the missing handler on every boot. * fix: allowlist cron-parser dep and document OMNIROUTE_RUNNOW_TIMEOUT_MS env var Co-authored-by: diegosouzapw <diegosouza.pw@gmail.com> --------- Signed-off-by: Minxi Hou <houminxi@gmail.com> Co-authored-by: Minxi Hou <houminxi@gmail.com> Co-authored-by: diegosouzapw <diegosouzapw@users.noreply.github.com>
184 lines
7.2 KiB
TypeScript
184 lines
7.2 KiB
TypeScript
// Server startup script
|
|
import initializeCloudSync from "./shared/services/initializeCloudSync";
|
|
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 { 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";
|
|
import { applyRuntimeSettings } from "./lib/config/runtimeSettings";
|
|
import { setSystemPromptConfig } from "@omniroute/open-sse/services/systemPrompt.ts";
|
|
import { hydrateThinkingBudgetConfig } from "@omniroute/open-sse/services/thinkingBudget.ts";
|
|
import { startRuntimeConfigHotReload } from "./lib/config/hotReload";
|
|
import { startSpendBatchWriter } from "./lib/spend/batchWriter";
|
|
import { registerDefaultGuardrails } from "./lib/guardrails";
|
|
import { ensurePersistentManagementPasswordHash } from "./lib/auth/managementPassword";
|
|
import { skillExecutor } from "./lib/skills/executor";
|
|
import { registerBuiltinSkills } from "./lib/skills/builtins";
|
|
import { createLogger } from "./shared/utils/logger";
|
|
|
|
const startupLog = createLogger("server-init");
|
|
|
|
function getErrorMessage(error: unknown) {
|
|
return error instanceof Error ? error.message : String(error);
|
|
}
|
|
|
|
async function startServer() {
|
|
// Trigger request-log layout migration during startup, before serving requests.
|
|
await import("./lib/usage/migrations");
|
|
|
|
// Console interceptor: capture all console output to log file (must be first)
|
|
initConsoleInterceptor();
|
|
|
|
// FASE-01: Validate required secrets before anything else (fail-fast)
|
|
enforceSecrets();
|
|
enforceWebRuntimeEnv();
|
|
|
|
// Compliance: Initialize audit_log table
|
|
try {
|
|
initAuditLog();
|
|
startupLog.info("Audit log table initialized");
|
|
} catch (err) {
|
|
startupLog.warn({ err }, "Could not initialize audit log");
|
|
}
|
|
|
|
// Compliance: One-time cleanup of expired logs
|
|
try {
|
|
const cleanup = await cleanupExpiredLogs();
|
|
if (
|
|
cleanup.deletedUsage ||
|
|
cleanup.deletedCallLogs ||
|
|
cleanup.deletedProxyLogs ||
|
|
cleanup.deletedRequestDetailLogs ||
|
|
cleanup.deletedAuditLogs ||
|
|
cleanup.deletedMcpAuditLogs
|
|
) {
|
|
startupLog.info({ cleanup }, "Expired log cleanup completed");
|
|
}
|
|
} catch (err) {
|
|
startupLog.warn({ err }, "Log cleanup failed");
|
|
}
|
|
|
|
startupLog.info("Starting server with cloud sync");
|
|
|
|
try {
|
|
let settings = await getSettings();
|
|
const passwordState = await ensurePersistentManagementPasswordHash({
|
|
logger: { log: (message: string) => startupLog.info(message) },
|
|
settings,
|
|
source: "startup",
|
|
});
|
|
settings = passwordState.settings;
|
|
const runtimeChanges = await applyRuntimeSettings(settings, { force: true, source: "startup" });
|
|
if (runtimeChanges.length > 0) {
|
|
startupLog.info(
|
|
{ sections: runtimeChanges.map((entry) => entry.section) },
|
|
"Runtime settings hydrated"
|
|
);
|
|
}
|
|
|
|
// Restore the Global System Prompt into the in-memory config. It lives in the
|
|
// `settings.systemPrompt` key but is NOT covered by applyRuntimeSettings, so without
|
|
// this the toggle/prompt revert to defaults on every restart (#2470).
|
|
if (settings.systemPrompt) {
|
|
setSystemPromptConfig(settings.systemPrompt);
|
|
startupLog.info("Global System Prompt restored from settings");
|
|
}
|
|
|
|
// Restore the proxy-level Thinking-Budget config (#5312). It lives in
|
|
// `settings.thinkingBudget` and is NOT covered by applyRuntimeSettings, so
|
|
// without this the dashboard mode (auto/custom/adaptive) silently reverts to
|
|
// the passthrough default on every restart.
|
|
if (hydrateThinkingBudgetConfig(settings)) {
|
|
startupLog.info("Thinking-Budget config restored from settings");
|
|
}
|
|
|
|
// Initialize cloud sync
|
|
startSpendBatchWriter();
|
|
registerDefaultGuardrails();
|
|
registerBuiltinSkills(skillExecutor);
|
|
startupLog.info("Spend batch writer started");
|
|
startupLog.info("Guardrail registry initialized");
|
|
startupLog.info("Builtin skill handlers registered");
|
|
|
|
// Load active plugins on startup so they survive restarts
|
|
try {
|
|
const { pluginManager } = await import("./lib/plugins/manager");
|
|
await pluginManager.loadAll();
|
|
startupLog.info("Plugin manager loaded active plugins");
|
|
} catch (err) {
|
|
startupLog.warn({ err }, "Plugin manager loadAll failed (non-fatal)");
|
|
}
|
|
|
|
await initializeCloudSync();
|
|
// 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();
|
|
startupLog.info("Server started with cloud sync initialized");
|
|
|
|
// Log server start event to audit log
|
|
logAuditEvent({
|
|
action: "server.start",
|
|
actor: "system",
|
|
target: "server-runtime",
|
|
resourceType: "maintenance",
|
|
status: "success",
|
|
details: { timestamp: new Date().toISOString() },
|
|
});
|
|
} catch (error) {
|
|
startupLog.error({ err: error }, "Error initializing cloud sync");
|
|
process.exit(1);
|
|
}
|
|
|
|
// Pricing sync: opt-in external pricing data (non-blocking, never fatal)
|
|
if (process.env.PRICING_SYNC_ENABLED === "true") {
|
|
try {
|
|
const { initPricingSync } = await import("./lib/pricingSync");
|
|
await initPricingSync();
|
|
} catch (err) {
|
|
startupLog.warn({ error: getErrorMessage(err) }, "Pricing sync could not initialize");
|
|
}
|
|
}
|
|
|
|
// Arena ELO sync: model intelligence from leaderboard data (non-blocking, never fatal).
|
|
// On by default; opt out with Dashboard Feature Flags or ARENA_ELO_SYNC_ENABLED=false.
|
|
try {
|
|
const { initArenaEloSync } = await import("./lib/arenaEloSync");
|
|
await initArenaEloSync();
|
|
} catch (err) {
|
|
startupLog.warn({ error: getErrorMessage(err) }, "Arena ELO sync could not initialize");
|
|
}
|
|
|
|
// Radar daily feed sync: only arms itself when RADAR_ENABLED AND the user
|
|
// opt-in are already on (a flag-off boot stays timer-free — Radar inertia
|
|
// contract). Non-blocking, never fatal.
|
|
try {
|
|
const { initRadarSyncScheduler } = await import("./lib/radar/scheduler");
|
|
initRadarSyncScheduler();
|
|
} catch (err) {
|
|
startupLog.warn({ error: getErrorMessage(err) }, "Radar sync scheduler could not initialize");
|
|
}
|
|
}
|
|
|
|
// Start the server initialization
|
|
startServer().catch((err) => {
|
|
startupLog.error({ err }, "Server initialization failed");
|
|
process.exit(1);
|
|
});
|
|
|
|
// Export for use as module if needed
|
|
export default startServer;
|