From ca586db060dd4525a8bbfdb2e0574e4e4e418461 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sun, 4 Oct 2026 23:20:40 +0200 Subject: [PATCH] fix(metrics): correct CPU and queue alerts with bounded compile statistics --- server/routes/status.routes.ts | 5 +- server/services/compilation-worker-pool.ts | 9 ++- server/services/server-metrics.ts | 25 ++++---- .../routes/status-compile-capacity.test.ts | 37 ++++++++++++ .../services/server-metrics-cpu.test.ts | 60 +++++++++++++++++++ tests/server/worker-pool.test.ts | 53 ++++++++++++++++ 6 files changed, 172 insertions(+), 17 deletions(-) create mode 100644 tests/server/services/server-metrics-cpu.test.ts diff --git a/server/routes/status.routes.ts b/server/routes/status.routes.ts index 228ed64ca..e7c45c9a6 100644 --- a/server/routes/status.routes.ts +++ b/server/routes/status.routes.ts @@ -38,8 +38,8 @@ statusRouter.get("/api/status", (_req, res) => { ? { max: compilerStats.maxWorkers, live: compilerStats.liveWorkers } : { max: 0, live: 0 }; let compileCapacity = { maxConcurrent: workers.max, active: compilerStats.activeWorkers }; + const gatekeeper = getUnifiedGatekeeper().getStats(); if (workers.live === 0) { - const gatekeeper = getUnifiedGatekeeper().getStats(); compileCapacity = { maxConcurrent: gatekeeper.maxConcurrentCompiles, active: gatekeeper.activeCompiles }; } @@ -132,7 +132,8 @@ statusRouter.get("/api/status", (_req, res) => { }, observabilityAlerts: evaluateObservabilityAlerts({ compileMetrics, - compileQueueDepth: sandboxStart.queueLength, + // Both native sandbox starts and REST compiles can wait for compile capacity. + compileQueueDepth: sandboxStart.queueLength + compilerStats.queuedTasks + (gatekeeper.queuedCompiles ?? 0), runnerQueueDepth: poolStats.queuedRequests, runnerCapacity: poolStats.maxRunners, processMetrics, diff --git a/server/services/compilation-worker-pool.ts b/server/services/compilation-worker-pool.ts index 9f59f9560..64253230e 100644 --- a/server/services/compilation-worker-pool.ts +++ b/server/services/compilation-worker-pool.ts @@ -104,7 +104,7 @@ export class CompilationWorkerPool { totalTasks: 0, completedTasks: 0, failedTasks: 0, - compileTimes: [] as number[], + totalCompileTimeMs: 0, }; constructor(numWorkers?: number, options: CompilationWorkerPoolOptions = {}) { @@ -334,7 +334,7 @@ export class CompilationWorkerPool { !payload.result.success && `${payload.result.stderr ?? ""} ${payload.result.errors.map((err) => err.message).join(" ")}`.toLowerCase().includes("timeout"), ); this.stats.completedTasks++; - this.stats.compileTimes.push(compileTimeMs); + this.stats.totalCompileTimeMs += compileTimeMs; this.logger.info( `[Worker ${workerId}] Compiled in ${compileTimeMs}ms`, ); @@ -412,10 +412,9 @@ export class CompilationWorkerPool { * Get pool statistics */ getStats(): PoolStats { - const compileTimes = this.stats.compileTimes; const avgCompileTimeMs = - compileTimes.length > 0 - ? compileTimes.reduce((a, b) => a + b, 0) / compileTimes.length + this.stats.completedTasks > 0 + ? this.stats.totalCompileTimeMs / this.stats.completedTasks : 0; return { diff --git a/server/services/server-metrics.ts b/server/services/server-metrics.ts index 3663b9574..ee05bdc06 100644 --- a/server/services/server-metrics.ts +++ b/server/services/server-metrics.ts @@ -24,6 +24,7 @@ interface CpuUsageSample { } let previousCpuSample: CpuUsageSample | null = null; +let previousCpuPercent = 0; /** * Get current CPU and memory usage of the Node.js process @@ -39,26 +40,30 @@ export function getProcessMetrics(): ProcessMetrics { // CPU usage (requires two samples) const cpuUsage = process.cpuUsage(); - let cpuPercent = 0; + let cpuPercent = previousCpuPercent; - if (previousCpuSample !== null) { + if (previousCpuSample !== null && now > previousCpuSample.timestamp) { const elapsedMs = now - previousCpuSample.timestamp; - const elapsedNs = elapsedMs * 1_000_000; // Convert to nanoseconds + const elapsedUs = elapsedMs * 1_000; // cpuUsage() reports microseconds const userDiff = cpuUsage.user - previousCpuSample.user; const systemDiff = cpuUsage.system - previousCpuSample.system; const totalDiff = userDiff + systemDiff; - // CPU percentage across all cores - cpuPercent = (totalDiff / elapsedNs) * 100; + // Aggregate process CPU: 100% is one fully used core, not host utilization. + cpuPercent = (totalDiff / elapsedUs) * 100; } // Store sample for next calculation - previousCpuSample = { - user: cpuUsage.user, - system: cpuUsage.system, - timestamp: now, - }; + // Equal/backward timestamps have no valid interval; keep the last baseline. + if (previousCpuSample === null || now > previousCpuSample.timestamp) { + previousCpuSample = { + user: cpuUsage.user, + system: cpuUsage.system, + timestamp: now, + }; + previousCpuPercent = cpuPercent; + } return { cpuPercent, diff --git a/tests/server/routes/status-compile-capacity.test.ts b/tests/server/routes/status-compile-capacity.test.ts index 120ec7461..7980531af 100644 --- a/tests/server/routes/status-compile-capacity.test.ts +++ b/tests/server/routes/status-compile-capacity.test.ts @@ -16,6 +16,7 @@ vi.mock("../../../server/services/unified-gatekeeper", () => ({ })); import { registerStatusRoutes } from "../../../server/routes/status.routes"; +import { getSandboxStartSemaphore } from "../../../server/services/sandbox/docker-compile-semaphore"; let server: http.Server | undefined; afterEach(() => new Promise((resolve) => (server ? server.close(() => resolve()) : resolve()))); @@ -30,6 +31,42 @@ async function status(): Promise> { } describe("/api/status compile capacity", () => { + it("retains sandbox-start queue alerts and aliases without inventing REST queue work", async () => { + const queued = vi.spyOn(getSandboxStartSemaphore(), "queueLength", "get").mockReturnValue(2); + Object.assign(compilerStats, { queuedTasks: 0 }); + Object.assign(gatekeeperStats, { queuedCompiles: 0 }); + try { + const body = await status(); + expect(body.observabilityAlerts.map((alert: { code: string }) => alert.code)).toContain("compile_queue_nonempty"); + expect(body.compile.queued).toBe(2); + expect(body.compileSlots.queued).toBe(2); + expect(body.compileWorkerPool.queued).toBe(0); + } finally { + queued.mockRestore(); + } + }); + + it("does not alert on compile queues in an idle snapshot", async () => { + Object.assign(compilerStats, { queuedTasks: 0 }); + Object.assign(gatekeeperStats, { queuedCompiles: 0 }); + const body = await status(); + expect(body.observabilityAlerts.map((alert: { code: string }) => alert.code)).not.toContain("compile_queue_nonempty"); + }); + + it.each(["worker", "gatekeeper"])("alerts on queued REST compiles in the %s path without changing sandbox aliases", async (queue) => { + Object.assign(compilerStats, { maxWorkers: queue === "worker" ? 3 : 0, liveWorkers: queue === "worker" ? 3 : 0, activeWorkers: 0, queuedTasks: queue === "worker" ? 1 : 0 }); + Object.assign(gatekeeperStats, { queuedCompiles: queue === "gatekeeper" ? 1 : 0 }); + try { + const body = await status(); + expect(body.observabilityAlerts.map((alert: { code: string }) => alert.code)).toContain("compile_queue_nonempty"); + expect(body.compile.queued).toBe(0); + expect(body.compileSlots.queued).toBe(0); + } finally { + Object.assign(compilerStats, { queuedTasks: 0 }); + Object.assign(gatekeeperStats, { queuedCompiles: 0 }); + } + }); + it("reports the running compile workers when the worker pool serves compiles", async () => { Object.assign(compilerStats, { maxWorkers: 3, liveWorkers: 3, activeWorkers: 1 }); diff --git a/tests/server/services/server-metrics-cpu.test.ts b/tests/server/services/server-metrics-cpu.test.ts new file mode 100644 index 000000000..20df696af --- /dev/null +++ b/tests/server/services/server-metrics-cpu.test.ts @@ -0,0 +1,60 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +let metrics: typeof import("../../../server/services/server-metrics"); +let now: number; +let cpu: NodeJS.CpuUsage; + +beforeEach(async () => { + vi.resetModules(); + metrics = await import("../../../server/services/server-metrics"); + now = 1_000; + cpu = { user: 0, system: 0 }; + vi.spyOn(Date, "now").mockImplementation(() => now); + vi.spyOn(process, "cpuUsage").mockImplementation(() => ({ ...cpu })); +}); + +afterEach(() => vi.restoreAllMocks()); + +describe("process CPU sampling", () => { + it.each([ + { user: 900_000, system: 50_000, expected: 95 }, + { user: 1_800_000, system: 200_000, expected: 200 }, + ])("converts microseconds to elapsed time: $expected percent", ({ user, system, expected }) => { + expect(metrics.getProcessMetrics().cpuPercent).toBe(0); + now = 2_000; + cpu = { user, system }; + const sampled = metrics.getProcessMetrics(); + expect(sampled.cpuPercent).toBe(expected); + const alerts = metrics.evaluateObservabilityAlerts({ + compileMetrics: metrics.compileMetricsTracker.getMetrics(), + compileQueueDepth: 0, + runnerQueueDepth: 0, + runnerCapacity: 1, + processMetrics: sampled, + }); + expect(alerts.map(({ code }) => code)).toContain("process_cpu_high"); + }); + + it("keeps a finite value and the complete interval for same-timestamp polls", () => { + metrics.getProcessMetrics(); + cpu = { user: 50_000, system: 0 }; + expect(metrics.getProcessMetrics().cpuPercent).toBe(0); + now = 2_000; + cpu = { user: 1_000_000, system: 0 }; + expect(metrics.getProcessMetrics().cpuPercent).toBe(100); + expect(metrics.getProcessMetrics().cpuPercent).toBe(100); + }); + + it("retains the last valid sample through a backwards wall-clock adjustment", () => { + metrics.getProcessMetrics(); + now = 2_000; + cpu = { user: 950_000, system: 0 }; + expect(metrics.getProcessMetrics().cpuPercent).toBe(95); + now = 1_999; + cpu = { user: 1_000_000, system: 0 }; + expect(metrics.getProcessMetrics().cpuPercent).toBe(95); + now = 3_000; + cpu = { user: 1_900_000, system: 0 }; + expect(metrics.getProcessMetrics().cpuPercent).toBe(95); + }); +}); diff --git a/tests/server/worker-pool.test.ts b/tests/server/worker-pool.test.ts index 255ad24d5..9b867a390 100644 --- a/tests/server/worker-pool.test.ts +++ b/tests/server/worker-pool.test.ts @@ -135,6 +135,59 @@ describe("CompilationWorkerPool", () => { }); }); + it("keeps lifetime average statistics without retaining a growing duration history", async () => { + const pool = createPool(1); + let now = 1_000; + const clock = vi.spyOn(Date, "now").mockImplementation(() => now); + const retainedSamples = () => Object.values( + (pool as unknown as { stats: Record }).stats, + ).reduce((count, value) => count + (Array.isArray(value) ? value.length : 0), 0); + try { + for (let index = 0; index < 2_000; index++) { + const pending = pool.compile({ code: "statistics-fixture" }); + now += index % 2 === 0 ? 10 : 30; + succeed(0); + await pending; + } + expect(pool.getStats()).toMatchObject({ completedTasks: 2_000, avgCompileTimeMs: 20 }); + // Storage invariant: lifetime counters need no individual duration samples. + expect(retainedSamples()).toBe(0); + } finally { + clock.mockRestore(); + } + }); + + it("averages completed result durations including diagnostics but excludes structured worker errors", async () => { + const pool = createPool(1); + let now = 1_000; + const clock = vi.spyOn(Date, "now").mockImplementation(() => now); + try { + expect(pool.getStats().avgCompileTimeMs).toBe(0); + const success = pool.compile({ code: "success" }); + now += 10; + succeed(0); + await success; + const diagnostic = pool.compile({ code: "diagnostic" }); + now += 30; + worker(0).emit("message", { + type: "compile_result", + payload: { result: { ...successfulResult, success: false } }, + }); + await diagnostic; + const failure = pool.compile({ code: "worker-error" }); + const rejected = expect(failure).rejects.toThrow("worker error"); + now += 50; + worker(0).emit("message", { + type: "compile_result", + payload: { error: { message: "worker error" } }, + }); + await rejected; + expect(pool.getStats()).toMatchObject({ completedTasks: 2, failedTasks: 1, avgCompileTimeMs: 20 }); + } finally { + clock.mockRestore(); + } + }); + it("uses all workers before applying backpressure", async () => { const pool = createPool(2);