Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242import { 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<typeof setTimeout>; 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<T>(operation: () => Promise<T>): Promise<T> { this.remaining(); let listener: () => void; const abort = new Promise<never>((_, 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<string>; acquired(sessionId: string, expiresAt: number): Promise<void>; closed(sessionId: string): Promise<void>;}
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<unknown>;}
/** 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<T>( browser: BrowserBinding, deadline: WebDeadline, keepAlive: (promise: Promise<unknown>) => void, run: (browser: ResearchBrowser) => Promise<T>, lease?: BrowserLeaseHooks,): Promise<T> { let session: CdpSession | undefined; let sessionId: string | undefined; let closing: Promise<void> | 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(); }}