H-385: address CDP spike review blockers

Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
2026-08-20 22:53:26 +08:00
co-authored by multica-agent
parent 8095392f7f
commit 1330585ccc
5 changed files with 150 additions and 98 deletions
+84 -45
View File
@@ -48,23 +48,35 @@ export function safeUrl(raw) {
export function boundedBody(body, base64Encoded, maxBytes) {
const bytes = Buffer.from(body, base64Encoded ? "base64" : "utf8");
if (bytes.length > maxBytes) {
return { state: "omitted", reason: "size_limit", bytes: bytes.length };
return { state: "size_limit", storage: "omitted", bytes: bytes.length };
}
return {
state: "stored",
state: "fetched",
storage: "omitted",
reason: "default_body_policy",
bytes: bytes.length,
encoding: "utf8",
text: bytes.toString("utf8"),
};
}
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() {
@@ -74,10 +86,12 @@ class CDP {
this.ws.onerror = () => reject(new Error("CDP WebSocket connection failed"));
});
this.ws.onmessage = ({ data }) => this.onMessage(JSON.parse(data));
this.ws.onclose = () => {
for (const { reject } of this.pending.values()) reject(new Error("CDP WebSocket closed"));
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();
this.resolveClosed();
for (const listener of this.closeListeners) listener(details);
this.resolveClosed(details);
};
return this;
}
@@ -86,6 +100,17 @@ class CDP {
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",
@@ -102,6 +127,8 @@ class CDP {
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();
@@ -109,13 +136,13 @@ class CDP {
for (const listener of this.listeners) listener(message.method, message.params ?? {});
}
call(method, 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}`));
}, 20_000);
}, timeoutMs);
this.pending.set(id, {
resolve: (value) => { clearTimeout(timer); resolve(value); },
reject: (error) => { clearTimeout(timer); reject(error); },
@@ -192,6 +219,7 @@ function acceptWebsocket(request, socket) {
async function startFixtures(tempDir) {
const secrets = {
authorization: randomUUID(),
body: randomUUID(),
cookie: randomUUID(),
setCookie: randomUUID(),
};
@@ -248,6 +276,10 @@ async function startFixtures(tempDir) {
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) => {
@@ -261,6 +293,7 @@ async function startFixtures(tempDir) {
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) => {
@@ -367,15 +400,6 @@ function processSnapshot(rootPid) {
};
}
function browserPidForPort(port) {
const result = spawnSync("ps", ["-eo", "pid=,args="], { encoding: "utf8" });
assert.equal(result.status, 0, "ps must locate the browser process");
const marker = `--remote-debugging-port=${port}`;
const line = result.stdout.split("\n").find((row) => row.includes(marker) && !row.includes("scripts/m0-spike.mjs"));
assert(line, `browser process with ${marker} must exist`);
return Number(line.trim().split(/\s+/, 1)[0]);
}
function isRunning(pid) {
try {
process.kill(pid, 0);
@@ -486,14 +510,14 @@ async function collectBody(cdp, collector, requestId, maxBodyBytes) {
} catch (error) {
collector.push("http_body", {
request_id: requestId,
body: { state: "unavailable", reason: "cdp_error", error: error.message },
body: { state: "unavailable", storage: "omitted", reason: "cdp_error", error: error.message },
});
}
}
function assertCoverage(collector, secrets) {
const { records, requestState } = collector;
for (const pathname of ["/ok", "/large", "/secure"]) {
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"]);
@@ -509,8 +533,12 @@ function assertCoverage(collector, secrets) {
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");
assert(records.some((record) => record.kind === "http_body" && record.body.reason === "size_limit"), "oversize body is omitted");
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 = []) {
@@ -525,7 +553,7 @@ async function main(attemptFailures = []) {
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");
spawn("setsid", ["--fork", binary,
const browser = spawn(binary, [
"--headless=new",
"--no-sandbox",
"--use-mock-keychain",
@@ -537,32 +565,40 @@ async function main(attemptFailures = []) {
"--disable-gpu",
"--disable-features=WebGPU",
"about:blank",
], { stdio: ["ignore", "ignore", DEBUG ? "inherit" : "ignore"] });
], { stdio: ["ignore", "ignore", "pipe"] });
browser.stderr.pipe(process.stderr);
const sessionId = `m0-${randomUUID()}`;
let browserPid;
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,
);
browserPid = browserPidForPort(cdpPort);
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");
await new Promise((resolve) => setTimeout(resolve, 500));
await cdp.call("Page.navigate", { url: fixtures.url });
await waitFor(async () => {
const { result } = await cdp.call("Runtime.evaluate", { expression: "document.title", returnByValue: true });
return result.value === "M0_READY";
}, "fixture page readiness");
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: {
@@ -577,7 +613,7 @@ async function main(attemptFailures = []) {
}, "fixture completion");
await waitFor(() => [...collector.requestState.values()].some((state) => state.has("ws_closed")), "WebSocket closure");
for (const pathname of ["/ok", "/large", "/secure"]) {
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];
@@ -590,11 +626,13 @@ async function main(attemptFailures = []) {
collector.push("target_closed", { reason: "Target.closeTarget" });
assert(isRunning(browserPid), "closing the page does not exit clark-browser");
const disconnectTarget = await getJson(`http://127.0.0.1:${cdpPort}/json/new?about:blank`, { method: "PUT" });
const disconnectClient = await new CDP(disconnectTarget.webSocketDebuggerUrl).connect();
await disconnectClient.close();
collector.push("cdp_disconnected", { reason: "client_closed" });
assert(isRunning(browserPid), "disconnecting CDP 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 });
@@ -616,8 +654,9 @@ async function main(attemptFailures = []) {
response_body_size_limit: "pass",
response_body_unavailable: "pass",
target_close_survives: "pass",
cdp_disconnect_survives: "pass",
cdp_disconnect_reconnect: "pass",
sensitive_header_redaction: "pass",
sensitive_body_redaction: "pass",
},
resources,
event_count: collector.records.length,
@@ -626,9 +665,10 @@ async function main(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 (browserPid && isRunning(browserPid)) process.kill(browserPid, "SIGTERM");
if (browserPid) await waitFor(() => !isRunning(browserPid), "browser shutdown").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 });
}
@@ -641,11 +681,10 @@ async function runSpike() {
await main(failures);
return;
} catch (error) {
failures.push(error.message);
const retryable = error.message.startsWith("Inspector.detached:") ||
error.message === "CDP command timed out: Page.navigate";
if (!retryable || attempt === 3) throw error;
console.error(`M0 retry ${attempt}/3: ${error.message}`);
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}`);
}
}
}