Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions server/routes/status.routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 };
}

Expand Down Expand Up @@ -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,
Expand Down
9 changes: 4 additions & 5 deletions server/services/compilation-worker-pool.ts
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,7 @@ export class CompilationWorkerPool {
totalTasks: 0,
completedTasks: 0,
failedTasks: 0,
compileTimes: [] as number[],
totalCompileTimeMs: 0,
};

constructor(numWorkers?: number, options: CompilationWorkerPoolOptions = {}) {
Expand Down Expand Up @@ -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`,
);
Expand Down Expand Up @@ -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 {
Expand Down
25 changes: 15 additions & 10 deletions server/services/server-metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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,
Expand Down
37 changes: 37 additions & 0 deletions tests/server/routes/status-compile-capacity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((resolve) => (server ? server.close(() => resolve()) : resolve())));
Expand All @@ -30,6 +31,42 @@ async function status(): Promise<Record<string, any>> {
}

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 });

Expand Down
60 changes: 60 additions & 0 deletions tests/server/services/server-metrics-cpu.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
53 changes: 53 additions & 0 deletions tests/server/worker-pool.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> }).stats,
).reduce<number>((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);

Expand Down
Loading