Files
lume-ctrl/scripts/m0-spike.mjs
T

934 lines
40 KiB
JavaScript

#!/usr/bin/env node
import assert from "node:assert/strict";
import { spawn, spawnSync } from "node:child_process";
import { createHash, randomUUID } from "node:crypto";
import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises";
import http from "node:http";
import https from "node:https";
import net from "node:net";
import os from "node:os";
import path from "node:path";
import { fileURLToPath, pathToFileURL } from "node:url";
const REDACTED = "[REDACTED]";
const DEBUG = process.env.M0_DEBUG === "1";
const SENSITIVE_HEADERS = new Set([
"authorization",
"cookie",
"location",
"proxy-authorization",
"referer",
"sec-websocket-accept",
"sec-websocket-key",
"set-cookie",
"x-api-key",
]);
export function redactHeaders(headers = {}) {
return Object.fromEntries(
Object.entries(headers)
.sort(([a], [b]) => a.localeCompare(b))
.map(([name, value]) => [
name.toLowerCase(),
SENSITIVE_HEADERS.has(name.toLowerCase()) ? REDACTED : String(value),
]),
);
}
export function safeUrl(raw) {
try {
const value = new URL(raw);
return `${value.origin}${value.pathname}${value.search ? "?<redacted>" : ""}`;
} catch {
return "<invalid-url>";
}
}
export function boundedBody(body, base64Encoded, maxBytes, source = "network") {
const bytes = Buffer.from(body, base64Encoded ? "base64" : "utf8");
if (bytes.length > maxBytes) {
return { state: "size_limit", source, encoding: base64Encoded ? "base64" : "utf8", storage: "omitted", reason: "size_limit", bytes: bytes.length };
}
return {
state: "fetched",
source,
encoding: base64Encoded ? "base64" : "utf8",
storage: "omitted",
reason: "default_body_policy",
bytes: bytes.length,
};
}
export function queryRecords(records, sessionId, { cursor = 0, limit = 100 } = {}) {
assert(sessionId, "sessionId is required");
assert(Number.isSafeInteger(cursor) && cursor >= 0, "cursor must be a non-negative integer");
assert(Number.isSafeInteger(limit) && limit > 0, "limit must be a positive integer");
// ponytail: linear scan fits the single-instance Spike; index by session if retained volume grows.
const matching = records.filter((record) => record.session_id === sessionId);
const items = matching.slice(cursor, cursor + limit);
return { items, next_cursor: cursor + items.length < matching.length ? cursor + items.length : null };
}
export function exportSession(records, sessionId) {
return queryRecords(records, sessionId, { limit: Number.MAX_SAFE_INTEGER }).items
.map((record) => JSON.stringify(record))
.join("\n");
}
export function shouldRetryStartup(detachedReason) {
return detachedReason === "Render process gone.";
}
class StartupNavigateFailure extends Error {
constructor(cause, detachedReason) {
super(cause.message, { cause });
this.detachedReason = detachedReason;
}
}
class CDP {
constructor(url) {
this.url = url;
this.nextId = 1;
this.pending = new Map();
this.listeners = [];
this.closeListeners = [];
this.closed = new Promise((resolve) => (this.resolveClosed = resolve));
this.detached = new Promise((resolve) => (this.resolveDetached = resolve));
}
async connect(timeoutMs = 20_000) {
this.ws = new WebSocket(this.url);
await new Promise((resolve, reject) => {
const timer = setTimeout(() => {
this.ws.close();
reject(new Error("CDP WebSocket connection timed out"));
}, timeoutMs);
this.ws.onopen = () => { clearTimeout(timer); resolve(); };
this.ws.onerror = () => { clearTimeout(timer); reject(new Error("CDP WebSocket connection failed")); };
});
this.ws.onmessage = ({ data }) => this.onMessage(JSON.parse(data));
this.ws.onclose = ({ code, reason, wasClean }) => {
const details = { code, reason: reason || "socket_closed", was_clean: wasClean };
for (const { reject } of this.pending.values()) reject(new Error(`CDP WebSocket closed: ${code} ${details.reason}`));
this.pending.clear();
for (const listener of this.closeListeners) listener(details);
this.resolveClosed(details);
};
return this;
}
onEvent(listener) {
this.listeners.push(listener);
}
onClose(listener) {
this.closeListeners.push(listener);
}
waitForDetached(timeoutMs) {
return Promise.race([
this.detached,
new Promise((resolve) => setTimeout(resolve, timeoutMs)),
]);
}
onMessage(message) {
if (DEBUG) console.error(
"CDP",
message.id ?? message.method,
message.method === "Inspector.detached" ? message.params :
message.method === "Network.requestWillBeSent" ? safeUrl(message.params.request.url) : "",
);
if (message.id) {
const pending = this.pending.get(message.id);
if (!pending) return;
this.pending.delete(message.id);
if (message.error) pending.reject(new Error(`${message.error.code}: ${message.error.message}`));
else pending.resolve(message.result);
return;
}
if (message.method === "Inspector.detached") {
this.detachedReason = message.params.reason;
this.resolveDetached(message.params.reason);
const error = new Error(`Inspector.detached: ${message.params.reason}`);
for (const { reject } of this.pending.values()) reject(error);
this.pending.clear();
}
for (const listener of this.listeners) listener(message.method, message.params ?? {});
}
call(method, params = {}, timeoutMs = 20_000) {
const id = this.nextId++;
return new Promise((resolve, reject) => {
const timer = setTimeout(() => {
this.pending.delete(id);
reject(new Error(`CDP command timed out: ${method}`));
}, timeoutMs);
this.pending.set(id, {
resolve: (value) => { clearTimeout(timer); resolve(value); },
reject: (error) => { clearTimeout(timer); reject(error); },
});
this.ws.send(JSON.stringify({ id, method, params }));
});
}
async close() {
if (this.ws.readyState < WebSocket.CLOSING) this.ws.close();
await this.closed;
}
}
function listen(server, port = 0) {
return new Promise((resolve, reject) => {
server.once("error", reject);
server.listen(port, "127.0.0.1", () => resolve(server.address().port));
});
}
async function unusedPort() {
const server = net.createServer();
const port = await listen(server);
await new Promise((resolve) => server.close(resolve));
return port;
}
function websocketFrame(opcode, payload = "") {
const body = Buffer.from(payload);
assert(body.length <= 65_535, "fixture frame must fit a 16-bit length");
const header = body.length < 126
? Buffer.from([0x80 | opcode, body.length])
: Buffer.from([0x80 | opcode, 126, body.length >> 8, body.length & 0xff]);
return Buffer.concat([header, body]);
}
function acceptWebsocket(request, socket) {
const accept = createHash("sha1")
.update(`${request.headers["sec-websocket-key"]}258EAFA5-E914-47DA-95CA-C5AB0DC85B11`)
.digest("base64");
socket.write(
"HTTP/1.1 101 Switching Protocols\r\n" +
"Upgrade: websocket\r\n" +
"Connection: Upgrade\r\n" +
`Sec-WebSocket-Accept: ${accept}\r\n\r\n`,
);
let buffer = Buffer.alloc(0);
socket.on("data", (chunk) => {
buffer = Buffer.concat([buffer, chunk]);
while (buffer.length >= 2) {
const masked = Boolean(buffer[1] & 0x80);
let length = buffer[1] & 0x7f;
let offset = 2;
if (length === 126) {
if (buffer.length < 4) return;
length = buffer.readUInt16BE(2);
offset = 4;
}
const maskBytes = masked ? 4 : 0;
if (buffer.length < offset + maskBytes + length) return;
const opcode = buffer[0] & 0x0f;
const mask = buffer.subarray(offset, offset + maskBytes);
const payload = Buffer.from(buffer.subarray(offset + maskBytes, offset + maskBytes + length));
if (masked) for (let i = 0; i < payload.length; i++) payload[i] ^= mask[i % 4];
buffer = buffer.subarray(offset + maskBytes + length);
if (opcode === 1) socket.write(websocketFrame(1, "server-message"));
if (opcode === 8) {
socket.end(websocketFrame(8));
return;
}
}
});
}
async function startFixtures(tempDir) {
const secrets = {
authorization: randomUUID(),
body: randomUUID(),
cookie: randomUUID(),
setCookie: randomUUID(),
};
const send = (response, headers, body) => {
response.writeHead(200, {
...headers,
connection: "close",
"content-length": Buffer.byteLength(body),
});
response.end(body);
};
const key = path.join(tempDir, "fixture-key.pem");
const cert = path.join(tempDir, "fixture-cert.pem");
const openssl = spawnSync("openssl", [
"req", "-x509", "-newkey", "rsa:2048", "-nodes", "-days", "1",
"-subj", "/CN=127.0.0.1", "-addext", "subjectAltName=IP:127.0.0.1",
"-keyout", key, "-out", cert,
], { stdio: "ignore" });
assert.equal(openssl.status, 0, "openssl must generate the local HTTPS fixture certificate");
const httpsServer = https.createServer(
{ key: await readFile(key), cert: await readFile(cert) },
(_request, response) => {
send(response, {
"access-control-allow-origin": "*",
"access-control-allow-headers": "authorization",
"content-type": "application/json",
}, '{"secure":true}');
},
);
const httpsPort = await listen(httpsServer);
const failurePort = await unusedPort();
const pageHtml = "<!doctype html><title>M0_READY</title>";
const httpServer = http.createServer((request, response) => {
if (DEBUG) console.error("HTTP", request.url);
const cors = {
"access-control-allow-origin": "*",
"access-control-allow-headers": "authorization",
};
if (request.method === "OPTIONS") {
send(response, cors, "");
return;
}
if (request.url === "/ok") {
send(response, {
...cors,
"content-type": "application/json",
"set-cookie": `m0_response=${secrets.setCookie}; SameSite=Lax`,
}, '{"ok":true}');
return;
}
if (request.url === "/large") {
send(response, { ...cors, "content-type": "text/plain" }, "x".repeat(2_048));
return;
}
if (request.url === "/binary") {
send(response, { ...cors, "content-type": "application/octet-stream" }, Buffer.from([0, 255, 1, 254]));
return;
}
if (request.url === "/cache") {
send(response, { ...cors, "cache-control": "public, max-age=3600", "content-type": "application/json" }, '{"cached":true}');
return;
}
if (request.url === "/secret") {
send(response, { ...cors, "content-type": "application/json" }, JSON.stringify({ secret: secrets.body }));
return;
}
send(response, { "content-type": "text/html" }, pageHtml);
});
httpServer.on("upgrade", (request, socket) => {
if (DEBUG) console.error("HTTP upgrade", request.url);
if (request.url === "/ws") acceptWebsocket(request, socket);
else socket.destroy();
});
const httpPort = await listen(httpServer, Number(process.env.M0_HTTP_PORT ?? 0));
const pageScript = `
(async () => {
document.title = "M0_RUNNING";
try {
await (await fetch("http://127.0.0.1:${httpPort}/ok")).text();
await (await fetch("http://127.0.0.1:${httpPort}/secret")).text();
await (await fetch("http://127.0.0.1:${httpPort}/large")).text();
await (await fetch("http://127.0.0.1:${httpPort}/binary")).arrayBuffer();
await (await fetch("http://127.0.0.1:${httpPort}/cache", { cache: "force-cache" })).text();
await (await fetch("http://127.0.0.1:${httpPort}/cache", { cache: "force-cache" })).text();
await (await fetch("https://127.0.0.1:${httpsPort}/secure")).text();
await new Promise((resolve, reject) => {
const ws = new WebSocket("ws://127.0.0.1:${httpPort}/ws");
let received = 0;
ws.onopen = () => {
ws.send("x".repeat(2048));
for (let index = 0; index < 20; index++) ws.send("client-message");
};
ws.onmessage = () => { if (++received === 21) ws.close(); };
ws.onclose = resolve;
ws.onerror = reject;
});
try { await fetch("http://127.0.0.1:${failurePort}/missing"); } catch {}
document.title = "M0_DONE";
} catch (error) {
document.title = "M0_ERROR:" + error.name;
}
})()`;
return {
pageHtml,
pageScript,
secrets,
url: `http://127.0.0.1:${httpPort}/`,
close: async () => Promise.all([
new Promise((resolve) => httpServer.close(resolve)),
new Promise((resolve) => httpsServer.close(resolve)),
]),
};
}
async function fixtureChild() {
const fixtures = await startFixtures(process.env.M0_FIXTURE_DIR);
const { close, ...config } = fixtures;
process.stdout.write(`${JSON.stringify(config)}\n`);
const stop = async () => {
await close();
process.exit(0);
};
process.on("SIGINT", stop);
process.on("SIGTERM", stop);
}
async function startFixtureProcess(tempDir) {
const child = spawn(process.execPath, [fileURLToPath(import.meta.url), "--fixture-child"], {
env: { ...process.env, M0_FIXTURE_DIR: tempDir },
stdio: ["ignore", "pipe", DEBUG ? "inherit" : "ignore"],
});
const config = await new Promise((resolve, reject) => {
let output = "";
const onData = (chunk) => {
output += chunk;
const newline = output.indexOf("\n");
if (newline < 0) return;
child.stdout.off("data", onData);
resolve(JSON.parse(output.slice(0, newline)));
};
child.stdout.on("data", onData);
child.once("exit", (code) => reject(new Error(`fixture process exited before readiness: ${code}`)));
});
return {
...config,
close: async () => {
if (child.exitCode !== null) return;
child.kill("SIGTERM");
await new Promise((resolve) => child.once("exit", resolve));
},
};
}
async function getJson(url, options = {}) {
const response = await fetch(url, options);
if (!response.ok) throw new Error(`HTTP ${response.status} from CDP discovery endpoint`);
return response.json();
}
async function waitFor(check, label, timeoutMs = 10_000) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
const value = await check();
if (value) return value;
await new Promise((resolve) => setTimeout(resolve, 100));
}
throw new Error(`Timed out waiting for ${label}`);
}
function processSnapshot(rootPid) {
const result = spawnSync("ps", ["-eo", "pid=,ppid=,rss=,pcpu=,comm="], { encoding: "utf8" });
assert.equal(result.status, 0, "ps must report the browser process tree");
const rows = result.stdout.trim().split("\n").map((line) => {
const [pid, ppid, rss, cpu, ...command] = line.trim().split(/\s+/);
return { pid: Number(pid), ppid: Number(ppid), rss: Number(rss), cpu: Number(cpu), command: command.join(" ") };
});
const pids = new Set([rootPid]);
for (let changed = true; changed;) {
changed = false;
for (const row of rows) if (pids.has(row.ppid) && !pids.has(row.pid)) {
pids.add(row.pid);
changed = true;
}
}
const tree = rows.filter((row) => pids.has(row.pid));
return {
process_count: tree.length,
rss_kib: tree.reduce((total, row) => total + row.rss, 0),
cpu_percent_snapshot: Number(tree.reduce((total, row) => total + row.cpu, 0).toFixed(1)),
};
}
export function summarizeProcessSamples(samples, intervalMs) {
assert(samples.length > 1, "resource baseline requires multiple samples");
const summary = (name) => {
const values = samples.map((sample) => sample[name]).sort((a, b) => a - b);
return {
min: values[0],
max: values.at(-1),
mean: Number((values.reduce((total, value) => total + value, 0) / values.length).toFixed(1)),
p95: values[Math.ceil(values.length * 0.95) - 1],
};
};
return {
sample_count: samples.length,
interval_ms: intervalMs,
process_count: { min: Math.min(...samples.map((sample) => sample.process_count)), max: Math.max(...samples.map((sample) => sample.process_count)) },
rss_kib: summary("rss_kib"),
cpu_percent: summary("cpu_percent_snapshot"),
};
}
async function processBaseline(rootPid, sampleCount, intervalMs) {
const samples = [];
for (let index = 0; index < sampleCount; index++) {
if (index) await new Promise((resolve) => setTimeout(resolve, intervalMs));
samples.push(processSnapshot(rootPid));
}
return summarizeProcessSamples(samples, intervalMs);
}
function isRunning(pid) {
try {
process.kill(pid, 0);
return true;
} catch {
return false;
}
}
export function eventCollector(sessionId, targetId, {
maxWebSocketFrameBytes = 1_024,
maxWebSocketEventsPerSecond = 16,
maxWebSocketEventsPerConnection = 16,
} = {}) {
for (const [name, value] of Object.entries({ maxWebSocketFrameBytes, maxWebSocketEventsPerSecond, maxWebSocketEventsPerConnection })) {
assert(Number.isSafeInteger(value) && value > 0, `${name} must be a positive integer`);
}
const records = [];
const requestState = new Map();
const bodySources = new Map();
const webSocketLimits = new Map();
const push = (kind, fields = {}) => {
const record = {
schema_version: "1.0",
session_id: sessionId,
target_id: targetId,
kind,
...fields,
};
records.push(record);
return record;
};
const dropWebSocketFrame = (connectionId, reason, payloadBytes) => {
let state = webSocketLimits.get(connectionId);
if (!state) {
state = { window: null, windowEvents: 0, emitted: 0 };
webSocketLimits.set(connectionId, state);
}
if (!state.dropRecord) state.dropRecord = push("websocket_frames_dropped", {
connection_id: connectionId,
dropped: {
count: 0,
reasons: {},
max_payload_bytes: 0,
limits: {
frame_bytes: maxWebSocketFrameBytes,
events_per_second: maxWebSocketEventsPerSecond,
events_per_connection: maxWebSocketEventsPerConnection,
},
},
});
const dropped = state.dropRecord.dropped;
dropped.count++;
dropped.reasons[reason] = (dropped.reasons[reason] ?? 0) + 1;
dropped.max_payload_bytes = Math.max(dropped.max_payload_bytes, payloadBytes);
};
const listener = (method, params) => {
const requestId = params.requestId;
if (requestId && !requestState.has(requestId)) requestState.set(requestId, new Set());
const mark = (value) => requestState.get(requestId)?.add(value);
switch (method) {
case "Network.requestWillBeSent":
mark("request");
bodySources.set(requestId, "network");
push("http_request", {
request_id: requestId,
at: params.wallTime,
request: {
method: params.request.method,
url: safeUrl(params.request.url),
headers: redactHeaders(params.request.headers),
},
});
break;
case "Network.requestServedFromCache":
bodySources.set(requestId, "memory_cache");
push("http_cache_hit", { request_id: requestId, source: "memory_cache" });
break;
case "Network.requestWillBeSentExtraInfo":
push("http_request_headers", { request_id: requestId, headers: redactHeaders(params.headers) });
break;
case "Network.responseReceived":
mark("response");
if (params.response.fromServiceWorker) bodySources.set(requestId, "service_worker");
else if (params.response.fromDiskCache) bodySources.set(requestId, "disk_cache");
else if (params.response.fromPrefetchCache) bodySources.set(requestId, "prefetch_cache");
push("http_response", {
request_id: requestId,
response: {
url: safeUrl(params.response.url),
status: params.response.status,
protocol: params.response.protocol,
mime_type: params.response.mimeType,
headers: redactHeaders(params.response.headers),
cache: bodySources.has(requestId) ? {
state: bodySources.get(requestId) === "network" ? "miss" : "hit",
source: bodySources.get(requestId),
} : { state: "unknown" },
},
});
break;
case "Network.responseReceivedExtraInfo":
push("http_response_headers", { request_id: requestId, status: params.statusCode, headers: redactHeaders(params.headers) });
break;
case "Network.loadingFinished":
mark("finished");
push("http_finished", { request_id: requestId, encoded_bytes: params.encodedDataLength });
break;
case "Network.loadingFailed":
mark("failed");
push("http_failed", {
request_id: requestId,
error: { text: params.errorText, canceled: params.canceled, blocked_reason: params.blockedReason },
});
break;
case "Network.webSocketCreated":
mark("ws_created");
push("websocket_created", { connection_id: requestId, url: safeUrl(params.url) });
break;
case "Network.webSocketWillSendHandshakeRequest":
push("websocket_handshake_request", { connection_id: requestId, headers: redactHeaders(params.request.headers) });
break;
case "Network.webSocketHandshakeResponseReceived":
push("websocket_handshake_response", { connection_id: requestId, status: params.response.status, headers: redactHeaders(params.response.headers) });
break;
case "Network.webSocketFrameSent":
case "Network.webSocketFrameReceived": {
const direction = method.endsWith("Sent") ? "sent" : "received";
mark(`ws_${direction}`);
const encodedPayloadBytes = Buffer.byteLength(params.response.payloadData);
const payloadBytes = params.response.opcode === 2
? Buffer.from(params.response.payloadData, "base64").length
: encodedPayloadBytes;
let state = webSocketLimits.get(requestId);
if (!state) {
state = { window: null, windowEvents: 0, emitted: 0 };
webSocketLimits.set(requestId, state);
}
const window = Math.floor(params.timestamp ?? 0);
if (state.window !== window) {
state.window = window;
state.windowEvents = 0;
}
state.windowEvents++;
let dropReason;
if (payloadBytes > maxWebSocketFrameBytes) dropReason = "frame_size_limit";
else if (state.windowEvents > maxWebSocketEventsPerSecond) dropReason = "rate_limit";
else if (state.emitted >= maxWebSocketEventsPerConnection) dropReason = "event_limit";
if (dropReason) {
dropWebSocketFrame(requestId, dropReason, payloadBytes);
break;
}
state.emitted++;
push("websocket_frame", {
connection_id: requestId,
direction,
opcode: params.response.opcode,
payload_bytes: payloadBytes,
...(params.response.opcode === 2 && { encoded_payload_bytes: encodedPayloadBytes }),
payload: { state: "omitted", reason: "frame_policy" },
});
break;
}
case "Network.webSocketClosed":
mark("ws_closed");
push("websocket_closed", { connection_id: requestId });
break;
}
};
return { records, requestState, bodySources, push, listener };
}
function findRequest(records, pathname) {
return records.find((record) => record.kind === "http_request" && new URL(record.request.url.replace("?<redacted>", "")).pathname === pathname)?.request_id;
}
function findRequests(records, pathname) {
return records.filter((record) => record.kind === "http_request" && new URL(record.request.url.replace("?<redacted>", "")).pathname === pathname).map((record) => record.request_id);
}
async function collectBody(cdp, collector, requestId, maxBodyBytes) {
const source = collector.bodySources.get(requestId) ?? "network";
try {
const result = await cdp.call("Network.getResponseBody", { requestId });
collector.push("http_body", { request_id: requestId, body: boundedBody(result.body, result.base64Encoded, maxBodyBytes, source) });
} catch (error) {
collector.push("http_body", {
request_id: requestId,
body: { state: "unavailable", source, encoding: "unknown", storage: "omitted", reason: "cdp_error", error: error.message },
});
}
}
function assertCoverage(collector, secrets) {
const { records, requestState } = collector;
for (const pathname of ["/ok", "/secret", "/large", "/binary", "/cache", "/secure"]) {
const id = findRequest(records, pathname);
assert(id, `request observed for ${pathname}`);
assert.deepEqual([...requestState.get(id)].filter((value) => ["request", "response", "finished"].includes(value)).sort(), ["finished", "request", "response"]);
}
const failed = [...requestState.entries()].find(([, state]) => state.has("failed"));
assert(failed?.[1].has("request"), "failed request correlates to its requestId");
const websocket = [...requestState.entries()].find(([, state]) => state.has("ws_created"));
assert(websocket, "WebSocket connection observed");
for (const event of ["ws_sent", "ws_received", "ws_closed"]) assert(websocket[1].has(event), `${event} correlates to the WebSocket connection`);
const sample = records.map((record) => JSON.stringify(record)).join("\n");
for (const secret of Object.values(secrets)) assert(!sample.includes(secret), "sample does not leak a fixture credential");
assert(records.some((record) => record.kind === "http_request_headers" && record.headers.cookie === REDACTED), "Cookie is redacted");
assert(records.some((record) => record.kind === "http_request_headers" && record.headers.authorization === REDACTED), "Authorization is redacted");
assert(records.some((record) => record.kind === "http_response_headers" && record.headers["set-cookie"] === REDACTED), "Set-Cookie is redacted");
const secretBody = records.find((record) => record.kind === "http_body" && record.request_id === findRequest(records, "/secret"));
assert.deepEqual(secretBody?.body, { state: "fetched", source: "network", encoding: "utf8", storage: "omitted", reason: "default_body_policy", bytes: Buffer.byteLength(JSON.stringify({ secret: secrets.body })) });
assert(records.some((record) => record.kind === "http_body" && record.body.state === "size_limit" && record.body.storage === "omitted"), "oversize body is omitted");
assert(records.some((record) => record.kind === "http_body" && record.body.state === "unavailable"), "unavailable body is structured");
assert(records.some((record) => record.kind === "http_body" && record.body.source.endsWith("_cache")), "cached body source is explicit");
assert(records.some((record) => record.kind === "http_body" && record.body.encoding === "base64"), "binary body encoding is explicit");
assert(records.some((record) => record.kind === "websocket_frames_dropped" && record.dropped.reasons.frame_size_limit), "oversize WebSocket frame drop is observable");
assert(records.some((record) => record.kind === "websocket_frames_dropped" && record.dropped.count > 0), "high-rate WebSocket drops are aggregated");
assert(records.some((record) => record.kind === "cdp_disconnected" && record.source === "socket"), "real CDP socket close is recorded");
assert(records.some((record) => record.kind === "cdp_reconnected" && record.verified), "CDP reconnect is verified");
}
function assertIsolation(records, sessions) {
for (const { sessionId, targetId, otherSessionId } of sessions) {
const queried = queryRecords(records, sessionId, { limit: 1_000 }).items;
assert(queried.length > 0, `records exist for ${sessionId}`);
assert(queried.every((record) => record.session_id === sessionId && record.target_id === targetId), `query stays inside ${sessionId}`);
const exported = exportSession(records, sessionId);
assert(!exported.includes(otherSessionId), `export does not leak ${otherSessionId}`);
assert(queried.some((record) => record.kind === "cdp_disconnected"), `disconnect stays observable for ${sessionId}`);
}
}
async function main(attemptFailures = [], startup = { requested_samples: 1, successful_samples: 1, max_renderer_failure_rate: 1 }, resourceSamples = []) {
const binary = process.env.CLARK_BINARY_PATH;
assert(binary, "set CLARK_BINARY_PATH to the clark-browser Chromium binary");
const outputDir = path.resolve(process.env.M0_OUTPUT_DIR ?? "artifacts");
const maxBodyBytes = Number(process.env.M0_MAX_BODY_BYTES ?? 1_024);
assert(Number.isSafeInteger(maxBodyBytes) && maxBodyBytes > 0, "M0_MAX_BODY_BYTES must be a positive integer");
const limits = {
maxWebSocketFrameBytes: Number(process.env.M0_MAX_WS_FRAME_BYTES ?? 1_024),
maxWebSocketEventsPerSecond: Number(process.env.M0_MAX_WS_EVENTS_PER_SECOND ?? 16),
maxWebSocketEventsPerConnection: Number(process.env.M0_MAX_WS_EVENTS_PER_CONNECTION ?? 16),
};
const resourceSampleCount = Number(process.env.M0_RESOURCE_SAMPLES ?? 5);
const resourceSampleIntervalMs = Number(process.env.M0_RESOURCE_SAMPLE_INTERVAL_MS ?? 200);
const maxRssKib = Number(process.env.M0_MAX_RSS_KIB ?? 1_500_000);
const maxProcessCount = Number(process.env.M0_MAX_PROCESS_COUNT ?? 20);
assert(Number.isSafeInteger(resourceSampleCount) && resourceSampleCount > 1, "M0_RESOURCE_SAMPLES must be an integer greater than one");
assert(Number.isSafeInteger(resourceSampleIntervalMs) && resourceSampleIntervalMs > 0, "M0_RESOURCE_SAMPLE_INTERVAL_MS must be a positive integer");
assert(Number.isSafeInteger(maxRssKib) && maxRssKib > 0, "M0_MAX_RSS_KIB must be a positive integer");
assert(Number.isSafeInteger(maxProcessCount) && maxProcessCount > 0, "M0_MAX_PROCESS_COUNT must be a positive integer");
const tempDir = await mkdtemp(path.join(os.tmpdir(), "lume-ctrl-m0-"));
const profile = path.join(tempDir, "profile");
const cdpPort = Number(process.env.M0_CDP_PORT ?? 0) || await unusedPort();
const fixtures = await startFixtureProcess(tempDir);
assert.equal((await fetch(fixtures.url)).status, 200, "fixture process must be reachable before Clark starts");
const browser = spawn(binary, [
"--headless=new",
"--no-sandbox",
"--use-mock-keychain",
`--remote-debugging-port=${cdpPort}`,
"--remote-debugging-address=127.0.0.1",
"--remote-allow-origins=*",
`--user-data-dir=${profile}`,
"--ignore-certificate-errors",
"--disable-gpu",
"--disable-features=WebGPU",
"about:blank",
], { stdio: ["ignore", "ignore", "pipe"] });
browser.stderr.pipe(process.stderr);
const sessionId = `m0-${randomUUID()}`;
const browserPid = browser.pid;
let page;
let browserClient;
let reconnectClient;
try {
const version = await waitFor(
() => getJson(`http://127.0.0.1:${cdpPort}/json/version`).catch(() => null),
"clark CDP endpoint",
20_000,
);
const targets = await getJson(`http://127.0.0.1:${cdpPort}/json/list`);
page = targets.find((target) => target.type === "page" && target.url === "about:blank");
assert(page, "clark must expose its startup page through /json/list");
const cdp = await new CDP(page.webSocketDebuggerUrl).connect();
const collector = eventCollector(sessionId, page.id, limits);
cdp.onClose((details) => collector.push("cdp_disconnected", { source: "socket", ...details }));
await cdp.call("Page.enable");
await cdp.call("Runtime.enable");
try {
await cdp.call("Page.navigate", { url: fixtures.url }, 60_000);
await waitFor(async () => {
const { result } = await cdp.call("Runtime.evaluate", { expression: "document.title", returnByValue: true });
return result.value === "M0_READY";
}, "fixture page readiness");
} catch (error) {
const detachedReason = cdp.detachedReason ?? await cdp.waitForDetached(1_000);
if (DEBUG) console.error("startup failure", { detached_reason: detachedReason, error: error.message });
if (shouldRetryStartup(detachedReason)) throw new StartupNavigateFailure(error, detachedReason);
throw error;
}
cdp.onEvent(collector.listener);
await cdp.call("Network.enable", { maxTotalBufferSize: 10_000_000, maxResourceBufferSize: 2_000_000 });
await cdp.call("Network.setExtraHTTPHeaders", { headers: {
Authorization: `Bearer ${fixtures.secrets.authorization}`,
Cookie: `m0_cookie=${fixtures.secrets.cookie}`,
} });
await cdp.call("Runtime.evaluate", { expression: fixtures.pageScript });
await waitFor(async () => {
const { result } = await cdp.call("Runtime.evaluate", { expression: "document.title", returnByValue: true });
if (String(result.value).startsWith("M0_ERROR")) throw new Error(result.value);
return result.value === "M0_DONE";
}, "fixture completion");
await waitFor(() => [...collector.requestState.values()].some((state) => state.has("ws_closed")), "WebSocket closure");
for (const pathname of ["/ok", "/secret", "/large", "/binary", "/secure"]) {
await collectBody(cdp, collector, findRequest(collector.records, pathname), maxBodyBytes);
}
const cachedId = findRequests(collector.records, "/cache").find((requestId) => {
const source = collector.bodySources.get(requestId);
return source && source !== "network";
});
assert(cachedId, "a repeated force-cache request must be served from browser cache");
await collectBody(cdp, collector, cachedId, maxBodyBytes);
const failedId = [...collector.requestState.entries()].find(([, state]) => state.has("failed"))?.[0];
await collectBody(cdp, collector, failedId, maxBodyBytes);
const resources = await processBaseline(browserPid, resourceSampleCount, resourceSampleIntervalMs);
assert(resources.rss_kib.max <= maxRssKib, `browser RSS ${resources.rss_kib.max} KiB exceeds ${maxRssKib} KiB`);
assert(resources.process_count.max <= maxProcessCount, `browser process count ${resources.process_count.max} exceeds ${maxProcessCount}`);
resourceSamples.push({ sample: startup.successful_samples, resources });
browserClient = await new CDP(version.webSocketDebuggerUrl).connect();
await browserClient.call("Target.closeTarget", { targetId: page.id });
await cdp.closed;
collector.push("target_closed", { reason: "Target.closeTarget" });
assert(isRunning(browserPid), "closing the page does not exit clark-browser");
assert(collector.records.some((record) => record.kind === "cdp_disconnected" && record.source === "socket"), "target closure emits a socket disconnect event");
const reconnectTarget = await getJson(`http://127.0.0.1:${cdpPort}/json/new?about:blank`, { method: "PUT" });
reconnectClient = await new CDP(reconnectTarget.webSocketDebuggerUrl).connect();
const reconnectSessionId = `m0-${randomUUID()}`;
const reconnectCollector = eventCollector(reconnectSessionId, reconnectTarget.id, limits);
reconnectClient.onClose((details) => reconnectCollector.push("cdp_disconnected", { source: "socket", ...details }));
const { result: reconnectResult } = await reconnectClient.call("Runtime.evaluate", { expression: "1 + 1", returnByValue: true });
assert.equal(reconnectResult.value, 2, "CDP reconnect can execute a command");
reconnectCollector.push("cdp_reconnected", { verified: true });
assert(isRunning(browserPid), "disconnecting and reconnecting CDP does not exit clark-browser");
await browserClient.call("Target.closeTarget", { targetId: reconnectTarget.id });
await reconnectClient.closed;
reconnectCollector.push("target_closed", { reason: "Target.closeTarget" });
assert(isRunning(browserPid), "closing a second session target does not exit clark-browser");
const records = [...collector.records, ...reconnectCollector.records];
assertCoverage({ ...collector, records }, fixtures.secrets);
assertIsolation(records, [
{ sessionId, targetId: page.id, otherSessionId: reconnectSessionId },
{ sessionId: reconnectSessionId, targetId: reconnectTarget.id, otherSessionId: sessionId },
]);
const rendererFailures = attemptFailures.filter((failure) => failure.detached_reason === "Render process gone.").length;
const startupAttempts = startup.successful_samples + attemptFailures.length;
const rendererFailureRate = Number((rendererFailures / startupAttempts).toFixed(3));
const rendererConclusion = rendererFailureRate <= startup.max_renderer_failure_rate ? "within_threshold" : "exceeds_threshold";
if (startup.successful_samples === startup.requested_samples) {
assert.equal(rendererConclusion, "within_threshold", `renderer failure rate ${rendererFailureRate} exceeds ${startup.max_renderer_failure_rate}`);
}
await mkdir(outputDir, { recursive: true });
await writeFile(path.join(outputDir, "m0-events.sample.jsonl"), `${records.map((record) => JSON.stringify(record)).join("\n")}\n`);
await writeFile(path.join(outputDir, "m0-report.json"), `${JSON.stringify({
result: "pass",
clark: {
wrapper_version: "0.2.1",
chromium_version: version.Browser,
protocol_version: version["Protocol-Version"],
cdp_bind: "127.0.0.1",
},
coverage: {
http: "pass",
failed_request: "pass",
https_without_mitm: "pass",
websocket: "pass",
response_body_success: "pass",
response_body_size_limit: "pass",
response_body_unavailable: "pass",
cached_body_semantics: "pass",
binary_body_semantics: "pass",
websocket_limits: "pass",
multi_session_isolation: "pass",
target_close_survives: "pass",
cdp_disconnect_reconnect: "pass",
sensitive_header_redaction: "pass",
sensitive_body_redaction: "pass",
},
resources,
resource_samples: resourceSamples,
resource_thresholds: { max_rss_kib: maxRssKib, max_process_count: maxProcessCount, conclusion: "within_threshold" },
startup: {
...startup,
attempts: startupAttempts,
renderer_failures: rendererFailures,
renderer_failure_rate: rendererFailureRate,
conclusion: rendererConclusion,
},
event_count: records.length,
max_body_bytes: maxBodyBytes,
websocket_limits: {
frame_bytes: limits.maxWebSocketFrameBytes,
events_per_second: limits.maxWebSocketEventsPerSecond,
events_per_connection: limits.maxWebSocketEventsPerConnection,
},
attempt_failures: attemptFailures,
}, null, 2)}\n`);
console.log(`M0 PASS: ${records.length} sanitized events; RSS p95 ${resources.rss_kib.p95} KiB`);
} finally {
if (reconnectClient) await reconnectClient.close().catch(() => {});
if (browserClient) await browserClient.close().catch(() => {});
if (browser.exitCode === null) browser.kill("SIGTERM");
if (browser.exitCode === null) await new Promise((resolve) => browser.once("exit", resolve));
await fixtures.close();
await rm(tempDir, { recursive: true, force: true });
}
}
async function runSpike() {
const failures = [];
const resourceSamples = [];
const startup = {
requested_samples: Number(process.env.M0_STARTUP_SAMPLES ?? 3),
successful_samples: 0,
max_renderer_failure_rate: Number(process.env.M0_MAX_RENDERER_FAILURE_RATE ?? 0.34),
};
assert(Number.isSafeInteger(startup.requested_samples) && startup.requested_samples > 1, "M0_STARTUP_SAMPLES must be an integer greater than one");
assert(Number.isFinite(startup.max_renderer_failure_rate) && startup.max_renderer_failure_rate >= 0 && startup.max_renderer_failure_rate <= 1, "M0_MAX_RENDERER_FAILURE_RATE must be between zero and one");
for (let sample = 1; sample <= startup.requested_samples; sample++) {
for (let attempt = 1; attempt <= 3; attempt++) {
startup.successful_samples = sample;
try {
await main(failures, startup, resourceSamples);
break;
} catch (error) {
const original = error.cause ?? error;
failures.push({ sample, attempt, error: original.message, detached_reason: error.detachedReason ?? null });
if (!(error instanceof StartupNavigateFailure) || attempt === 3) throw original;
console.error(`M0 startup sample ${sample}/${startup.requested_samples} retry ${attempt}/3: ${original.message}`);
}
}
}
}
if (process.argv[1] && import.meta.url === pathToFileURL(process.argv[1]).href) {
const run = process.argv[2] === "--fixture-child" ? fixtureChild : runSpike;
run().catch((error) => {
console.error(`M0 FAIL: ${error.message}`);
process.exitCode = 1;
});
}