Add pi-browser-harness ext (vendored 0.11.0, sharp dropped)
This commit is contained in:
368
extensions/pi-browser-harness/src/daemon/bridge.ts
Normal file
368
extensions/pi-browser-harness/src/daemon/bridge.ts
Normal file
@@ -0,0 +1,368 @@
|
||||
import WebSocket from "ws";
|
||||
import { discoverWsUrl } from "../cdp/discovery";
|
||||
import { isCdpRawMessage } from "../cdp/types";
|
||||
import type { CdpRawMessage } from "../cdp/types";
|
||||
import type { WireRequest, WireResponse, WireEvent } from "./protocol";
|
||||
import { CDP_CONNECT_TIMEOUT_MS, CDP_COMMAND_TIMEOUT_MS } from "./protocol";
|
||||
import { asString, isRecord } from "../util/guards";
|
||||
|
||||
export type SendToClient = (clientId: string, msg: WireResponse | WireEvent) => void;
|
||||
|
||||
export type EventHandler = (event: WireEvent, targetClientIds: string[]) => void;
|
||||
|
||||
export type CloseHandler = () => void;
|
||||
|
||||
export type CdpBridge = {
|
||||
start(): Promise<void>;
|
||||
stop(): Promise<void>;
|
||||
handleRequest(req: WireRequest, clientId: string, send: SendToClient): Promise<void>;
|
||||
isAlive(): boolean;
|
||||
getSessionOwner(sessionId: string): string | undefined;
|
||||
removeClient(clientId: string): void;
|
||||
onEvent(handler: EventHandler): void;
|
||||
onClose(handler: CloseHandler): void;
|
||||
};
|
||||
|
||||
export type InFlight = {
|
||||
readonly clientId: string;
|
||||
readonly localId: number;
|
||||
readonly send: SendToClient;
|
||||
readonly timer: ReturnType<typeof setTimeout>;
|
||||
readonly isAttach: boolean;
|
||||
};
|
||||
|
||||
// One map, keyed by the daemon-side id: the client that asked, the id it used, and everything needed to answer it.
|
||||
export type IdMultiplexer = {
|
||||
allocate(entry: Omit<InFlight, "timer">, arm: (daemonId: number) => InFlight["timer"]): number;
|
||||
take(daemonId: number): InFlight | undefined;
|
||||
takeAll(): ReadonlyArray<InFlight>;
|
||||
clearClient(clientId: string): void;
|
||||
};
|
||||
|
||||
export const createIdMultiplexer = (): IdMultiplexer => {
|
||||
let nextId = 1;
|
||||
const inFlight = new Map<number, InFlight>();
|
||||
|
||||
return {
|
||||
allocate(entry, arm) {
|
||||
const daemonId = nextId++;
|
||||
inFlight.set(daemonId, { ...entry, timer: arm(daemonId) });
|
||||
return daemonId;
|
||||
},
|
||||
take(daemonId) {
|
||||
const e = inFlight.get(daemonId);
|
||||
if (!e) return undefined;
|
||||
inFlight.delete(daemonId);
|
||||
clearTimeout(e.timer);
|
||||
return e;
|
||||
},
|
||||
takeAll() {
|
||||
const all = [...inFlight.values()];
|
||||
inFlight.clear();
|
||||
for (const e of all) clearTimeout(e.timer);
|
||||
return all;
|
||||
},
|
||||
// A disconnected client must not be answered later: dropping its entries also cancels their timeouts.
|
||||
clearClient(clientId) {
|
||||
for (const [daemonId, e] of inFlight) {
|
||||
if (e.clientId !== clientId) continue;
|
||||
clearTimeout(e.timer);
|
||||
inFlight.delete(daemonId);
|
||||
}
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
export type EventRouter = {
|
||||
record(clientId: string, sessionId: string): void;
|
||||
release(sessionId: string): void;
|
||||
getOwner(sessionId: string): string | undefined;
|
||||
removeClient(clientId: string): void;
|
||||
route(event: WireEvent): string[];
|
||||
};
|
||||
|
||||
export const createEventRouter = (): EventRouter => {
|
||||
const owners = new Map<string, string>();
|
||||
const clientSessions = new Map<string, Set<string>>();
|
||||
|
||||
return {
|
||||
record(clientId, sessionId) {
|
||||
const prev = owners.get(sessionId);
|
||||
if (prev && prev !== clientId) clientSessions.get(prev)?.delete(sessionId);
|
||||
owners.set(sessionId, clientId);
|
||||
let s = clientSessions.get(clientId);
|
||||
if (!s) { s = new Set(); clientSessions.set(clientId, s); }
|
||||
s.add(sessionId);
|
||||
},
|
||||
release(sessionId) {
|
||||
const owner = owners.get(sessionId);
|
||||
if (owner) clientSessions.get(owner)?.delete(sessionId);
|
||||
owners.delete(sessionId);
|
||||
},
|
||||
getOwner(sessionId) {
|
||||
return owners.get(sessionId);
|
||||
},
|
||||
removeClient(clientId) {
|
||||
const sessions = clientSessions.get(clientId);
|
||||
if (sessions) { for (const sid of sessions) owners.delete(sid); }
|
||||
clientSessions.delete(clientId);
|
||||
},
|
||||
route(event) {
|
||||
if (event.sessionId) {
|
||||
const owner = owners.get(event.sessionId);
|
||||
return owner ? [owner] : [];
|
||||
}
|
||||
return [...clientSessions.keys()];
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
const CONNECT_WAIT_MS = 15_000;
|
||||
|
||||
export const createCdpBridge = (): CdpBridge => {
|
||||
let ws: WebSocket | null = null;
|
||||
let wsUrl: string | null = null;
|
||||
const mux = createIdMultiplexer();
|
||||
const router = createEventRouter();
|
||||
let eventHandler: EventHandler | null = null;
|
||||
let closeHandler: CloseHandler | null = null;
|
||||
|
||||
const connectedWaiters: Array<() => void> = [];
|
||||
const signalConnected = (): void => {
|
||||
for (const w of connectedWaiters.splice(0)) w();
|
||||
};
|
||||
|
||||
const waitForConnection = (timeoutMs: number): Promise<void> =>
|
||||
new Promise((resolve) => {
|
||||
let settled = false;
|
||||
const finish = (): void => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
clearTimeout(timer);
|
||||
resolve();
|
||||
};
|
||||
const timer = setTimeout(finish, timeoutMs);
|
||||
connectedWaiters.push(finish);
|
||||
});
|
||||
|
||||
const onChromeMessage = (raw: string): void => {
|
||||
let parsed: unknown;
|
||||
try { parsed = JSON.parse(raw); } catch { return; }
|
||||
if (!isCdpRawMessage(parsed)) return;
|
||||
const msg: CdpRawMessage = parsed;
|
||||
|
||||
if (msg.id !== undefined) {
|
||||
const cb = mux.take(msg.id);
|
||||
if (!cb) return;
|
||||
const localId = cb.localId;
|
||||
|
||||
if (msg.error) {
|
||||
cb.send(cb.clientId, {
|
||||
type: "response",
|
||||
id: localId,
|
||||
error: { code: msg.error.code ?? -1, message: msg.error.message },
|
||||
});
|
||||
} else {
|
||||
if (cb.isAttach && msg.result) {
|
||||
const sessionId = isRecord(msg.result) ? asString(msg.result["sessionId"]) : undefined;
|
||||
if (sessionId) {
|
||||
router.record(cb.clientId, sessionId);
|
||||
}
|
||||
}
|
||||
cb.send(cb.clientId, { type: "response", id: localId, result: msg.result });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (!msg.method) return;
|
||||
|
||||
if (msg.method === "Inspector.detached" && msg.sessionId) {
|
||||
router.release(msg.sessionId);
|
||||
}
|
||||
|
||||
const wireEvent: WireEvent = {
|
||||
type: "event",
|
||||
method: msg.method,
|
||||
...(msg.params !== undefined ? { params: msg.params } : {}),
|
||||
...(msg.sessionId !== undefined ? { sessionId: msg.sessionId } : {}),
|
||||
};
|
||||
|
||||
if (eventHandler) {
|
||||
const targets = router.route(wireEvent);
|
||||
if (targets.length > 0) {
|
||||
eventHandler(wireEvent, targets);
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
let reconnectTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
let reconnectAttempt = 0;
|
||||
let stopped = false;
|
||||
|
||||
const tryConnect = (): void => {
|
||||
if (stopped) return;
|
||||
if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = null; }
|
||||
|
||||
const attempt = async (): Promise<void> => {
|
||||
if (ws && ws.readyState === WebSocket.OPEN) return;
|
||||
|
||||
const cachedUrl = wsUrl;
|
||||
let url: string;
|
||||
if (cachedUrl === null || cachedUrl === "" || reconnectAttempt > 0) {
|
||||
const d = await discoverWsUrl();
|
||||
if (!d.success) { scheduleRetry(); return; }
|
||||
url = d.data;
|
||||
wsUrl = url;
|
||||
} else {
|
||||
url = cachedUrl;
|
||||
}
|
||||
|
||||
const settledPromise = new Promise<void>((settle) => {
|
||||
let settled = false;
|
||||
const settleOnce = (): void => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
settle();
|
||||
};
|
||||
|
||||
let sock: WebSocket;
|
||||
try {
|
||||
sock = new WebSocket(url, { perMessageDeflate: false });
|
||||
} catch {
|
||||
scheduleRetry();
|
||||
return;
|
||||
}
|
||||
|
||||
const timer = setTimeout(() => {
|
||||
sock.close();
|
||||
settleOnce();
|
||||
}, CDP_CONNECT_TIMEOUT_MS);
|
||||
|
||||
sock.on("open", () => {
|
||||
clearTimeout(timer);
|
||||
ws = sock;
|
||||
reconnectAttempt = 0;
|
||||
ws.send(JSON.stringify({ id: 0, method: "Target.setDiscoverTargets", params: { discover: true } }));
|
||||
console.log("[pi-browser-daemon] Connected to Chrome ✓");
|
||||
signalConnected();
|
||||
settleOnce();
|
||||
});
|
||||
|
||||
sock.on("message", (data: WebSocket.Data) => {
|
||||
onChromeMessage(typeof data === "string" ? data : data.toString());
|
||||
});
|
||||
|
||||
sock.on("error", () => {
|
||||
clearTimeout(timer);
|
||||
settleOnce();
|
||||
});
|
||||
|
||||
sock.on("close", () => {
|
||||
clearTimeout(timer);
|
||||
ws = null;
|
||||
for (const cb of mux.takeAll()) {
|
||||
cb.send(cb.clientId, {
|
||||
type: "response",
|
||||
id: cb.localId,
|
||||
error: { code: -32000, message: "Chrome disconnected" },
|
||||
});
|
||||
}
|
||||
closeHandler?.();
|
||||
scheduleRetry();
|
||||
});
|
||||
});
|
||||
|
||||
await settledPromise;
|
||||
if (!ws || ws.readyState !== WebSocket.OPEN) {
|
||||
scheduleRetry();
|
||||
}
|
||||
};
|
||||
|
||||
attempt().catch(() => scheduleRetry());
|
||||
};
|
||||
|
||||
const scheduleRetry = (): void => {
|
||||
if (stopped) return;
|
||||
const delay = Math.min(1000 * Math.pow(2, Math.min(reconnectAttempt, 6)), 60_000);
|
||||
reconnectAttempt++;
|
||||
wsUrl = null;
|
||||
reconnectTimer = setTimeout(tryConnect, delay);
|
||||
};
|
||||
|
||||
const start = async (): Promise<void> => {
|
||||
stopped = false;
|
||||
tryConnect();
|
||||
};
|
||||
|
||||
const stop = async (): Promise<void> => {
|
||||
stopped = true;
|
||||
if (reconnectTimer) { clearTimeout(reconnectTimer); reconnectTimer = null; }
|
||||
mux.takeAll();
|
||||
signalConnected();
|
||||
if (ws) { try { ws.close(1000, "Shutdown"); } catch {} ws = null; }
|
||||
wsUrl = null;
|
||||
};
|
||||
|
||||
const handleRequest = async (
|
||||
req: WireRequest,
|
||||
clientId: string,
|
||||
send: SendToClient,
|
||||
): Promise<void> => {
|
||||
// Every queued request awaits the same connect signal rather than each spinning its own poll loop.
|
||||
if (!ws || ws.readyState !== WebSocket.OPEN) await waitForConnection(CONNECT_WAIT_MS);
|
||||
|
||||
if (!ws || ws.readyState !== WebSocket.OPEN) {
|
||||
send(clientId, {
|
||||
type: "response",
|
||||
id: req.id,
|
||||
error: { code: -32000, message: "Chrome not connected" },
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
const isDetach = req.method === "Target.detachFromTarget";
|
||||
if (isDetach && req.sessionId) router.release(req.sessionId);
|
||||
|
||||
const daemonId = mux.allocate(
|
||||
{ clientId, localId: req.id, send, isAttach: req.method === "Target.attachToTarget" },
|
||||
(id) =>
|
||||
setTimeout(() => {
|
||||
mux.take(id);
|
||||
send(clientId, {
|
||||
type: "response",
|
||||
id: req.id,
|
||||
error: { code: -32000, message: `Timeout after ${CDP_COMMAND_TIMEOUT_MS}ms: ${req.method}` },
|
||||
});
|
||||
}, CDP_COMMAND_TIMEOUT_MS),
|
||||
);
|
||||
|
||||
const payload: Record<string, unknown> = {
|
||||
id: daemonId,
|
||||
method: req.method,
|
||||
params: req.params ?? {},
|
||||
};
|
||||
if (req.sessionId) payload["sessionId"] = req.sessionId;
|
||||
|
||||
try {
|
||||
ws.send(JSON.stringify(payload));
|
||||
} catch (e) {
|
||||
mux.take(daemonId);
|
||||
send(clientId, {
|
||||
type: "response",
|
||||
id: req.id,
|
||||
error: { code: -32000, message: e instanceof Error ? e.message : String(e) },
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
return {
|
||||
start,
|
||||
stop,
|
||||
handleRequest,
|
||||
isAlive: () => ws !== null && ws.readyState === WebSocket.OPEN,
|
||||
getSessionOwner: (sid) => router.getOwner(sid),
|
||||
removeClient: (cid) => { mux.clearClient(cid); router.removeClient(cid); },
|
||||
onEvent: (h) => { eventHandler = h; },
|
||||
onClose: (h) => { closeHandler = h; },
|
||||
};
|
||||
};
|
||||
86
extensions/pi-browser-harness/src/daemon/index.ts
Normal file
86
extensions/pi-browser-harness/src/daemon/index.ts
Normal file
@@ -0,0 +1,86 @@
|
||||
import { createIpcServer } from "./server";
|
||||
import { createCdpBridge, type SendToClient } from "./bridge";
|
||||
import { DAEMON_IDLE_TIMEOUT_MS } from "./protocol";
|
||||
|
||||
async function main() {
|
||||
console.log("[pi-browser-daemon] Starting...");
|
||||
|
||||
const ipcServer = createIpcServer();
|
||||
const cdpBridge = createCdpBridge();
|
||||
|
||||
ipcServer.onMessage((msg, client) => {
|
||||
if (msg.type !== "request") return;
|
||||
|
||||
const send: SendToClient = (cid, resp) => {
|
||||
ipcServer.send(cid, resp);
|
||||
};
|
||||
|
||||
cdpBridge.handleRequest(msg, client.id, send);
|
||||
});
|
||||
|
||||
cdpBridge.onEvent((event, targetClientIds) => {
|
||||
for (const cid of targetClientIds) {
|
||||
ipcServer.send(cid, event);
|
||||
}
|
||||
});
|
||||
|
||||
ipcServer.onDisconnect((client) => {
|
||||
cdpBridge.removeClient(client.id);
|
||||
});
|
||||
|
||||
cdpBridge.onClose(() => {
|
||||
console.log("[pi-browser-daemon] Chrome disconnected");
|
||||
ipcServer.broadcast({
|
||||
type: "control",
|
||||
action: "shutdown",
|
||||
reason: "chrome_disconnected",
|
||||
});
|
||||
});
|
||||
|
||||
let idleTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
|
||||
const resetIdleTimer = (): void => {
|
||||
if (idleTimer) clearTimeout(idleTimer);
|
||||
if (ipcServer.clientCount() === 0) {
|
||||
idleTimer = setTimeout(() => {
|
||||
console.log("[pi-browser-daemon] Idle timeout — no clients for " +
|
||||
`${DAEMON_IDLE_TIMEOUT_MS / 60000} minutes. Shutting down.`);
|
||||
shutdown();
|
||||
}, DAEMON_IDLE_TIMEOUT_MS);
|
||||
}
|
||||
};
|
||||
|
||||
const cancelIdleTimer = (): void => {
|
||||
if (idleTimer) { clearTimeout(idleTimer); idleTimer = null; }
|
||||
};
|
||||
|
||||
ipcServer.onConnect(() => { cancelIdleTimer(); });
|
||||
ipcServer.onDisconnect(() => { resetIdleTimer(); });
|
||||
|
||||
try {
|
||||
await ipcServer.start();
|
||||
console.log("[pi-browser-daemon] IPC server listening");
|
||||
} catch (e) {
|
||||
console.error("[pi-browser-daemon] Failed to start IPC server:", e);
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
await cdpBridge.start();
|
||||
|
||||
resetIdleTimer();
|
||||
|
||||
const shutdown = async () => {
|
||||
console.log("[pi-browser-daemon] Shutting down...");
|
||||
await cdpBridge.stop();
|
||||
await ipcServer.stop();
|
||||
process.exit(0);
|
||||
};
|
||||
|
||||
process.on("SIGINT", shutdown);
|
||||
process.on("SIGTERM", shutdown);
|
||||
}
|
||||
|
||||
main().catch((e) => {
|
||||
console.error("[pi-browser-daemon] Fatal:", e);
|
||||
process.exit(1);
|
||||
});
|
||||
110
extensions/pi-browser-harness/src/daemon/protocol.ts
Normal file
110
extensions/pi-browser-harness/src/daemon/protocol.ts
Normal file
@@ -0,0 +1,110 @@
|
||||
import { isRecord } from "../util/guards";
|
||||
|
||||
export type WireRequest = {
|
||||
readonly type: "request";
|
||||
readonly id: number;
|
||||
readonly method: string;
|
||||
readonly params?: Record<string, unknown>;
|
||||
readonly sessionId?: string;
|
||||
};
|
||||
|
||||
export type WireResponse = {
|
||||
readonly type: "response";
|
||||
readonly id: number;
|
||||
readonly result?: unknown;
|
||||
readonly error?: {
|
||||
readonly code: number;
|
||||
readonly message: string;
|
||||
readonly data?: string;
|
||||
};
|
||||
};
|
||||
|
||||
export type WireEvent = {
|
||||
readonly type: "event";
|
||||
readonly method: string;
|
||||
readonly params?: Record<string, unknown>;
|
||||
readonly sessionId?: string;
|
||||
};
|
||||
|
||||
export type WireControl = {
|
||||
readonly type: "control";
|
||||
readonly action: "register" | "registered" | "deregister" | "shutdown";
|
||||
readonly clientId?: string;
|
||||
readonly reason?: string;
|
||||
};
|
||||
|
||||
export type WireMessage = WireRequest | WireResponse | WireEvent | WireControl;
|
||||
|
||||
const isWireMessage = (v: unknown): v is WireMessage => {
|
||||
if (!isRecord(v)) return false;
|
||||
const t = v["type"];
|
||||
if (typeof t !== "string") return false;
|
||||
|
||||
switch (t) {
|
||||
case "request": {
|
||||
if (typeof v["id"] !== "number") return false;
|
||||
if (typeof v["method"] !== "string") return false;
|
||||
const sid = v["sessionId"];
|
||||
if (sid !== undefined && typeof sid !== "string") return false;
|
||||
return true;
|
||||
}
|
||||
case "response": {
|
||||
if (typeof v["id"] !== "number") return false;
|
||||
const errVal = v["error"];
|
||||
if (errVal !== undefined) {
|
||||
if (!isRecord(errVal)) return false;
|
||||
if (typeof errVal["message"] !== "string") return false;
|
||||
if (typeof errVal["code"] !== "number") return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
case "event": {
|
||||
if (typeof v["method"] !== "string") return false;
|
||||
const sid = v["sessionId"];
|
||||
if (sid !== undefined && typeof sid !== "string") return false;
|
||||
return true;
|
||||
}
|
||||
case "control": {
|
||||
const action = v["action"];
|
||||
if (typeof action !== "string") return false;
|
||||
if (!["register", "registered", "deregister", "shutdown"].includes(action)) return false;
|
||||
const cid = v["clientId"];
|
||||
if (
|
||||
(action === "register" || action === "registered" || action === "deregister") &&
|
||||
typeof cid !== "string"
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
default:
|
||||
return false;
|
||||
}
|
||||
};
|
||||
|
||||
export const serialize = (msg: WireMessage): string => JSON.stringify(msg);
|
||||
|
||||
export const deserialize = (line: string): WireMessage | null => {
|
||||
let parsed: unknown;
|
||||
try {
|
||||
parsed = JSON.parse(line);
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
if (!isWireMessage(parsed)) return null;
|
||||
return parsed;
|
||||
};
|
||||
|
||||
// Windows requires a named pipe (`\\.\pipe\<name>`); a Unix path is not a valid `net` listen/connect target there.
|
||||
export const DAEMON_SOCKET_PATH =
|
||||
process.platform === "win32"
|
||||
? "\\\\.\\pipe\\pi-browser-daemon"
|
||||
: "/tmp/pi-browser-daemon.sock";
|
||||
|
||||
export const DAEMON_IDLE_TIMEOUT_MS = 30 * 60 * 1000;
|
||||
|
||||
export const DAEMON_MAX_CLIENTS = 16;
|
||||
|
||||
export const CDP_CONNECT_TIMEOUT_MS = 10_000;
|
||||
|
||||
export const CDP_COMMAND_TIMEOUT_MS = 10_000;
|
||||
85
extensions/pi-browser-harness/src/daemon/queue.ts
Normal file
85
extensions/pi-browser-harness/src/daemon/queue.ts
Normal file
@@ -0,0 +1,85 @@
|
||||
import { type Result, err } from "../util/result";
|
||||
import { type CdpError, cdpError } from "../cdp/errors";
|
||||
import type { CdpEvent } from "../cdp/types";
|
||||
|
||||
export type EventQueue = {
|
||||
readonly push: (e: CdpEvent) => void;
|
||||
readonly end: () => void;
|
||||
readonly iter: AsyncIterable<CdpEvent>;
|
||||
};
|
||||
|
||||
export const makeEventQueue = (): EventQueue => {
|
||||
const buf: CdpEvent[] = [];
|
||||
const waiters: Array<(v: IteratorResult<CdpEvent, undefined>) => void> = [];
|
||||
let ended = false;
|
||||
return {
|
||||
push(e) {
|
||||
if (ended) return;
|
||||
const w = waiters.shift();
|
||||
if (w) w({ value: e, done: false });
|
||||
else buf.push(e);
|
||||
},
|
||||
end() {
|
||||
ended = true;
|
||||
for (const w of waiters.splice(0)) w({ value: undefined, done: true });
|
||||
},
|
||||
iter: {
|
||||
[Symbol.asyncIterator](): AsyncIterator<CdpEvent, undefined, undefined> {
|
||||
return {
|
||||
next: (): Promise<IteratorResult<CdpEvent, undefined>> =>
|
||||
new Promise((resolve) => {
|
||||
const next = buf.shift();
|
||||
if (next) resolve({ value: next, done: false });
|
||||
else if (ended) resolve({ value: undefined, done: true });
|
||||
else waiters.push(resolve);
|
||||
}),
|
||||
};
|
||||
},
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
export type CdpResult = Result<unknown, CdpError>;
|
||||
|
||||
export type Pending = {
|
||||
readonly resolve: (v: Result<unknown, CdpError>) => void;
|
||||
readonly timer: ReturnType<typeof setTimeout>;
|
||||
readonly method: string;
|
||||
};
|
||||
|
||||
export const rejectAllPending = (pending: Map<number, Pending>, reason: string): void => {
|
||||
for (const [, p] of pending) {
|
||||
clearTimeout(p.timer);
|
||||
p.resolve(err(cdpError("transport_closed", reason, p.method)));
|
||||
}
|
||||
pending.clear();
|
||||
};
|
||||
|
||||
export const sendWithTimeout = (
|
||||
pending: Map<number, Pending>,
|
||||
id: number,
|
||||
method: string,
|
||||
timeoutMs: number,
|
||||
timeoutLabel: string,
|
||||
send: () => void,
|
||||
): Promise<CdpResult> =>
|
||||
new Promise((resolve) => {
|
||||
const timer = setTimeout(() => {
|
||||
pending.delete(id);
|
||||
resolve(err(cdpError("timeout", `${timeoutLabel} timeout after ${timeoutMs}ms: ${method}`, method)));
|
||||
}, timeoutMs);
|
||||
pending.set(id, { resolve, timer, method });
|
||||
try {
|
||||
send();
|
||||
} catch (e) {
|
||||
clearTimeout(timer);
|
||||
pending.delete(id);
|
||||
resolve(err(cdpError("transport_closed", e instanceof Error ? e.message : String(e), method)));
|
||||
}
|
||||
});
|
||||
|
||||
export const makeOnClose = (closeListeners: Set<() => void>) =>
|
||||
(cb: () => void): (() => void) => {
|
||||
closeListeners.add(cb);
|
||||
return () => closeListeners.delete(cb);
|
||||
};
|
||||
170
extensions/pi-browser-harness/src/daemon/server.ts
Normal file
170
extensions/pi-browser-harness/src/daemon/server.ts
Normal file
@@ -0,0 +1,170 @@
|
||||
import { createServer, type Server, type Socket } from "node:net";
|
||||
import { createInterface, type Interface } from "node:readline";
|
||||
import { unlinkSync } from "node:fs";
|
||||
import {
|
||||
type WireMessage,
|
||||
type WireControl,
|
||||
deserialize,
|
||||
serialize,
|
||||
DAEMON_SOCKET_PATH,
|
||||
DAEMON_MAX_CLIENTS,
|
||||
} from "./protocol";
|
||||
|
||||
export type ClientSocket = {
|
||||
id: string;
|
||||
readonly socket: Socket;
|
||||
readonly rl: Interface;
|
||||
registered: boolean;
|
||||
};
|
||||
|
||||
export type MessageHandler = (msg: WireMessage, client: ClientSocket) => void;
|
||||
|
||||
export type ConnectionHandler = (client: ClientSocket) => void;
|
||||
|
||||
export type IpcServer = {
|
||||
start(): Promise<void>;
|
||||
stop(): Promise<void>;
|
||||
onMessage(handler: MessageHandler): void;
|
||||
onConnect(handler: ConnectionHandler): void;
|
||||
onDisconnect(handler: ConnectionHandler): void;
|
||||
send(clientId: string, msg: WireMessage): void;
|
||||
broadcast(msg: WireMessage): void;
|
||||
clientCount(): number;
|
||||
};
|
||||
|
||||
export const createIpcServer = (): IpcServer => {
|
||||
let server: Server | null = null;
|
||||
const sockets = new Map<string, ClientSocket>();
|
||||
const clients = new Map<string, ClientSocket>();
|
||||
let messageHandler: MessageHandler | null = null;
|
||||
let connectHandler: ConnectionHandler | null = null;
|
||||
let disconnectHandler: ConnectionHandler | null = null;
|
||||
let nextTempId = 1;
|
||||
|
||||
const start = (): Promise<void> =>
|
||||
new Promise((resolve, reject) => {
|
||||
try {
|
||||
unlinkSync(DAEMON_SOCKET_PATH);
|
||||
} catch {}
|
||||
|
||||
server = createServer({ pauseOnConnect: false }, (socket: Socket) => {
|
||||
if (clients.size >= DAEMON_MAX_CLIENTS) {
|
||||
socket.destroy();
|
||||
return;
|
||||
}
|
||||
|
||||
const tempId = `anon-${nextTempId++}`;
|
||||
const rl = createInterface({ input: socket, crlfDelay: Infinity });
|
||||
let currentClient: ClientSocket = { id: tempId, socket, rl, registered: false };
|
||||
|
||||
sockets.set(tempId, currentClient);
|
||||
|
||||
rl.on("line", (line: string) => {
|
||||
const msg = deserialize(line);
|
||||
if (!msg) return;
|
||||
|
||||
if (msg.type === "control" && msg.action === "register" && !currentClient.registered) {
|
||||
const clientId = msg.clientId;
|
||||
if (!clientId) return;
|
||||
|
||||
if (clients.has(clientId)) {
|
||||
const reply: WireControl = {
|
||||
type: "control",
|
||||
action: "shutdown",
|
||||
reason: `clientId ${clientId} is already connected`,
|
||||
};
|
||||
socket.write(serialize(reply) + "\n");
|
||||
socket.destroy();
|
||||
return;
|
||||
}
|
||||
|
||||
// Mutate in place so the already-captured closure sees the promoted client.
|
||||
sockets.delete(tempId);
|
||||
currentClient.id = clientId;
|
||||
currentClient.registered = true;
|
||||
clients.set(clientId, currentClient);
|
||||
|
||||
const ack: WireControl = { type: "control", action: "registered", clientId };
|
||||
socket.write(serialize(ack) + "\n");
|
||||
|
||||
connectHandler?.(currentClient);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!currentClient.registered) return;
|
||||
|
||||
messageHandler?.(msg, currentClient);
|
||||
});
|
||||
|
||||
socket.on("error", () => {});
|
||||
|
||||
socket.on("close", () => {
|
||||
const found = clients.get(currentClient.id) ?? sockets.get(tempId);
|
||||
if (found) {
|
||||
clients.delete(found.id);
|
||||
sockets.delete(tempId);
|
||||
if (found.registered) {
|
||||
disconnectHandler?.(found);
|
||||
}
|
||||
}
|
||||
rl.close();
|
||||
});
|
||||
});
|
||||
|
||||
server.on("error", (err: NodeJS.ErrnoException) => {
|
||||
reject(err);
|
||||
});
|
||||
|
||||
server.listen(DAEMON_SOCKET_PATH, () => {
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
|
||||
const stop = (): Promise<void> =>
|
||||
new Promise((resolve) => {
|
||||
for (const [, client] of clients) {
|
||||
try { client.socket.destroy(); } catch {}
|
||||
}
|
||||
for (const [, client] of sockets) {
|
||||
try { client.socket.destroy(); } catch {}
|
||||
}
|
||||
clients.clear();
|
||||
sockets.clear();
|
||||
|
||||
if (server) {
|
||||
server.close(() => {
|
||||
try { unlinkSync(DAEMON_SOCKET_PATH); } catch {}
|
||||
resolve();
|
||||
});
|
||||
} else {
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
|
||||
const send = (clientId: string, msg: WireMessage): void => {
|
||||
const client = clients.get(clientId);
|
||||
if (!client) return;
|
||||
try {
|
||||
client.socket.write(serialize(msg) + "\n");
|
||||
} catch {}
|
||||
};
|
||||
|
||||
return {
|
||||
start,
|
||||
stop,
|
||||
onMessage(handler) {
|
||||
messageHandler = handler;
|
||||
},
|
||||
onConnect(handler) {
|
||||
connectHandler = handler;
|
||||
},
|
||||
onDisconnect(handler) {
|
||||
disconnectHandler = handler;
|
||||
},
|
||||
send,
|
||||
broadcast(msg) {
|
||||
for (const [clientId] of clients) send(clientId, msg);
|
||||
},
|
||||
clientCount: (): number => clients.size,
|
||||
};
|
||||
};
|
||||
120
extensions/pi-browser-harness/src/daemon/spawn.ts
Normal file
120
extensions/pi-browser-harness/src/daemon/spawn.ts
Normal file
@@ -0,0 +1,120 @@
|
||||
import { spawn, type ChildProcess } from "node:child_process";
|
||||
import { access, unlink } from "node:fs/promises";
|
||||
import { existsSync } from "node:fs";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { dirname, join } from "node:path";
|
||||
import { createConnection } from "node:net";
|
||||
import { DAEMON_SOCKET_PATH, serialize, type WireControl } from "./protocol";
|
||||
|
||||
const moduleDir = dirname(fileURLToPath(import.meta.url));
|
||||
|
||||
const isWindows = process.platform === "win32";
|
||||
|
||||
// No-op on Windows: a named pipe is not a filesystem entry and unlinking a pipe path throws.
|
||||
const cleanupStaleSocket = (): void => {
|
||||
if (isWindows) return;
|
||||
unlink(DAEMON_SOCKET_PATH).catch(() => {});
|
||||
};
|
||||
|
||||
export const daemonSocketMayExist = async (
|
||||
platform: NodeJS.Platform = process.platform,
|
||||
socketPath: string = DAEMON_SOCKET_PATH,
|
||||
): Promise<boolean> => {
|
||||
if (platform === "win32") return true;
|
||||
try {
|
||||
await access(socketPath);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
};
|
||||
|
||||
export const isDaemonRunning = async (timeoutMs = 2_000): Promise<boolean> => {
|
||||
if (!(await daemonSocketMayExist())) return false;
|
||||
|
||||
return new Promise((resolve) => {
|
||||
let settled = false;
|
||||
const done = (alive: boolean): void => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
resolve(alive);
|
||||
};
|
||||
|
||||
const sock = createConnection(DAEMON_SOCKET_PATH);
|
||||
const timer = setTimeout(() => {
|
||||
sock.destroy();
|
||||
cleanupStaleSocket();
|
||||
done(false);
|
||||
}, timeoutMs);
|
||||
|
||||
sock.on("connect", () => {
|
||||
const regMsg: WireControl = { type: "control", action: "register", clientId: "_liveness_probe" };
|
||||
sock.write(serialize(regMsg) + "\n");
|
||||
});
|
||||
|
||||
sock.on("data", (chunk: Buffer) => {
|
||||
const text = chunk.toString();
|
||||
if (text.includes('"registered"')) {
|
||||
clearTimeout(timer);
|
||||
const deregMsg: WireControl = { type: "control", action: "deregister", clientId: "_liveness_probe" };
|
||||
try { sock.write(serialize(deregMsg) + "\n"); } catch {}
|
||||
sock.destroy();
|
||||
done(true);
|
||||
}
|
||||
});
|
||||
|
||||
sock.on("error", () => {
|
||||
clearTimeout(timer);
|
||||
sock.destroy();
|
||||
cleanupStaleSocket();
|
||||
done(false);
|
||||
});
|
||||
|
||||
// Closed without a 'registered' ack: the daemon is alive but rejected us, so do NOT unlink its socket.
|
||||
sock.on("close", () => {
|
||||
clearTimeout(timer);
|
||||
done(false);
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
export const spawnDaemon = (): ChildProcess | null => {
|
||||
const daemonScript = join(moduleDir, "index.ts");
|
||||
|
||||
// Windows npm binaries are `.cmd` shims, which only launch through a shell — spawning one directly yields EINVAL.
|
||||
const tsxBinBase = join(moduleDir, "..", "..", "node_modules", ".bin", "tsx");
|
||||
const tsxBin = isWindows ? `${tsxBinBase}.cmd` : tsxBinBase;
|
||||
const hasLocalTsx = existsSync(tsxBin);
|
||||
|
||||
const cmd = hasLocalTsx ? tsxBin : isWindows ? "npx.cmd" : "npx";
|
||||
const args = hasLocalTsx ? [daemonScript] : ["tsx", daemonScript];
|
||||
|
||||
// With shell: true Node passes the command line verbatim, so a path containing spaces must be quoted or it splits into tokens.
|
||||
const quote = (s: string): string => (isWindows ? `"${s}"` : s);
|
||||
|
||||
const child = spawn(quote(cmd), args.map(quote), {
|
||||
detached: true,
|
||||
stdio: "ignore",
|
||||
env: { ...process.env },
|
||||
shell: isWindows,
|
||||
windowsVerbatimArguments: isWindows,
|
||||
});
|
||||
|
||||
child.unref();
|
||||
return child;
|
||||
};
|
||||
|
||||
export const ensureDaemon = async (timeoutMs = 10_000): Promise<boolean> => {
|
||||
if (await isDaemonRunning()) return true;
|
||||
|
||||
const child = spawnDaemon();
|
||||
if (!child) return false;
|
||||
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (Date.now() < deadline) {
|
||||
if (await isDaemonRunning()) return true;
|
||||
await new Promise((r) => setTimeout(r, 200));
|
||||
}
|
||||
|
||||
return false;
|
||||
};
|
||||
178
extensions/pi-browser-harness/src/daemon/transport.ts
Normal file
178
extensions/pi-browser-harness/src/daemon/transport.ts
Normal file
@@ -0,0 +1,178 @@
|
||||
import { createConnection, type Socket } from "node:net";
|
||||
import { createInterface } from "node:readline";
|
||||
import { type Result, err, ok } from "../util/result";
|
||||
import { type CdpError, cdpError, classifyRemoteError } from "../cdp/errors";
|
||||
import { DEFAULT_TIMEOUT_MS, type CdpTransport } from "../cdp/types";
|
||||
import type { CdpEvent } from "../cdp/types";
|
||||
import { type Pending, makeEventQueue, makeOnClose, rejectAllPending, sendWithTimeout } from "./queue";
|
||||
import {
|
||||
DAEMON_SOCKET_PATH,
|
||||
type WireRequest,
|
||||
type WireControl,
|
||||
deserialize,
|
||||
serialize,
|
||||
} from "./protocol";
|
||||
|
||||
export const createDaemonTransport = (clientId: string): CdpTransport => {
|
||||
let socket: Socket | null = null;
|
||||
let rl: ReturnType<typeof createInterface> | null = null;
|
||||
let queue = makeEventQueue();
|
||||
const closeListeners = new Set<() => void>();
|
||||
const pending = new Map<number, Pending>();
|
||||
let registered = false;
|
||||
let nextRequestId = 1;
|
||||
|
||||
const cleanup = (reason: string): void => {
|
||||
rejectAllPending(pending, reason);
|
||||
queue.end();
|
||||
queue = makeEventQueue();
|
||||
registered = false;
|
||||
|
||||
if (rl) { rl.close(); rl = null; }
|
||||
if (socket) { try { socket.destroy(); } catch {} socket = null; }
|
||||
|
||||
for (const cb of closeListeners) cb();
|
||||
};
|
||||
|
||||
const connect = (_url: string, opts?: { timeoutMs?: number }): Promise<Result<void, CdpError>> => {
|
||||
const timeoutMs = opts?.timeoutMs ?? 10_000;
|
||||
|
||||
if (socket && !socket.destroyed && registered) {
|
||||
return Promise.resolve(ok(undefined));
|
||||
}
|
||||
|
||||
return new Promise((resolve) => {
|
||||
let settled = false;
|
||||
const settle = (r: Result<void, CdpError>): void => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
resolve(r);
|
||||
};
|
||||
|
||||
let sock: Socket;
|
||||
try {
|
||||
sock = createConnection(DAEMON_SOCKET_PATH);
|
||||
} catch (e) {
|
||||
settle(err(cdpError("transport_closed", e instanceof Error ? e.message : String(e))));
|
||||
return;
|
||||
}
|
||||
socket = sock;
|
||||
|
||||
const connectTimer = setTimeout(() => {
|
||||
cleanup("Connection timeout");
|
||||
settle(err(cdpError("timeout", `Daemon connection timed out after ${timeoutMs}ms`)));
|
||||
}, timeoutMs);
|
||||
|
||||
sock.on("connect", () => {
|
||||
clearTimeout(connectTimer);
|
||||
rl = createInterface({ input: sock, crlfDelay: Infinity });
|
||||
|
||||
rl.on("line", (line: string) => {
|
||||
const msg = deserialize(line);
|
||||
if (!msg) return;
|
||||
|
||||
if (msg.type === "control" && msg.action === "registered" && msg.clientId === clientId) {
|
||||
registered = true;
|
||||
settle(ok(undefined));
|
||||
return;
|
||||
}
|
||||
|
||||
if (msg.type === "control" && msg.action === "shutdown") {
|
||||
cleanup(`Daemon shutting down: ${msg.reason ?? "unknown"}`);
|
||||
return;
|
||||
}
|
||||
|
||||
if (!registered) return;
|
||||
|
||||
if (msg.type === "response") {
|
||||
const p = pending.get(msg.id);
|
||||
if (!p) return;
|
||||
pending.delete(msg.id);
|
||||
clearTimeout(p.timer);
|
||||
|
||||
if (msg.error) {
|
||||
p.resolve(err(cdpError(classifyRemoteError(msg.error.message), msg.error.message, p.method)));
|
||||
} else {
|
||||
p.resolve(ok(msg.result));
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (msg.type === "event") {
|
||||
queue.push({
|
||||
method: msg.method,
|
||||
params: msg.params,
|
||||
...(msg.sessionId !== undefined ? { sessionId: msg.sessionId } : {}),
|
||||
});
|
||||
return;
|
||||
}
|
||||
});
|
||||
|
||||
sock.on("error", () => {});
|
||||
|
||||
sock.on("close", () => {
|
||||
clearTimeout(connectTimer);
|
||||
cleanup("Daemon socket closed");
|
||||
settle(err(cdpError("transport_closed", "Daemon socket closed before registration")));
|
||||
});
|
||||
|
||||
const regMsg: WireControl = { type: "control", action: "register", clientId };
|
||||
sock.write(serialize(regMsg) + "\n");
|
||||
});
|
||||
|
||||
sock.on("error", (e: NodeJS.ErrnoException) => {
|
||||
clearTimeout(connectTimer);
|
||||
const msg = e.code === "ENOENT"
|
||||
? "Daemon not running — socket not found"
|
||||
: e.message;
|
||||
settle(err(cdpError("transport_closed", msg)));
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
const close = async (): Promise<void> => {
|
||||
if (registered && socket && !socket.destroyed) {
|
||||
const dereg: WireControl = { type: "control", action: "deregister", clientId };
|
||||
try { socket.write(serialize(dereg) + "\n"); } catch {}
|
||||
}
|
||||
cleanup("close() called");
|
||||
};
|
||||
|
||||
const request = (
|
||||
method: string,
|
||||
params: Record<string, unknown>,
|
||||
opts?: { sessionId?: string | null; timeoutMs?: number },
|
||||
): Promise<Result<unknown, CdpError>> => {
|
||||
const timeoutMs = opts?.timeoutMs ?? DEFAULT_TIMEOUT_MS;
|
||||
const sessionId = opts?.sessionId ?? null;
|
||||
|
||||
const sock = socket;
|
||||
if (sock === null || sock.destroyed || !registered) {
|
||||
return Promise.resolve(err(cdpError("transport_closed", "Daemon not connected", method)));
|
||||
}
|
||||
|
||||
const id = nextRequestId++;
|
||||
const req: WireRequest = {
|
||||
type: "request",
|
||||
id,
|
||||
method,
|
||||
params,
|
||||
...(sessionId !== null && sessionId !== undefined ? { sessionId } : {}),
|
||||
};
|
||||
|
||||
return sendWithTimeout(pending, id, method, timeoutMs, "Daemon", () => sock.write(serialize(req) + "\n"));
|
||||
};
|
||||
|
||||
const events = (): AsyncIterable<CdpEvent> => queue.iter;
|
||||
|
||||
const state = (): "open" | "closed" | "connecting" => {
|
||||
if (!socket) return "closed";
|
||||
if (!registered) return "connecting";
|
||||
if (socket.destroyed) return "closed";
|
||||
return "open";
|
||||
};
|
||||
|
||||
const onClose = makeOnClose(closeListeners);
|
||||
|
||||
return { connect, close, request, events, state, onClose };
|
||||
};
|
||||
Reference in New Issue
Block a user