import { createBrowserSession, connectBrowserSession, deleteBrowserSession, BrowserRenderingError, type BrowserBinding, type CdpSession, } from "agents/browser"; import { WebFailure } from "./web-errors"; /** One absolute deadline for acquisition, navigation and extraction. */ export class WebDeadline { readonly controller = new AbortController(); readonly signal = this.controller.signal; private timer: ReturnType; readonly expires: number; private readonly abort = () => this.controller.abort(new WebFailure("cancelled")); constructor( ms: number, private caller?: AbortSignal, ) { this.expires = Date.now() + ms; this.timer = setTimeout( () => this.controller.abort(new WebFailure("timeout")), ms, ); caller?.addEventListener("abort", this.abort, { once: true }); if (caller?.aborted) this.abort(); } remaining() { if (Date.now() >= this.expires && !this.signal.aborted) this.controller.abort(new WebFailure("timeout")); this.signal.throwIfAborted(); return Math.max(1, this.expires - Date.now()); } async run(operation: () => Promise): Promise { this.remaining(); let listener: () => void; const abort = new Promise((_, reject) => { listener = () => reject(this.signal.reason); this.signal.addEventListener("abort", listener, { once: true }); }); try { return await Promise.race([operation(), abort]); } finally { this.signal.removeEventListener("abort", listener!); } } dispose() { clearTimeout(this.timer); this.caller?.removeEventListener("abort", this.abort); } } export interface BrowserLeaseHooks { create?(expiresAt: number): Promise; acquired(sessionId: string, expiresAt: number): Promise; closed(sessionId: string): Promise; } export async function closeResearchBrowser( browser: BrowserBinding, sessionId: string, ) { const cleanup = new WebDeadline(5_000); try { await cleanup.run(() => deleteBrowserSession( { fetch: (url, init) => browser.fetch(url, { ...init, signal: cleanup.signal }), }, sessionId, ), ); } catch { throw new WebFailure("cleanup_failed"); } finally { cleanup.dispose(); } } export type ResearchBrowserEvent = | { method: "Fetch.requestPaused"; sessionId?: string; params: { requestId: string; request: { url: string }; resourceType: string; }; } | { method: "Network.responseReceived"; sessionId?: string; params: { type: string; frameId: string; response: { status: number } }; } | { method: "Target.attachedToTarget"; sessionId?: string; params: { targetInfo: { targetId: string } }; }; export interface ResearchBrowser { onEvent(listener: (event: ResearchBrowserEvent) => void): void; send(method: string, params?: unknown, target?: string): Promise; } /** Uses only public native CDP primitives. Each call owns its session in its * closure; no shared "current browser" can be overwritten by a parallel call. */ export async function withResearchBrowser( browser: BrowserBinding, deadline: WebDeadline, keepAlive: (promise: Promise) => void, run: (browser: ResearchBrowser) => Promise, lease?: BrowserLeaseHooks, ): Promise { let session: CdpSession | undefined; let sessionId: string | undefined; let closing: Promise | undefined; const listeners: ((event: ResearchBrowserEvent) => void)[] = []; const close = () => (closing ??= (async () => { session?.close(); // connectBrowserSession's close disconnects only. if (!sessionId) return; await closeResearchBrowser(browser, sessionId); if (lease) { const cleanup = new WebDeadline(5_000); const release = lease.closed(sessionId); keepAlive(release.catch(() => {})); try { await cleanup.run(() => release); } catch { /* Remote deletion succeeded; metadata can reconcile on wake. */ } finally { cleanup.dispose(); } } })()); // Creation has no native abort parameter. Keep its late continuation alive so // a returned ID is deleted even when the caller already stopped waiting. try { const acquisition = (async () => { deadline.remaining(); const workBinding: BrowserBinding = { fetch: async (url, init) => { const response = await browser.fetch(url, { ...init, signal: deadline.signal, }); // Native CdpSession owns command correlation. This narrow event tap // supports request policy before the browser continues a request. response.webSocket?.addEventListener("message", (message) => { if (typeof message.data !== "string") return; const event = JSON.parse(message.data) as ResearchBrowserEvent; if ( [ "Fetch.requestPaused", "Network.responseReceived", "Target.attachedToTarget", ].includes(event.method) ) for (const listener of listeners) listener(event); }); return response; }, }; const created = lease?.create ? { sessionId: await lease.create(deadline.expires) } : await createBrowserSession(workBinding, { keepAliveMs: 60_000 }); sessionId = created.sessionId; try { if (lease) { const registration = lease.acquired(sessionId, deadline.expires); keepAlive(registration.catch(() => {})); await deadline.run(() => registration); } if (deadline.signal.aborted) { await close(); throw deadline.signal.reason; } deadline.remaining(); const connected = await connectBrowserSession( workBinding, sessionId, deadline.remaining(), ); session = connected; if (deadline.signal.aborted) { connected.close(); await close(); throw deadline.signal.reason; } return connected; } catch (error) { // This continuation may run after the outer finally saw no session ID. // Registration failure must still close the now-known remote session. await close(); throw error; } })(); keepAlive( acquisition.then( () => {}, () => {}, ), ); await deadline.run(() => acquisition); return await deadline.run(() => run({ onEvent: (listener) => { listeners.push(listener); }, send: (method, params, target) => deadline.run(() => session!.send(method, params, { timeoutMs: deadline.remaining(), ...(target ? { sessionId: target } : {}), }), ), }), ); } catch (error) { // The parent's expiry alarm can close CDP before this isolate dispatches // its timer callback. Classify against the absolute clock as well. deadline.remaining(); if (error instanceof WebFailure || deadline.signal.aborted) throw error; if ( error instanceof Error && /^CDP command timed out after \d+ms: /.test(error.message) ) throw new WebFailure("timeout"); if (error instanceof BrowserRenderingError && error.status === 429) throw new WebFailure("search_blocked", 429); throw new WebFailure("browser_unavailable"); } finally { // Do not memoize a no-ID close: acquisition may still return an ID later. if (sessionId) await close(); } }