Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262import { BrowserRenderingError, connectBrowserSession, createBrowserSession, deleteBrowserSession, type BrowserBinding, type CdpSession,} from "agents/browser";import * as Effect from "effect/Effect";import { WebFailure, webValidation } 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()); } run<T>(operation: () => PromiseLike<T>): Effect.Effect<T, WebFailure> { return webValidation(() => this.remaining()).pipe( Effect.andThen(this.limit(browserCall(operation))), ); } limit<A, E>(program: Effect.Effect<A, E>) { const interrupted = Effect.callback<never, WebFailure>((resume) => { const abort = () => resume(Effect.fail(browserFailure(this.signal.reason))); this.signal.addEventListener("abort", abort, { once: true }); if (this.signal.aborted) abort(); return Effect.sync(() => this.signal.removeEventListener("abort", abort)); }); return Effect.raceFirst(program, interrupted); } dispose() { clearTimeout(this.timer); this.controller.abort(new WebFailure("cancelled")); 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>;}
function browserFailure(error: unknown): WebFailure { if (error instanceof WebFailure) return error; if ( error instanceof Error && /^CDP command timed out after \d+ms: /.test(error.message) ) return new WebFailure("timeout"); if (error instanceof BrowserRenderingError && error.status === 429) return new WebFailure("search_blocked", 429); return new WebFailure("browser_unavailable");}
function browserCall<A>(operation: () => PromiseLike<A>) { return Effect.tryPromise({ try: operation, catch: browserFailure });}
export function closeResearchBrowser( browser: BrowserBinding, sessionId: string,) { return Effect.acquireUseRelease( Effect.sync(() => new WebDeadline(5_000)), (cleanup) => cleanup .run(() => deleteBrowserSession( { fetch: (url, init) => browser.fetch(url, { ...init, signal: cleanup.signal }), }, sessionId, ), ) .pipe(Effect.mapError(() => new WebFailure("cleanup_failed"))), (cleanup) => Effect.sync(() => 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, ): Effect.Effect<unknown, WebFailure>;}
// Each invocation owns its CDP connection and remote session. Creation and late// settlement outlive the caller; the parent lease is the durable cleanup owner.export function withResearchBrowser<T, E>( browser: BrowserBinding, deadline: WebDeadline, keepAlive: (promise: Promise<unknown>) => void, run: (browser: ResearchBrowser) => Effect.Effect<T, E>, lease?: BrowserLeaseHooks,): Effect.Effect<T, E | WebFailure> { return Effect.gen(function* () { let session: CdpSession | undefined; let sessionId: string | undefined; let closing: Promise<void> | undefined; const listeners: ((event: ResearchBrowserEvent) => void)[] = []; const close = (): Effect.Effect<void, WebFailure> => Effect.suspend(() => { if (!sessionId) return Effect.void; const id = sessionId; closing ??= Effect.runPromise( Effect.gen(function* () { session?.close(); // Disconnecting CDP alone does not delete the browser. yield* closeResearchBrowser(browser, id); if (lease) { const release = lease.closed(id); keepAlive(release.catch(() => {})); yield* browserCall(() => release).pipe( Effect.timeoutOrElse({ duration: 5_000, orElse: () => Effect.void, }), Effect.catchTag("WebFailure", () => Effect.void), ); } }), ); return browserCall(() => closing!); }); const workBinding: BrowserBinding = { fetch: (url, init) => Effect.runPromise( browserCall(() => browser.fetch(url, { ...init, signal: deadline.signal }), ).pipe( Effect.tap((response) => Effect.sync(() => { 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); }); }), ), ), ), }; const acquire = Effect.gen(function* () { yield* webValidation(() => deadline.remaining()); sessionId = lease?.create ? yield* browserCall(() => lease.create!(deadline.expires)) : (yield* browserCall(() => createBrowserSession(workBinding, { keepAliveMs: 60_000 }), )).sessionId; const id = sessionId; if (lease) { const registration = lease.acquired(id, deadline.expires); keepAlive(registration.catch(() => {})); yield* deadline.run(() => registration); } yield* webValidation(() => deadline.remaining()); const connected = yield* browserCall(() => connectBrowserSession(workBinding, id, deadline.remaining()), ); session = connected; if (deadline.signal.aborted) { connected.close(); return yield* Effect.fail(browserFailure(deadline.signal.reason)); } return connected; }).pipe(Effect.onError(() => close().pipe(Effect.orDie))); return yield* Effect.acquireUseRelease( Effect.sync(() => { const acquisition = Effect.runPromise(acquire); keepAlive( acquisition.then( () => {}, () => {}, ), ); return acquisition; }), (acquisition) => Effect.gen(function* () { const connected = yield* deadline.run(() => acquisition); return yield* deadline.limit( run({ onEvent: (listener) => { listeners.push(listener); }, send: (method, params, target) => deadline.run(() => connected.send(method, params, { timeoutMs: deadline.remaining(), ...(target ? { sessionId: target } : {}), }), ), }), ); }), () => close().pipe(Effect.orDie), ).pipe( Effect.catch((error) => Effect.gen(function* () { yield* webValidation(() => deadline.remaining()); return yield* Effect.fail(error); }), ), ); });}