repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/chrome-cdp/browser.ts
Rawimport { type ChildProcess, execFileSync, spawn } from "node:child_process";
import { randomUUID } from "node:crypto";
import { existsSync } from "node:fs";
import {
appendFile,
lstat,
mkdir,
mkdtemp,
open,
readdir,
readFile,
rm,
writeFile,
} from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import {
artifactPath,
binaryKey,
dataRoot,
downloads,
type Frame,
findExecutable,
type Json,
profileName,
ref,
TMP_PREFIX,
validateFlags,
} from "./lib.ts";
import type { Debug, Span } from "./src/debug.ts";
const COMMAND_MS = 30_000,
INLINE = 64 * 1024,
SPOOL_LIMIT = 16 * 1024 * 1024;
const EVENTS = 200,
VIDEO_MAX = 5 * 60_000,
VIDEO_BYTES = 256 * 1024 * 1024;
const traceCategories =
"devtools.timeline,blink.user_timing,loading,disabled-by-default-devtools.timeline";
export type Operation = {
kind:
| "cdp"
| "wait"
| "status"
| "close"
| "screenshot"
| "videoStart"
| "videoStop"
| "traceStart"
| "traceStop";
method?: string;
params?: Json;
sessionId?: string;
label?: string;
after?: string;
targetId?: string;
event?: string;
timeoutMs?: number;
filename?: string;
fullPage?: boolean;
};
export type Options = {
headed?: boolean;
executablePath?: string;
profile?: string;
flags?: string[];
artifactDir?: string;
diagnostic?: Debug;
operationSpan?: Span;
};
type Event = {
seq: number;
time: string;
sessionId?: string;
method: string;
frame?: Record<string, Json>;
spooled?: string;
bytes?: number;
lost?: boolean;
};
type Watcher = {
method: string;
sessionId: string;
since: number;
hit?: Event;
};
type Video = {
process: ChildProcess;
sessionId: string;
file: string;
size: number;
droppedFrames: number;
timer: NodeJS.Timeout;
sampleTimer?: NodeJS.Timeout;
lastFrame?: Buffer;
samples?: number;
done: Promise<number | null>;
};
type Trace = {
file: string;
done: Promise<Record<string, Json>>;
finish: (v: Record<string, Json>) => void;
timer: NodeJS.Timeout;
};
function alive(pid: number): boolean {
if (!Number.isSafeInteger(pid) || pid <= 0) return false;
try {
process.kill(pid, 0);
return true;
} catch (e) {
return (e as NodeJS.ErrnoException).code === "EPERM";
}
}
async function markOwned(dir: string, browserPid = 0): Promise<void> {
await writeFile(
join(dir, ".pi-chrome-cdp-owner.json"),
JSON.stringify({
tag: "pi-chrome-cdp-v1",
ownerPid: process.pid,
browserPid,
}),
{ mode: 0o600 },
);
}
async function cleanupStale(): Promise<void> {
for (const name of await readdir(tmpdir())) {
if (!/^pi-chrome-cdp-[A-Za-z0-9_-]+$/.test(name)) continue;
const path = join(tmpdir(), name);
try {
if (!(await lstat(path)).isDirectory()) continue;
const marker = JSON.parse(
await readFile(join(path, ".pi-chrome-cdp-owner.json"), "utf8"),
) as { tag?: string; ownerPid?: number; browserPid?: number };
if (
marker.tag !== "pi-chrome-cdp-v1" ||
alive(marker.ownerPid || 0) ||
alive(marker.browserPid || 0)
)
continue;
await rm(path, {
recursive: true,
force: true,
maxRetries: 5,
retryDelay: 100,
});
} catch {
/* unknown ownership or active filesystem use: leave untouched */
}
}
}
async function saveNew(file: string, bytes: Buffer | string): Promise<void> {
const f = await open(file, "wx", 0o600);
try {
await f.writeFile(bytes);
} finally {
await f.close();
}
}
async function waitExit(child: ChildProcess, ms = 3000): Promise<boolean> {
if (child.exitCode !== null || child.signalCode !== null) return true;
return new Promise((resolve) => {
const t = setTimeout(() => {
child.off("exit", onExit);
resolve(false);
}, ms);
const onExit = () => {
clearTimeout(t);
resolve(true);
};
child.once("exit", onExit);
});
}
export class Browser {
child!: ChildProcess;
socket!: WebSocket;
port = 0;
endpoint = "";
profileDir = "";
spoolDir = "";
artifactDir = "";
persistent = false;
ownedLock = false;
lockToken = randomUUID();
headed = false;
binary = "";
version = "";
viewport = "unknown";
observedTargets: Json[] = [];
generation = randomUUID();
initialTarget = "";
initialSession = "";
nextId = 1;
seq = 0;
cursor = 0;
uncertain = false;
monitor = new Map<string, boolean>();
autoAttach = true;
sessions = new Map<string, string>();
results = new Map<string, Frame>();
pending = new Map<number, (f: Frame) => void>();
events: Event[] = [];
watchers: Watcher[] = [];
lateResponses: Frame[] = [];
gaps = 0;
notices: string[] = [];
spoolFiles: {
file: string;
bytes: number;
firstEventSeq?: number;
lastEventSeq?: number;
}[] = [];
spoolBytes = 0;
eventChunk = 0;
spoolWrites: Promise<void> = Promise.resolve();
queuedSpoolBytes = 0;
private protectedSpoolFiles = new Set<string>();
private batchSpooling = false;
video?: Video;
videoStopping?: Promise<string>;
trace?: Trace;
traceStopping?: Promise<string>;
busy: Promise<void> = Promise.resolve();
closed = false;
diagnostic?: Debug;
operationSpan?: Span;
private eventCount = 0;
private flushEventCounts(): void {
if (this.eventCount)
this.diagnostic?.("cdp.events", { count: this.eventCount });
this.eventCount = 0;
}
static async launch(opts: Options): Promise<Browser> {
await cleanupStale();
const b = new Browser();
b.diagnostic = opts.diagnostic;
b.operationSpan = opts.operationSpan;
b.diagnostic?.("browser.launch.start");
b.binary = findExecutable(opts.executablePath);
b.headed = opts.headed ?? false;
validateFlags(opts.flags || []);
b.artifactDir = opts.artifactDir || downloads();
b.persistent = !!opts.profile;
try {
if (opts.profile) {
const name = profileName(opts.profile);
b.profileDir = join(dataRoot(), binaryKey(b.binary), name);
await mkdir(b.profileDir, { recursive: true, mode: 0o700 });
const marker = join(b.profileDir, ".pi-chrome-cdp-version");
const lock = join(b.profileDir, ".pi-chrome-cdp-lock");
if (existsSync(lock))
throw new Error(`profile locked: ${name}; close other owner first`);
const version = await Browser.binaryMajor(b.binary);
const entries = await readdir(b.profileDir);
if (!existsSync(marker) && entries.length)
throw new Error(
"profile version marker absent; refusing existing profile",
);
await saveNew(
lock,
JSON.stringify({ pid: process.pid, token: b.lockToken }),
);
b.ownedLock = true;
if (existsSync(marker)) {
const previous = (await readFile(marker, "utf8")).trim();
if (!/^\d+\.\d+\.\d+\.\d+$/.test(previous))
throw new Error("corrupt profile version marker");
const a = previous.split(".").map(Number),
bver = version.split(".").map(Number);
if (
a.some(
(n, i) =>
n > bver[i] && a.slice(0, i).every((v, j) => v === bver[j]),
)
)
throw new Error(
`profile downgrade refused: ${previous} → ${version}`,
);
}
const staged = `${marker}.${randomUUID()}`;
await writeFile(staged, version, { flag: "wx", mode: 0o600 });
await import("node:fs/promises").then((fs) =>
fs.rename(staged, marker),
);
} else {
b.profileDir = await mkdtemp(TMP_PREFIX);
await markOwned(b.profileDir);
}
b.spoolDir = await mkdtemp(TMP_PREFIX);
await markOwned(b.spoolDir);
const args = [
"--remote-debugging-port=0",
"--remote-debugging-address=127.0.0.1",
`--user-data-dir=${b.profileDir}`,
"--no-first-run",
"--window-size=1280,720",
...(b.headed ? [] : ["--headless=new"]),
...(opts.flags || []),
"about:blank",
];
const active = join(b.profileDir, "DevToolsActivePort");
await rm(active, { force: true });
b.child = spawn(b.binary, args, { stdio: "ignore" });
let spawnFailure: Error | undefined;
b.child.on("error", (e) => {
spawnFailure = e;
});
if (!b.persistent) await markOwned(b.profileDir, b.child.pid);
await markOwned(b.spoolDir, b.child.pid);
let found = false;
for (let i = 0; i < 200; i++) {
if (existsSync(active)) {
const [port] = (await readFile(active, "utf8")).split("\n");
b.port = Number(port);
if (Number.isInteger(b.port) && b.port > 0) {
found = true;
break;
}
}
if (spawnFailure) throw spawnFailure;
if (b.child.exitCode !== null || b.child.signalCode !== null)
throw new Error(
`Chrome exited during launch (${b.child.exitCode ?? b.child.signalCode})`,
);
await new Promise((r) => setTimeout(r, 100));
}
if (!found)
throw new Error("Chrome did not publish DevToolsActivePort within 20s");
const address = `http://127.0.0.1:${b.port}`;
let meta: { webSocketDebuggerUrl: string; Browser: string } | undefined;
for (let attempt = 0; attempt < 50; attempt++) {
try {
meta = (await (
await fetch(`${address}/json/version`, {
signal: AbortSignal.timeout(1000),
})
).json()) as typeof meta;
break;
} catch {
if (spawnFailure) throw spawnFailure;
if (b.child.exitCode !== null || b.child.signalCode !== null)
throw new Error("Chrome exited before HTTP endpoint became ready");
await new Promise((r) => setTimeout(r, 100));
}
}
if (!meta) throw new Error("Chrome HTTP discovery endpoint unavailable");
if (
!meta.webSocketDebuggerUrl?.startsWith(
`ws://127.0.0.1:${b.port}/devtools/browser/`,
) &&
!meta.webSocketDebuggerUrl?.startsWith(
`ws://localhost:${b.port}/devtools/browser/`,
)
)
throw new Error("unexpected Chrome debug endpoint");
b.endpoint = address;
b.version = meta.Browser;
b.socket = new WebSocket(meta.webSocketDebuggerUrl);
await new Promise<void>((resolve, reject) => {
const timer = setTimeout(
() => reject(new Error("Chrome WebSocket timeout")),
5000,
);
b.socket.onopen = () => {
clearTimeout(timer);
resolve();
};
b.socket.onerror = () => {
clearTimeout(timer);
reject(new Error("Chrome WebSocket failed"));
};
});
b.socket.onmessage = (e) => b.onMessage(String(e.data));
b.socket.onclose = () => {
for (const [id, pending] of b.pending)
pending({
id,
error: { code: -1, message: "Chrome connection closed" },
transportError: true,
});
b.pending.clear();
b.uncertain = true;
};
const dir = await mkdtemp(TMP_PREFIX);
// Browser downloads belong to a private disposable workspace.
await markOwned(dir, b.child.pid);
b.spoolFiles.push({ file: dir, bytes: 0 });
const download = await b.send("Browser.setDownloadBehavior", {
behavior: "allow",
downloadPath: dir,
eventsEnabled: true,
});
if (download.error)
throw new Error(`cannot contain downloads: ${download.error.message}`);
const auto = await b.send("Target.setAutoAttach", {
autoAttach: true,
waitForDebuggerOnStart: false,
flatten: true,
});
if (auto.error)
throw new Error(`cannot attach targets: ${auto.error.message}`);
const target = await b.send("Target.createTarget", {
url: "about:blank",
});
if (target.error) throw new Error(target.error.message);
b.initialTarget = String(
(target.result as Record<string, Json>).targetId,
);
const attach = await b.send("Target.attachToTarget", {
targetId: b.initialTarget,
flatten: true,
});
if (attach.error) throw new Error(attach.error.message);
b.initialSession = String(
(attach.result as Record<string, Json>).sessionId,
);
b.sessions.set(b.initialSession, b.initialTarget);
for (const method of ["Runtime.enable", "Log.enable", "Page.enable"]) {
const result = await b.send(method, {}, b.initialSession);
b.monitor.set(
`${b.initialSession}:${method.split(".")[0]}`,
!result.error,
);
}
await b.refreshStatus();
b.diagnostic?.("browser.launch.finish", { status: "ok" });
return b;
} catch (e) {
b.diagnostic?.("browser.launch.error");
await b.close();
throw e;
}
}
static async binaryMajor(file: string): Promise<string> {
const text = execFileSync(file, ["--version"], {
encoding: "utf8",
timeout: 3000,
});
const version = text.match(/\b(\d+\.\d+\.\d+\.\d+)\b/)?.[1];
if (!version)
throw new Error(`cannot parse Chrome version: ${text.trim()}`);
return version;
}
onMessage(text: string): void {
let frame: Record<string, Json>;
try {
frame = JSON.parse(text);
} catch {
return;
}
if (typeof frame.id === "number") {
const callback = this.pending.get(frame.id);
if (callback) {
this.pending.delete(frame.id);
callback(frame as Frame);
} else {
const size = Buffer.byteLength(text);
this.lateResponses.push(
size > 16384
? {
id: frame.id,
error: {
code: -5,
message: `late response of ${size} bytes omitted from memory`,
},
transportError: true,
}
: (frame as Frame),
);
if (size > 16384) this.gaps++;
while (this.lateResponses.length > 4) this.lateResponses.shift();
}
return;
}
const method = frame.method;
const params = frame.params as Record<string, Json> | undefined;
if (method !== "Page.screencastFrame" && this.diagnostic) {
const name = String(method);
const subtype = name.startsWith("Runtime.")
? "runtime"
: name.startsWith("Log.")
? "log"
: name.startsWith("Target.")
? "target"
: name.startsWith("Page.")
? "page"
: name.startsWith("Network.")
? "network"
: name.startsWith("Tracing.")
? "tracing"
: undefined;
if (subtype) this.diagnostic("cdp.event", { subtype });
else if (++this.eventCount >= 200) this.flushEventCounts();
}
if (
method === "Page.screencastFrame" &&
this.video &&
frame.sessionId === this.video.sessionId &&
params
) {
const data = String(params.data || "");
const bytes = Buffer.from(data, "base64");
this.video.size += bytes.length;
const video = this.video;
if (video.size > VIDEO_BYTES && !this.videoStopping)
void this.stopVideo().catch((e) => {
this.notices.push(`video finalize failed: ${String(e)}`);
});
else if (!this.videoStopping) video.lastFrame = bytes;
if (this.video === video)
void this.send(
"Page.screencastFrameAck",
{ sessionId: params.sessionId },
video.sessionId,
).catch(() => {});
return;
}
if (method === "Tracing.tracingComplete" && this.trace && params)
this.trace.finish(params);
if (method === "Target.attachedToTarget" && params) {
const sid = String(params.sessionId || "");
const info = params.targetInfo as Record<string, Json> | undefined;
if (sid && info?.targetId) {
this.sessions.set(sid, String(info.targetId));
if (
info.type === "page" ||
info.type === "worker" ||
info.type === "service_worker" ||
info.type === "iframe"
) {
void this.send(
"Target.setAutoAttach",
{ autoAttach: true, waitForDebuggerOnStart: false, flatten: true },
sid,
).catch(() => {});
for (const m of [
...(info.type === "page" ? ["Page.enable"] : []),
"Runtime.enable",
"Log.enable",
])
void this.send(m, {}, sid)
.then((r) =>
this.monitor.set(`${sid}:${m.split(".")[0]}`, !r.error),
)
.catch(() => {});
}
}
}
if (method === "Target.detachedFromTarget" && params)
this.sessions.delete(String(params.sessionId || ""));
if (this.events.length >= EVENTS) {
const removed = this.events.shift();
if (removed) this.queueEventSpool(removed);
}
const raw = JSON.stringify(frame);
const size = Buffer.byteLength(raw);
const event: Event = {
seq: ++this.seq,
time: new Date().toISOString(),
sessionId: frame.sessionId as string | undefined,
method: String(method),
...(size > 16384 ? { bytes: size } : { frame }),
};
if (size > 16384) this.queueEventSpool(event, frame);
for (const watcher of this.watchers)
if (
!watcher.hit &&
event.seq > watcher.since &&
event.method === watcher.method &&
event.sessionId === watcher.sessionId
)
watcher.hit = event;
this.events.push(event);
}
private queueSpool<T>(write: () => Promise<T>): Promise<T> {
const task = this.spoolWrites.then(write);
this.spoolWrites = task.then(
() => {},
() => {},
);
return task;
}
queueEventSpool(e: Event, frame?: Record<string, Json>): void {
const item = frame ? { ...e, frame } : e;
const bytes = Buffer.byteLength(JSON.stringify(item));
if (this.queuedSpoolBytes + bytes > SPOOL_LIMIT) {
e.lost = true;
this.gaps++;
return;
}
this.queuedSpoolBytes += bytes;
void this.queueSpool(async () => {
e.spooled = await this.writeSpoolEvent(item);
})
.catch(() => {
e.lost = true;
this.gaps++;
})
.finally(() => {
this.queuedSpoolBytes -= bytes;
});
}
beginBatchSpool(startSeq: number): void {
this.protectedSpoolFiles.clear();
this.batchSpooling = true;
for (const event of this.events)
if (event.seq > startSeq && event.spooled && !event.lost)
this.protectSpoolFile(event.spooled);
}
private async makeSpoolRoom(bytes: number): Promise<boolean> {
while (this.spoolBytes + bytes > SPOOL_LIMIT) {
const oldest = this.spoolFiles.find(
(f) => f.bytes > 0 && !this.protectedSpoolFiles.has(f.file),
);
if (!oldest) return false;
this.spoolFiles.splice(this.spoolFiles.indexOf(oldest), 1);
await rm(oldest.file, { force: true });
if (oldest.firstEventSeq !== undefined) {
this.eventChunk++;
for (const event of this.events)
if (event.spooled === oldest.file) {
delete event.spooled;
event.lost = true;
}
}
this.spoolBytes -= oldest.bytes;
this.gaps++;
}
return true;
}
protectSpoolFile(file: string): void {
if (this.batchSpooling) this.protectedSpoolFiles.add(file);
}
private async writeSpoolEvent(e: Event): Promise<string> {
const bytes = JSON.stringify(e) + "\n";
const size = Buffer.byteLength(bytes);
if (!(await this.makeSpoolRoom(size))) {
e.lost = true;
this.gaps++;
return "";
}
const path = join(
this.spoolDir,
`events-${this.generation}-${this.eventChunk}.jsonl`,
);
await appendFile(path, bytes, { mode: 0o600 });
let entry = this.spoolFiles.find((f) => f.file === path);
if (!entry) {
entry = {
file: path,
bytes: 0,
firstEventSeq: e.seq,
lastEventSeq: e.seq,
};
this.spoolFiles.push(entry);
}
entry.bytes += size;
entry.firstEventSeq ??= e.seq;
entry.lastEventSeq = e.seq;
this.spoolBytes += size;
this.protectSpoolFile(path);
return path;
}
spoolEvent(e: Event): Promise<string> {
return this.queueSpool(() => this.writeSpoolEvent(e));
}
send(
method: string,
params: Json = {},
sessionId?: string,
signal?: AbortSignal,
timeoutMs = COMMAND_MS,
): Promise<Frame> {
if (!this.socket || this.socket.readyState !== WebSocket.OPEN)
return Promise.reject(new Error("Chrome socket not open"));
const id = this.nextId++;
const end =
method === "Page.screencastFrameAck"
? undefined
: this.operationSpan?.("cdp.request", { timeoutMs });
return new Promise((resolve) => {
let done = false;
const finish = (frame: Frame) => {
if (done) return;
done = true;
clearTimeout(timer);
signal?.removeEventListener("abort", abort);
this.pending.delete(id);
end?.(frame.transportError || frame.error ? "error" : "finish", {
status: frame.transportError
? "transport"
: frame.error
? "cdp-error"
: "ok",
});
resolve(frame);
};
const abort = () => {
this.uncertain = true;
finish({
id,
error: { code: -2, message: "canceled; effects may have applied" },
transportError: true,
});
};
const timer = setTimeout(() => {
this.uncertain = true;
finish({
id,
error: {
code: -3,
message: `timeout after ${timeoutMs}ms; effects may have applied`,
},
transportError: true,
});
}, timeoutMs);
if (signal?.aborted) {
abort();
return;
}
signal?.addEventListener("abort", abort, { once: true });
this.pending.set(id, finish);
try {
this.socket.send(
JSON.stringify({
id,
method,
params,
...(sessionId ? { sessionId } : {}),
}),
);
} catch (e) {
finish({
id,
error: { code: -4, message: String(e) },
transportError: true,
});
}
});
}
async refreshStatus(): Promise<void> {
const targets = await this.send("Target.getTargets");
if (targets.error)
throw new Error(`status unavailable: ${targets.error.message}`);
this.observedTargets =
((targets.result as Record<string, Json>).targetInfos as Json[]) || [];
const metrics = await this.send(
"Page.getLayoutMetrics",
{},
this.initialSession,
);
if (!metrics.error) {
const viewport = (metrics.result as Record<string, Json>)
.cssVisualViewport as Record<string, Json> | undefined;
if (viewport)
this.viewport = `${viewport.clientWidth}x${viewport.clientHeight}`;
}
this.uncertain = false;
}
status() {
return {
generation: this.generation,
pid: this.child?.pid,
port: this.port,
endpoint: this.endpoint,
protocol: `${this.endpoint}/json/protocol`,
version: this.version,
binary: this.binary,
headed: this.headed,
viewport: this.viewport,
profile: this.persistent ? "named" : "temporary",
initialTarget: this.initialTarget,
initialSession: this.initialSession,
targets: this.observedTargets,
sessions: Object.fromEntries(this.sessions),
monitor: Object.fromEntries(this.monitor),
autoAttach: this.autoAttach,
uncertain: this.uncertain,
notices: this.notices.slice(-8),
lateResponses: this.lateResponses,
video: this.video?.file,
droppedVideoFrames: this.video?.droppedFrames,
trace: this.trace?.file,
gaps: this.gaps,
};
}
async outputFile(filename: string | undefined, ext: string): Promise<string> {
await mkdir(this.artifactDir, { recursive: true, mode: 0o700 });
const name =
filename ||
`chrome-cdp-${new Date().toISOString().replace(/[:.]/g, "-")}-${randomUUID()}.${ext}`;
const path = artifactPath(this.artifactDir, name);
if (existsSync(path)) throw new Error(`artifact already exists: ${path}`);
return path;
}
async screenshot(
sid: string,
filename?: string,
fullPage = false,
): Promise<{ file: string; data: string }> {
const response = await this.send(
"Page.captureScreenshot",
{ format: "png", captureBeyondViewport: fullPage, fromSurface: true },
sid,
);
if (response.error) throw new Error(response.error.message);
const data = String((response.result as Record<string, Json>).data || "");
const file = await this.outputFile(filename, "png");
await saveNew(file, Buffer.from(data, "base64"));
return { file, data };
}
async startVideo(sid: string, filename?: string): Promise<string> {
if (this.video) throw new Error("video already running");
const file = await this.outputFile(filename, "webm");
if (!file.endsWith(".webm"))
throw new Error("video filename must end in .webm");
const ffmpeg = spawn(
"ffmpeg",
[
"-loglevel",
"error",
"-f",
"image2pipe",
"-framerate",
"8",
"-vcodec",
"mjpeg",
"-i",
"pipe:0",
"-an",
"-c:v",
"libvpx-vp9",
"-n",
file,
],
{ stdio: ["pipe", "ignore", "pipe"] },
);
await new Promise<void>((resolve, reject) => {
ffmpeg.once("spawn", () => resolve());
ffmpeg.once("error", reject);
});
let error = "";
ffmpeg.stderr?.on("data", (data) => {
error = (error + String(data)).slice(-2048);
});
ffmpeg.stdin?.on("error", () => {
/* encoder exit is reported by stopVideo */
});
const done = new Promise<number | null>((resolve) =>
ffmpeg.once("exit", (code) => resolve(code)),
);
const timer = setTimeout(
() =>
void this.stopVideo().catch((e) => {
this.notices.push(`video finalize failed: ${String(e)}`);
}),
VIDEO_MAX,
);
const video: Video = {
file,
process: ffmpeg,
sessionId: sid,
size: 0,
droppedFrames: 0,
done,
timer,
};
this.video = video;
const r = await this.send(
"Page.startScreencast",
{ format: "jpeg", quality: 75, everyNthFrame: 1 },
sid,
);
if (r.error) {
await this.stopVideo();
throw new Error(r.error.message);
}
// ffmpeg's image pipe has no timestamps; sample at its declared 8 fps rather than
// treating every incoming (variable-rate) CDP frame as 125 ms of playback.
video.sampleTimer = setInterval(() => {
if (!video.lastFrame) return;
if (
video.process.stdin?.writable &&
video.process.stdin.writableLength < 2 * 1024 * 1024
) {
video.process.stdin.write(video.lastFrame);
video.samples = (video.samples || 0) + 1;
} else video.droppedFrames++;
}, 125);
return file;
}
async stopVideo(issueStop = true): Promise<string> {
if (this.videoStopping) return this.videoStopping;
const v = this.video;
if (!v) throw new Error("no managed video running");
this.videoStopping = (async () => {
clearTimeout(v.timer);
if (v.sampleTimer) clearInterval(v.sampleTimer);
if (issueStop && this.socket?.readyState === WebSocket.OPEN)
await this.send("Page.stopScreencast", {}, v.sessionId).catch(() => {});
v.process.stdin?.end();
let guard: NodeJS.Timeout | undefined;
const code = await Promise.race([
v.done,
new Promise<null>((resolve) => {
guard = setTimeout(() => resolve(null), 5000);
}),
]);
if (guard) clearTimeout(guard);
if (code === null) {
v.process.kill("SIGKILL");
await v.done;
}
if (this.video === v) this.video = undefined;
if (code !== 0 || !existsSync(v.file))
throw new Error(`ffmpeg failed (${code ?? "timeout"}); ${v.file}`);
return v.file;
})();
try {
return await this.videoStopping;
} finally {
this.videoStopping = undefined;
}
}
async startTrace(filename?: string): Promise<string> {
if (this.trace) throw new Error("trace already running");
const file = await this.outputFile(filename, "json");
let finish!: (v: Record<string, Json>) => void;
const done = new Promise<Record<string, Json>>((resolve) => {
finish = resolve;
});
const timer = setTimeout(
() =>
void this.stopTrace().catch((e) => {
this.notices.push(`trace finalize failed: ${String(e)}`);
}),
VIDEO_MAX,
);
this.trace = { file, done, finish, timer };
const result = await this.send("Tracing.start", {
categories: traceCategories,
transferMode: "ReturnAsStream",
});
if (result.error) {
clearTimeout(timer);
this.trace = undefined;
throw new Error(result.error.message);
}
return file;
}
async stopTrace(issueEnd = true): Promise<string> {
if (this.traceStopping) return this.traceStopping;
const trace = this.trace;
if (!trace) throw new Error("no managed trace running");
this.traceStopping = (async () => {
clearTimeout(trace.timer);
try {
if (issueEnd) {
const end = await this.send("Tracing.end");
if (end.error) throw new Error(end.error.message);
}
let guard: NodeJS.Timeout | undefined;
const completion = await Promise.race([
trace.done,
new Promise<null>((resolve) => {
guard = setTimeout(() => resolve(null), 15000);
}),
]);
if (guard) clearTimeout(guard);
if (!completion?.stream) throw new Error("trace stream not returned");
const stream = String(completion.stream);
let out: Awaited<ReturnType<typeof open>> | undefined;
let finished = false;
let bytes = 0;
try {
out = await open(trace.file, "wx", 0o600);
while (true) {
const r = await this.send("IO.read", {
handle: stream,
size: 1024 * 1024,
});
if (r.error) throw new Error(r.error.message);
const data = r.result as Record<string, Json>;
const chunk = data.base64Encoded
? Buffer.from(String(data.data), "base64")
: Buffer.from(String(data.data || ""));
bytes += chunk.length;
if (bytes > VIDEO_BYTES)
throw new Error("trace size limit reached");
await out.writeFile(chunk);
if (data.eof) break;
}
finished = true;
} finally {
try {
await out?.close();
} finally {
await this.send("IO.close", { handle: stream }).catch(() => {});
if (out && !finished) await rm(trace.file, { force: true });
}
}
return trace.file;
} finally {
if (this.trace === trace) this.trace = undefined;
}
})();
try {
return await this.traceStopping;
} finally {
this.traceStopping = undefined;
}
}
spoolPayload(text: string): Promise<string> {
return this.queueSpool(async () => {
const bytes = Buffer.byteLength(text);
if (bytes > SPOOL_LIMIT)
throw new Error("tool result exceeds spool quota");
if (!(await this.makeSpoolRoom(bytes)))
throw new Error(
"tool result exceeds available spool quota; current batch files are retained",
);
const file = join(this.spoolDir, `batch-${randomUUID()}.json`);
await saveNew(file, text);
this.spoolFiles.push({ file, bytes });
this.spoolBytes += bytes;
this.protectSpoolFile(file);
return file;
});
}
spoolFrame(
frame: Frame,
text = JSON.stringify(frame),
): Promise<{ path: string; bytes: number; preview: string } | Frame> {
return this.queueSpool(async () => {
const bytes = Buffer.byteLength(text);
if (bytes <= INLINE) return frame;
if (bytes > SPOOL_LIMIT)
throw new Error(`response ${bytes} bytes exceeds spool quota`);
if (!(await this.makeSpoolRoom(bytes)))
throw new Error(
`response ${bytes} bytes exceeds available spool quota; current batch files are retained`,
);
const file = join(
this.spoolDir,
`response-${frame.id}-${randomUUID()}.json`,
);
await saveNew(file, text);
this.spoolFiles.push({ file, bytes });
this.spoolBytes += bytes;
this.protectSpoolFile(file);
return { path: file, bytes, preview: text.slice(0, 512) };
});
}
async close(): Promise<void> {
if (this.closed) return;
this.flushEventCounts();
this.diagnostic?.("browser.close.start");
const failed = (type: string) => (e: unknown) => {
this.notices.push(`${type} finalize failed: ${String(e)}`);
};
if (this.videoStopping) await this.videoStopping.catch(failed("video"));
else if (this.video) await this.stopVideo().catch(failed("video"));
if (this.traceStopping) await this.traceStopping.catch(failed("trace"));
else if (this.trace) await this.stopTrace().catch(failed("trace"));
if (
this.child &&
this.child.exitCode === null &&
this.child.signalCode === null &&
this.socket?.readyState === WebSocket.OPEN
) {
try {
this.socket.send(
JSON.stringify({
id: this.nextId++,
method: "Browser.close",
params: {},
}),
);
} catch {
/* fall back to SIGTERM */
}
await waitExit(this.child);
}
this.socket?.close();
if (
this.child &&
this.child.exitCode === null &&
this.child.signalCode === null
) {
this.child.kill("SIGTERM");
if (!(await waitExit(this.child))) {
this.child.kill("SIGKILL");
if (!(await waitExit(this.child)))
throw new Error("Chrome remains alive; profile and lock preserved");
}
}
if (!this.persistent && this.profileDir)
await rm(this.profileDir, { recursive: true, force: true });
if (this.spoolDir)
await rm(this.spoolDir, { recursive: true, force: true });
for (const f of this.spoolFiles)
if (!f.bytes) await rm(f.file, { recursive: true, force: true });
if (this.persistent && this.ownedLock) {
const lock = join(this.profileDir, ".pi-chrome-cdp-lock");
const current = JSON.parse(await readFile(lock, "utf8")) as {
token?: string;
};
if (current.token !== this.lockToken)
throw new Error("profile lock ownership changed; refusing release");
await rm(lock);
this.ownedLock = false;
}
this.closed = true;
this.diagnostic?.("browser.close.finish", { status: "ok" });
}
}
export async function executeOperations(
b: Browser,
ops: Operation[],
signal?: AbortSignal,
replayFrom?: number,
onProgress?: (completed: number, active: string) => void,
): Promise<{
output: unknown[];
images: { data: string; mimeType: string }[];
events: Event[];
overflow: string[];
status: ReturnType<Browser["status"]>;
cursor: number;
failure?: string;
failedOperation?: { operation: number; kind: Operation["kind"] };
skipped?: { operation: number; kind: Operation["kind"] }[];
}> {
const results: unknown[] = [],
images: { data: string; mimeType: string }[] = [];
let failure: string | undefined;
let failedOperation:
| { operation: number; kind: Operation["kind"] }
| undefined;
const skipped = (from: number) =>
ops.slice(from).map((op, index) => ({
operation: from + index,
kind: op.kind,
}));
const startSeq = replayFrom ?? b.cursor;
const closeIndex = ops.findIndex((op) => op.kind === "close");
if (closeIndex >= 0 && ops.length !== 1)
return {
output: [],
images,
events: [],
overflow: [],
status: b.status(),
cursor: b.cursor,
failure: "close must be the only operation in its call",
failedOperation: { operation: closeIndex, kind: "close" },
skipped: skipped(closeIndex + 1),
};
const waits = new Map<string, Watcher>();
const knownLabels = new Set<string>();
for (const [index, op] of ops.entries()) {
if (op.kind === "cdp" && op.label) knownLabels.add(op.label);
if (op.kind === "wait") {
if (!op.after || !knownLabels.has(op.after) || !op.event || !op.sessionId)
return {
output: [],
images,
events: [],
overflow: [],
status: b.status(),
cursor: b.cursor,
failure:
"wait must follow a labeled cdp command in the same batch and include after, event, and sessionId",
failedOperation: { operation: index, kind: "wait" },
skipped: skipped(index + 1),
};
if (waits.has(op.after))
return {
output: [],
images,
events: [],
overflow: [],
status: b.status(),
cursor: b.cursor,
failure: `multiple waits use after label: ${op.after}; use distinct labeled cdp commands`,
failedOperation: { operation: index, kind: "wait" },
skipped: skipped(index + 1),
};
waits.set(op.after, {
method: op.event,
sessionId: op.sessionId,
since: 0,
});
}
}
b.beginBatchSpool(startSeq);
const session = (op: Operation): string | undefined => {
if (!op.sessionId) return undefined;
const sid = ref(op.sessionId, b.results);
if (typeof sid !== "string")
throw new Error("sessionId reference must resolve to a string");
return sid;
};
try {
for (let i = 0; i < ops.length; i++) {
const op = ops[i];
const end = b.operationSpan?.("operation", { index: i, kind: op.kind });
onProgress?.(i, op.kind === "cdp" ? `cdp ${op.method ?? ""}` : op.kind);
try {
if (signal?.aborted) {
b.uncertain = true;
throw new Error("canceled before next operation");
}
if (b.uncertain && op.kind !== "status" && op.kind !== "close")
throw new Error("browser state uncertain; call managed status first");
if (op.kind === "cdp") {
if (!op.method) throw new Error("method required");
const sid = session(op);
if (
!sid &&
/^(Page|DOM|Runtime|Log|Network|CSS|Input|Accessibility|Performance|Emulation|Fetch)\./.test(
op.method,
)
)
throw new Error("page command requires sessionId");
if (op.label) {
if (!/^[a-zA-Z][\w-]*$/.test(op.label))
throw new Error("invalid command label");
if (b.results.has(op.label))
throw new Error(`duplicate label: ${op.label}`);
}
if (op.label && waits.has(op.label)) {
const watcher = waits.get(op.label);
if (!watcher) throw new Error("wait not armed");
watcher.sessionId =
session({ ...op, sessionId: watcher.sessionId }) || "";
watcher.since = b.seq;
b.watchers.push(watcher);
}
const p = ref(op.params ?? {}, b.results);
const frame = await b.send(
op.method,
p,
sid,
signal,
op.timeoutMs || COMMAND_MS,
);
const serialized =
op.label && !frame.error ? JSON.stringify(frame) : undefined;
if (
op.label &&
serialized &&
Buffer.byteLength(serialized) <= INLINE
) {
b.results.set(op.label, frame);
if (b.results.size > 128) {
const oldest = b.results.keys().next().value;
if (oldest !== undefined) b.results.delete(oldest);
}
}
if (
sid &&
/^(Runtime|Log)\.(enable|disable)$/.test(op.method) &&
!frame.error
)
b.monitor.set(
`${sid}:${op.method.split(".")[0]}`,
op.method.endsWith("enable"),
);
const entry: Record<string, unknown> = {
operation: i,
command: op.method,
...(frame.transportError
? { transport: frame }
: { raw: await b.spoolFrame(frame, serialized) }),
};
results.push(entry);
if (frame.error) throw new Error(frame.error.message);
if (
op.method === "Target.setAutoAttach" &&
p &&
typeof p === "object" &&
!Array.isArray(p) &&
p.autoAttach === false
) {
b.autoAttach = false;
entry.interference = "managed descendant monitoring disabled";
}
if (
(op.method === "Page.stopScreencast" ||
op.method === "Page.startScreencast") &&
sid &&
b.video?.sessionId === sid
) {
entry.interference = "managed video interrupted";
entry.partialFile = await b.stopVideo(false);
}
if (op.method === "Tracing.end" && b.trace) {
entry.interference = "managed trace ended by raw command";
entry.partialFile = await b.stopTrace(false);
}
} else if (op.kind === "status") {
await b.refreshStatus();
results.push({
operation: i,
status: b.status(),
note: "observed status; earlier unconfirmed command may still take effect",
});
} else if (op.kind === "wait") {
const watcher = op.after ? waits.get(op.after) : undefined;
if (!watcher) throw new Error("wait not armed");
let hit = watcher.hit;
const end = Date.now() + (op.timeoutMs || 5000);
while (!hit && Date.now() < end && !signal?.aborted) {
await new Promise((r) => setTimeout(r, 20));
hit = watcher.hit;
}
if (signal?.aborted) {
b.uncertain = true;
throw new Error("wait canceled; command effects uncertain");
}
if (!hit)
throw new Error(
`wait timed out: ${op.event} on ${watcher.sessionId}`,
);
results.push({ operation: i, wait: hit });
} else if (op.kind === "screenshot") {
const sid = session(op);
if (!sid) throw new Error("screenshot requires sessionId");
const image = await b.screenshot(sid, op.filename, op.fullPage);
images.push({ data: image.data, mimeType: "image/png" });
results.push({ operation: i, file: image.file });
} else if (op.kind === "videoStart") {
const sid = session(op);
if (!sid) throw new Error("videoStart requires sessionId");
results.push({
operation: i,
file: await b.startVideo(sid, op.filename),
});
} else if (op.kind === "videoStop")
results.push({ operation: i, file: await b.stopVideo() });
else if (op.kind === "traceStart")
results.push({ operation: i, file: await b.startTrace(op.filename) });
else if (op.kind === "traceStop")
results.push({ operation: i, file: await b.stopTrace() });
else if (op.kind === "close") {
await b.close();
results.push({ operation: i, closed: true });
break;
} else throw new Error(`unknown operation: ${String(op.kind)}`);
} catch (error) {
failure = error instanceof Error ? error.message : String(error);
failedOperation = { operation: i, kind: op.kind };
end?.("error");
break;
}
end?.("finish");
onProgress?.(i + 1, "");
}
} finally {
const owned = new Set(waits.values());
b.watchers = b.watchers.filter((w) => !owned.has(w));
}
if (b.closed)
return {
output: results,
images,
events: [],
overflow: [],
status: b.status(),
cursor: b.seq,
...(failure
? {
failure,
failedOperation,
skipped: skipped((failedOperation?.operation || 0) + 1),
}
: {}),
};
await b.spoolWrites;
const candidates = b.events.filter((e) => e.seq > startSeq);
const events: Event[] = [];
let eventBytes = 0;
for (let i = candidates.length - 1; i >= 0; i--) {
const e = candidates[i];
const size = Buffer.byteLength(JSON.stringify(e));
if (eventBytes + size <= 65536) {
events.unshift(e);
eventBytes += size;
} else await b.spoolEvent(e);
}
const overflow = b.spoolFiles
.filter(
(f) =>
f.bytes > 0 &&
f.lastEventSeq !== undefined &&
f.lastEventSeq > startSeq,
)
.map((f) => {
b.protectSpoolFile(f.file);
return f.file;
});
b.cursor = b.seq;
return {
output: results,
images,
events,
overflow,
status: b.status(),
cursor: b.cursor,
...(failure
? {
failure,
failedOperation,
skipped: skipped((failedOperation?.operation || 0) + 1),
}
: {}),
};
}