60d71da90e
Cuts CPU on the `engine/v1/worker-actions/*` routes a managed supervisor calls, and adds the benchmark harness the numbers come from. Measured on a local stack: **on-CPU per completed run 9.07ms → 6.59ms (−27%)**, busy fraction 45.6% → 33.8%, with every worker-action p50 down 23–27%. Load was 5,000 runs / 24 virtual supervisors / 90s window / 30,120 requests / 0 errors. Query-count work from the same investigation is deliberately **not** here — it will follow as a separate PR. ## The three changes **1. Split the event-loop monitor in two (~14% of on-CPU, plus ~5pp of GC).** `eventLoopMonitor.server.ts` installs a global `async_hooks` hook: `init` writes a `Map` entry for *every* async resource the process creates, `before` calls `process.hrtime()` and `context.active()` on every one. Enabling any async hook also puts V8 on the slow path for promise instrumentation process-wide. `EVENT_LOOP_MONITOR_ENABLED` defaulted to `"1"`, so this was the shipping configuration. The blocked-loop detector is now opt-in (`EVENT_LOOP_MONITOR_ENABLED`, default `0`). The event-loop *utilization* gauge — a single interval timer with no per-request cost — moves to its own flag (`EVENT_LOOP_UTILIZATION_MONITOR_ENABLED`, default `1`) and stays on, so the useful half survives without the expensive half. A/B under identical load: | | monitor on | monitor off | change | |---|---|---|---| | on-CPU per run | 9.08ms | 7.25ms | −20% | | GC self time | 9.80% | 5.05% | −4.75pp | | dequeue p50 | 76.6ms | 62.8ms | −18% | | attempts/start p50 | 56.3ms | 43.5ms | −23% | **2. Bucket route matching by first static path segment (10.4% → 3.9% of on-CPU).** `patches/@remix-run__router@1.23.3.patch` already memoized flattened branches and compiled path regexes. What remained was the linear scan: `matchRouteBranch` walked the ranked branch list calling `matchPath` per branch across 521 route files, so every worker-action request paid a scan proportional to the whole route table. Branches are now indexed by their lowercased leading segment, with one always-considered list for branches whose leading segment is dynamic, splat or optional (and for root/pathless paths). A request walks only its own bucket merged with that list. Route-matching self time dropped 64% (3.6s → 1.3s over a 90s window). Ordering is preserved exactly: both lists hold indexes into the already rank-sorted branch array and are walked in ascending-index order, so the first match found is the same branch the full scan would have found. Bucketing lowercases on both sides, so case-insensitive matching still resolves and `caseSensitive: true` routes are still rejected by `matchPath` itself. A pathname whose own leading segment can't be bucketed falls back to the full scan. Verified equivalent to the unpatched matcher over 20,050 pathnames (literal, dynamic, splat, optional, case variants, basenames, percent-encoded) with zero mismatches. `apps/webapp/test/routeMatchingPatch.test.ts` pins the matching semantics rather than the optimisation, so it still passes without the patch. **3. Demote per-heartbeat and per-dequeue `info` logs to `debug`.** These are the two highest-rate engine calls and each wrote a synchronous structured log line on every request. Synchronous `console` writes can block the loop when stdout backs up, which costs more than the ~1.3% CPU share suggests. ## The harness Two benchmarks, neither in the default suite (they run for minutes, attach the V8 profiler, and report numbers rather than assert on them). See `apps/webapp/test/bench/README.md`. - `apps/webapp/test/bench/engineHttp.bench.test.ts` — spawns a real webapp against throwaway Postgres/Redis containers, seeds a production environment with a promoted managed deployment, and drives a closed-loop supervisor pool through the full lifecycle. Profiling runs over CDP rather than `--cpu-prof` so it covers only the measured window instead of being swamped by boot, and `performance.eventLoopUtilization()` is sampled *inside* the webapp process. - `internal-packages/run-engine/src/engine/bench/runEngineLifecycle.bench.test.ts` — drives `RunEngine` directly, profiling enqueue and lifecycle separately so engine cost isn't mixed with request-stack overhead. - `apps/webapp/test/bench/analyzeProfile.ts` — dependency-free `.cpuprofile` analyzer that symbolicates through the build's source maps and ranks CPU by package, self time and total time. Percentages are shares of on-CPU time (V8's `(idle)`/`(program)` excluded). `startWebapp` gains `overrideEnv`, applied after the worker-disable defaults, so the HTTP bench can re-enable the run engine worker that drains the master queue into the worker queues a supervisor dequeues from. The local OTel collector gains a traces pipeline. It only defined a metrics pipeline, so pointing `INTERNAL_OTEL_TRACE_EXPORTER_URL` at it locally failed and the webapp silently fell back to the console span logger. ## Configuration For operators upgrading: - `EVENT_LOOP_MONITOR_ENABLED` (now defaults to `0`) — the per-async-resource blocked-loop detector. Set to `1` to restore the previous behaviour and keep emitting `event-loop-blocked` spans. - `EVENT_LOOP_UTILIZATION_MONITOR_ENABLED` (new, defaults to `1`) — the `nodejs.event_loop.utilization` gauge. Unchanged in behaviour; it just has its own flag now so it survives turning the detector off. ## Notes for review - `pnpm-lock.yaml` changes only because the router patch content changed, which changes its patch hash. - One thing the profile ruled out: with a real OTLP collector receiving spans, tracing costs ~1.7% of on-CPU at 100% sampling and ~0.8% at the production rate. Span shipping is not a hidden cost, so nothing here touches it. - Caveats on the numbers: a laptop, not production hardware, so DB and Redis *latency* are unrepresentative (client-side CPU is what's ranked); single webapp process; throughput varies ~5% run to run, which is why the claims rest on on-CPU per run rather than req/s. ## Verification - 20,050-pathname router equivalence check vs the unpatched matcher, zero mismatches - `apps/webapp/test/routeMatchingPatch.test.ts` (12 cases) passes - webapp e2e smoke suite (68 tests) passes through the patched router - run-engine suites covering the snapshot/attempt paths pass - `typecheck`, `format`, `lint`, `knip` clean
289 lines
10 KiB
TypeScript
289 lines
10 KiB
TypeScript
/**
|
|
* CPU benchmark for the engine-facing HTTP surface: the
|
|
* `engine/v1/worker-actions/*` routes a managed supervisor calls.
|
|
*
|
|
* Unlike the run-engine bench, this measures the whole request stack — Remix
|
|
* routing, `createActionWorkerApiRoute`, worker-token auth, zod validation,
|
|
* JSON encode/decode — on top of the engine work, which is where a large share
|
|
* of the production engine service's event-loop time actually goes.
|
|
*
|
|
* The webapp runs as a child process with `--inspect`, and the profiler is
|
|
* driven over CDP so the profile covers only the measured window rather than
|
|
* boot. Artifacts land in `.bench/` at the repo root. Analyze one with:
|
|
*
|
|
* pnpm --filter webapp exec tsx test/bench/analyzeProfile.ts <path>
|
|
*
|
|
* Knobs, all optional:
|
|
* BENCH_RUNS, BENCH_SUPERVISORS, BENCH_HEARTBEATS, BENCH_DURATION_MS,
|
|
* BENCH_OUT_DIR, BENCH_SAMPLING_INTERVAL_US
|
|
*/
|
|
import { startTestServer, type TestServer } from "@internal/testcontainers/webapp";
|
|
import { createServer } from "node:net";
|
|
import { mkdir, writeFile } from "node:fs/promises";
|
|
import { join } from "node:path";
|
|
import { afterAll, beforeAll, describe, expect, it, vi } from "vitest";
|
|
import { WebappProfiler } from "./lib/cdp";
|
|
import { seedEngineFixtures, workerHeaders, type EngineFixtures } from "./lib/engineFixtures";
|
|
import { formatStatsTable, LatencyRecorder, runLoad } from "./lib/loadDriver";
|
|
|
|
vi.setConfig({ testTimeout: 900_000 });
|
|
|
|
const RUNS = Number(process.env.BENCH_RUNS ?? 1200);
|
|
const SUPERVISORS = Number(process.env.BENCH_SUPERVISORS ?? 16);
|
|
const HEARTBEATS_PER_RUN = Number(process.env.BENCH_HEARTBEATS ?? 2);
|
|
const DURATION_MS = Number(process.env.BENCH_DURATION_MS ?? 60_000);
|
|
const SAMPLING_INTERVAL_US = Number(process.env.BENCH_SAMPLING_INTERVAL_US ?? 200);
|
|
const OUT_DIR = process.env.BENCH_OUT_DIR ?? join(process.cwd(), "..", "..", ".bench");
|
|
const PROFILE_NAME = process.env.BENCH_PROFILE_NAME ?? "engine-http";
|
|
|
|
/**
|
|
* JSON object merged into the spawned webapp's environment, for A/B runs
|
|
* against a single flag, e.g.
|
|
* `BENCH_EXTRA_ENV='{"EVENT_LOOP_MONITOR_ENABLED":"0"}'`.
|
|
*/
|
|
const EXTRA_WEBAPP_ENV: Record<string, string> = process.env.BENCH_EXTRA_ENV
|
|
? JSON.parse(process.env.BENCH_EXTRA_ENV)
|
|
: {};
|
|
|
|
async function findFreePort(): Promise<number> {
|
|
return new Promise((resolve, reject) => {
|
|
const server = createServer();
|
|
server.listen(0, () => {
|
|
const { port } = server.address() as { port: number };
|
|
server.close((err) => (err ? reject(err) : resolve(port)));
|
|
});
|
|
});
|
|
}
|
|
|
|
let server: TestServer;
|
|
let profiler: WebappProfiler;
|
|
let fixtures: EngineFixtures;
|
|
let inspectPort: number;
|
|
|
|
beforeAll(async () => {
|
|
inspectPort = await findFreePort();
|
|
server = await startTestServer({
|
|
extraEnv: {
|
|
NODE_OPTIONS: `--inspect=${inspectPort}`,
|
|
API_RATE_LIMIT_MAX: "1000000",
|
|
API_RATE_LIMIT_REFILL_RATE: "1000000",
|
|
...EXTRA_WEBAPP_ENV,
|
|
},
|
|
overrideEnv: { RUN_ENGINE_WORKER_ENABLED: "1" },
|
|
});
|
|
profiler = await WebappProfiler.attach(inspectPort);
|
|
}, 300_000);
|
|
|
|
afterAll(async () => {
|
|
profiler?.detach();
|
|
await server?.stop();
|
|
}, 120_000);
|
|
|
|
/**
|
|
* Seeds the worker queue. Failures are counted by status and surfaced rather
|
|
* than swallowed: a partially-filled queue silently turns the measured window
|
|
* into mostly-empty dequeues, which reads as a fast server rather than a
|
|
* broken fixture.
|
|
*/
|
|
async function triggerRuns(count: number): Promise<number> {
|
|
let triggered = 0;
|
|
const failures = new Map<number, { count: number; sample: string }>();
|
|
|
|
const batches = Math.ceil(count / 50);
|
|
for (let batch = 0; batch < batches; batch++) {
|
|
const size = Math.min(50, count - batch * 50);
|
|
const responses = await Promise.all(
|
|
Array.from({ length: size }, (_, i) => {
|
|
const index = batch * 50 + i;
|
|
const taskId = fixtures.taskIdentifiers[index % fixtures.taskIdentifiers.length]!;
|
|
return server.webapp.fetch(`/api/v1/tasks/${taskId}/trigger`, {
|
|
method: "POST",
|
|
headers: {
|
|
Authorization: `Bearer ${fixtures.environmentApiKey}`,
|
|
"content-type": "application/json",
|
|
},
|
|
body: JSON.stringify({ payload: { index, message: "bench payload" } }),
|
|
});
|
|
})
|
|
);
|
|
|
|
for (const res of responses) {
|
|
if (res.ok) {
|
|
triggered += 1;
|
|
continue;
|
|
}
|
|
const existing = failures.get(res.status);
|
|
if (existing) existing.count += 1;
|
|
else failures.set(res.status, { count: 1, sample: (await res.text()).slice(0, 200) });
|
|
}
|
|
}
|
|
|
|
if (failures.size > 0) {
|
|
const detail = [...failures.entries()]
|
|
.map(([status, { count: n, sample }]) => `${status} x${n} (${sample})`)
|
|
.join("; ");
|
|
console.warn(`[engine-http] ${count - triggered}/${count} triggers failed: ${detail}`);
|
|
}
|
|
|
|
if (triggered < count * 0.9) {
|
|
throw new Error(`only ${triggered}/${count} runs were queued; the measured window would idle`);
|
|
}
|
|
|
|
return triggered;
|
|
}
|
|
|
|
describe("engine worker-action HTTP CPU benchmark", () => {
|
|
it("profiles the supervisor request loop", async () => {
|
|
await mkdir(OUT_DIR, { recursive: true });
|
|
|
|
fixtures = await seedEngineFixtures(server.prisma, { taskCount: 4 });
|
|
|
|
const triggered = await triggerRuns(RUNS);
|
|
expect(triggered).toBeGreaterThan(0);
|
|
|
|
const recorder = new LatencyRecorder();
|
|
let dequeuedRuns = 0;
|
|
let completedRuns = 0;
|
|
let emptyDequeues = 0;
|
|
|
|
await profiler.startCpuProfile(SAMPLING_INTERVAL_US);
|
|
await profiler.startEluSampling();
|
|
recorder.begin();
|
|
|
|
await runLoad({
|
|
concurrency: SUPERVISORS,
|
|
durationMs: DURATION_MS,
|
|
iteration: async (workerIndex) => {
|
|
const headers = workerHeaders(fixtures, `bench-instance-${workerIndex}`);
|
|
|
|
const dequeueResponse = await recorder.time("POST worker-actions/dequeue", async () => {
|
|
const res = await server.webapp.fetch("/engine/v1/worker-actions/dequeue", {
|
|
method: "POST",
|
|
headers,
|
|
body: JSON.stringify({}),
|
|
});
|
|
if (!res.ok) throw new Error(`dequeue ${res.status}`);
|
|
return (await res.json()) as Array<{
|
|
run: { friendlyId: string };
|
|
snapshot: { friendlyId: string };
|
|
}>;
|
|
});
|
|
|
|
const message = dequeueResponse?.[0];
|
|
if (!message) {
|
|
emptyDequeues += 1;
|
|
await new Promise((r) => setTimeout(r, 25));
|
|
return;
|
|
}
|
|
|
|
dequeuedRuns += 1;
|
|
const runId = message.run.friendlyId;
|
|
let snapshotId = message.snapshot.friendlyId;
|
|
|
|
const attempt = await recorder.time("POST attempts/start", async () => {
|
|
const res = await server.webapp.fetch(
|
|
`/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/attempts/start`,
|
|
{ method: "POST", headers, body: JSON.stringify({}) }
|
|
);
|
|
if (!res.ok) throw new Error(`start ${res.status}`);
|
|
return (await res.json()) as { snapshot: { friendlyId: string } };
|
|
});
|
|
|
|
if (!attempt) return;
|
|
snapshotId = attempt.snapshot.friendlyId;
|
|
|
|
for (let beat = 0; beat < HEARTBEATS_PER_RUN; beat++) {
|
|
await recorder.time("POST snapshots/heartbeat", async () => {
|
|
const res = await server.webapp.fetch(
|
|
`/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/heartbeat`,
|
|
{ method: "POST", headers, body: JSON.stringify({}) }
|
|
);
|
|
if (!res.ok) throw new Error(`heartbeat ${res.status}`);
|
|
return res.json();
|
|
});
|
|
}
|
|
|
|
await recorder.time("GET snapshots/latest", async () => {
|
|
const res = await server.webapp.fetch(
|
|
`/engine/v1/worker-actions/runs/${runId}/snapshots/latest`,
|
|
{ headers }
|
|
);
|
|
if (!res.ok) throw new Error(`latest ${res.status}`);
|
|
return res.json();
|
|
});
|
|
|
|
const completed = await recorder.time("POST attempts/complete", async () => {
|
|
const res = await server.webapp.fetch(
|
|
`/engine/v1/worker-actions/runs/${runId}/snapshots/${snapshotId}/attempts/complete`,
|
|
{
|
|
method: "POST",
|
|
headers,
|
|
body: JSON.stringify({
|
|
completion: {
|
|
ok: true,
|
|
id: runId,
|
|
outputType: "application/json",
|
|
output: JSON.stringify({ done: true }),
|
|
},
|
|
}),
|
|
}
|
|
);
|
|
if (!res.ok) throw new Error(`complete ${res.status}`);
|
|
return res.json();
|
|
});
|
|
|
|
if (completed) completedRuns += 1;
|
|
},
|
|
});
|
|
|
|
recorder.end();
|
|
const elu = profiler.stopEluSampling();
|
|
const profile = await profiler.stopCpuProfile(join(OUT_DIR, `${PROFILE_NAME}.cpuprofile`));
|
|
|
|
const stats = recorder.stats();
|
|
const totals = recorder.totals();
|
|
|
|
console.log(`\n${formatStatsTable(stats)}`);
|
|
console.log(
|
|
`\n[engine-http] ${totals.count} requests (${totals.errors} errors) at ` +
|
|
`${totals.throughputPerSecond.toFixed(1)} req/s | ` +
|
|
`${dequeuedRuns} dequeued, ${completedRuns} completed, ${emptyDequeues} empty dequeues`
|
|
);
|
|
console.log(
|
|
`[engine-http] webapp ELU mean ${(elu.stats.mean * 100).toFixed(1)}% ` +
|
|
`p50 ${(elu.stats.p50 * 100).toFixed(1)}% ` +
|
|
`p95 ${(elu.stats.p95 * 100).toFixed(1)}% ` +
|
|
`max ${(elu.stats.max * 100).toFixed(1)}% (${elu.stats.sampleCount} samples)`
|
|
);
|
|
console.log(`[engine-http] profile: ${profile.path} (${profile.sampleCount} samples)`);
|
|
|
|
await writeFile(
|
|
join(OUT_DIR, `${PROFILE_NAME}-bench-summary.json`),
|
|
JSON.stringify(
|
|
{
|
|
config: {
|
|
runs: RUNS,
|
|
supervisors: SUPERVISORS,
|
|
heartbeatsPerRun: HEARTBEATS_PER_RUN,
|
|
durationMs: DURATION_MS,
|
|
samplingIntervalUs: SAMPLING_INTERVAL_US,
|
|
},
|
|
triggered,
|
|
dequeuedRuns,
|
|
completedRuns,
|
|
emptyDequeues,
|
|
totals,
|
|
operations: stats,
|
|
elu: elu.stats,
|
|
eluSamples: elu.samples,
|
|
profilePath: profile.path,
|
|
},
|
|
null,
|
|
2
|
|
)
|
|
);
|
|
|
|
expect(dequeuedRuns).toBeGreaterThan(0);
|
|
});
|
|
});
|