Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/chrome-cdp/browser.ts

Raw
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<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),
				}
			: {}),
	};
}