H-385: validate Clark CDP network event loop (#1)
This commit was merged in pull request #1.
This commit is contained in:
@@ -0,0 +1,698 @@
|
||||
#!/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) {
|
||||
const bytes = Buffer.from(body, base64Encoded ? "base64" : "utf8");
|
||||
if (bytes.length > maxBytes) {
|
||||
return { state: "size_limit", storage: "omitted", bytes: bytes.length };
|
||||
}
|
||||
return {
|
||||
state: "fetched",
|
||||
storage: "omitted",
|
||||
reason: "default_body_policy",
|
||||
bytes: bytes.length,
|
||||
};
|
||||
}
|
||||
|
||||
export function shouldRetryStartup(detachedReason) {
|
||||
return detachedReason === "Render process gone.";
|
||||
}
|
||||
|
||||
class StartupNavigateFailure extends Error {
|
||||
constructor(cause) {
|
||||
super(cause.message, { cause });
|
||||
}
|
||||
}
|
||||
|
||||
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() {
|
||||
this.ws = new WebSocket(this.url);
|
||||
await new Promise((resolve, reject) => {
|
||||
this.ws.onopen = resolve;
|
||||
this.ws.onerror = () => 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 < 126, "fixture frame must stay small");
|
||||
return Buffer.concat([Buffer.from([0x80 | opcode, body.length]), 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 === "/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("https://127.0.0.1:${httpsPort}/secure")).text();
|
||||
await new Promise((resolve, reject) => {
|
||||
const ws = new WebSocket("ws://127.0.0.1:${httpPort}/ws");
|
||||
ws.onopen = () => ws.send("client-message");
|
||||
ws.onmessage = () => 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)),
|
||||
};
|
||||
}
|
||||
|
||||
function isRunning(pid) {
|
||||
try {
|
||||
process.kill(pid, 0);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
function eventCollector(sessionId, targetId) {
|
||||
const records = [];
|
||||
const requestState = new Map();
|
||||
const push = (kind, fields = {}) => records.push({
|
||||
schema_version: "1.0",
|
||||
session_id: sessionId,
|
||||
target_id: targetId,
|
||||
kind,
|
||||
...fields,
|
||||
});
|
||||
|
||||
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");
|
||||
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.requestWillBeSentExtraInfo":
|
||||
push("http_request_headers", { request_id: requestId, headers: redactHeaders(params.headers) });
|
||||
break;
|
||||
case "Network.responseReceived":
|
||||
mark("response");
|
||||
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),
|
||||
},
|
||||
});
|
||||
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}`);
|
||||
push("websocket_frame", {
|
||||
connection_id: requestId,
|
||||
direction,
|
||||
opcode: params.response.opcode,
|
||||
payload_bytes: Buffer.byteLength(params.response.payloadData),
|
||||
payload: { state: "omitted", reason: "frame_policy" },
|
||||
});
|
||||
break;
|
||||
}
|
||||
case "Network.webSocketClosed":
|
||||
mark("ws_closed");
|
||||
push("websocket_closed", { connection_id: requestId });
|
||||
break;
|
||||
}
|
||||
};
|
||||
return { records, requestState, 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;
|
||||
}
|
||||
|
||||
async function collectBody(cdp, collector, requestId, maxBodyBytes) {
|
||||
try {
|
||||
const result = await cdp.call("Network.getResponseBody", { requestId });
|
||||
collector.push("http_body", { request_id: requestId, body: boundedBody(result.body, result.base64Encoded, maxBodyBytes) });
|
||||
} catch (error) {
|
||||
collector.push("http_body", {
|
||||
request_id: requestId,
|
||||
body: { state: "unavailable", storage: "omitted", reason: "cdp_error", error: error.message },
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function assertCoverage(collector, secrets) {
|
||||
const { records, requestState } = collector;
|
||||
for (const pathname of ["/ok", "/secret", "/large", "/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", 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 === "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");
|
||||
}
|
||||
|
||||
async function main(attemptFailures = []) {
|
||||
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 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);
|
||||
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);
|
||||
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", "/secure"]) {
|
||||
await collectBody(cdp, collector, findRequest(collector.records, pathname), maxBodyBytes);
|
||||
}
|
||||
const failedId = [...collector.requestState.entries()].find(([, state]) => state.has("failed"))?.[0];
|
||||
await collectBody(cdp, collector, failedId, maxBodyBytes);
|
||||
|
||||
const resources = processSnapshot(browserPid);
|
||||
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 { result: reconnectResult } = await reconnectClient.call("Runtime.evaluate", { expression: "1 + 1", returnByValue: true });
|
||||
assert.equal(reconnectResult.value, 2, "CDP reconnect can execute a command");
|
||||
collector.push("cdp_reconnected", { target_id: reconnectTarget.id, verified: true });
|
||||
assert(isRunning(browserPid), "disconnecting and reconnecting CDP does not exit clark-browser");
|
||||
|
||||
assertCoverage(collector, fixtures.secrets);
|
||||
await mkdir(outputDir, { recursive: true });
|
||||
await writeFile(path.join(outputDir, "m0-events.sample.jsonl"), `${collector.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",
|
||||
target_close_survives: "pass",
|
||||
cdp_disconnect_reconnect: "pass",
|
||||
sensitive_header_redaction: "pass",
|
||||
sensitive_body_redaction: "pass",
|
||||
},
|
||||
resources,
|
||||
event_count: collector.records.length,
|
||||
max_body_bytes: maxBodyBytes,
|
||||
attempt_failures: attemptFailures,
|
||||
}, null, 2)}\n`);
|
||||
console.log(`M0 PASS: ${collector.records.length} sanitized events; RSS snapshot ${resources.rss_kib} 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 = [];
|
||||
for (let attempt = 1; attempt <= 3; attempt++) {
|
||||
try {
|
||||
await main(failures);
|
||||
return;
|
||||
} catch (error) {
|
||||
const original = error.cause ?? error;
|
||||
failures.push(original.message);
|
||||
if (!(error instanceof StartupNavigateFailure) || attempt === 3) throw original;
|
||||
console.error(`M0 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;
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user