import { 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; 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; }; type Trace = { file: string; done: Promise>; finish: (v: Record) => 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 { 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 { 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 { const f = await open(file, "wx", 0o600); try { await f.writeFile(bytes); } finally { await f.close(); } } async function waitExit(child: ChildProcess, ms = 3000): Promise { 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(); autoAttach = true; sessions = new Map(); results = new Map(); pending = new Map 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 = Promise.resolve(); queuedSpoolBytes = 0; private protectedSpoolFiles = new Set(); private batchSpooling = false; video?: Video; videoStopping?: Promise; trace?: Trace; traceStopping?: Promise; busy: Promise = 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 { 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((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).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).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 { 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; 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 | 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 | 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(write: () => Promise): Promise { const task = this.spoolWrites.then(write); this.spoolWrites = task.then( () => {}, () => {}, ); return task; } queueEventSpool(e: Event, frame?: Record): 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 { 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 { 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 { return this.queueSpool(() => this.writeSpoolEvent(e)); } send( method: string, params: Json = {}, sessionId?: string, signal?: AbortSignal, timeoutMs = COMMAND_MS, ): Promise { 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 { const targets = await this.send("Target.getTargets"); if (targets.error) throw new Error(`status unavailable: ${targets.error.message}`); this.observedTargets = ((targets.result as Record).targetInfos as Json[]) || []; const metrics = await this.send( "Page.getLayoutMetrics", {}, this.initialSession, ); if (!metrics.error) { const viewport = (metrics.result as Record) .cssVisualViewport as Record | 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 { 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).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 { 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((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((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 { 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((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 { if (this.trace) throw new Error("trace already running"); const file = await this.outputFile(filename, "json"); let finish!: (v: Record) => void; const done = new Promise>((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 { 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((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> | 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; 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 { 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 { 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; 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(); const knownLabels = new Set(); 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 = { 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), } : {}), }; }