Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158import type { Agent } from "agents";import { createBrowserSession, type BrowserBinding } from "agents/browser";import * as Effect from "effect/Effect";import { agentCall } from "./agent-io";import { closeResearchBrowser, WebDeadline } from "./browser-session";import type { PersonalAgent } from "./personal-agent";import { WebFailure } from "./web-errors";
type BrowserScheduler = Pick< PersonalAgent, "schedule" | "listSchedules" | "cancelSchedule">;
export class BrowserLeases { constructor( private readonly sql: Agent["sql"], private readonly browser: BrowserBinding, private readonly scheduler: BrowserScheduler, private readonly active: (id: string) => boolean, private readonly waitUntil: (work: Promise<unknown>) => void, ) {}
initialize() { this.sql`CREATE TABLE IF NOT EXISTS flarebot_browser_leases ( session_id TEXT PRIMARY KEY, conversation_id TEXT NOT NULL, expires_at INTEGER NOT NULL, attempts INTEGER NOT NULL DEFAULT 0 )`; }
closeConversation(id: string) { return Effect.suspend(() => Effect.forEach( this.sql<{ session_id: string; }>`SELECT session_id FROM flarebot_browser_leases WHERE conversation_id=${id}`, (lease) => this.expire(lease.session_id), { concurrency: "unbounded", discard: true }, ), ); }
recover() { return Effect.suspend(() => Effect.forEach( this.sql<{ session_id: string; }>`SELECT session_id FROM flarebot_browser_leases WHERE attempts < 3`, (row) => this.expire(row.session_id), { discard: true }, ), ); }
create(conversationId: string, expiresAt: number) { return Effect.acquireUseRelease( Effect.sync( () => new WebDeadline( Math.min(30_000, Math.max(0, expiresAt - Date.now())), ), ), (deadline) => Effect.gen({ self: this }, function* () { if (!this.active(conversationId)) return yield* Effect.fail(new WebFailure("cancelled")); if (Date.now() >= expiresAt) return yield* Effect.fail(new WebFailure("timeout")); // The surviving parent records a late ID even after the child stops waiting. const acquisition = Effect.runPromise( Effect.gen({ self: this }, function* () { const { sessionId } = yield* Effect.tryPromise({ try: () => createBrowserSession( { fetch: (url, init) => this.browser.fetch(url, { ...init, signal: deadline.signal, }), }, { keepAliveMs: 60_000 }, ), catch: () => new WebFailure("browser_unavailable"), }); yield* this.register(conversationId, sessionId, expiresAt).pipe( Effect.catch(() => closeResearchBrowser(this.browser, sessionId).pipe( Effect.andThen(this.release(sessionId)), Effect.catch(() => Effect.void), Effect.andThen(Effect.fail(new WebFailure("cancelled"))), ), ), ); return sessionId; }), ); this.waitUntil(acquisition.catch(() => {})); return yield* deadline.run(() => acquisition); }), (deadline) => Effect.sync(() => deadline.dispose()), ); }
register(conversationId: string, sessionId: string, expiresAt: number) { return Effect.gen({ self: this }, function* () { this .sql`INSERT INTO flarebot_browser_leases (session_id,conversation_id,expires_at) VALUES (${sessionId},${conversationId},${expiresAt})`; yield* agentCall(() => this.scheduler.schedule( new Date(Math.ceil(expiresAt / 1000) * 1000), "expireResearchBrowser", sessionId, { idempotent: true }, ), ); if (!this.active(conversationId) || Date.now() >= expiresAt) return yield* Effect.fail(new WebFailure("cancelled")); }); }
release(sessionId: string) { return Effect.gen({ self: this }, function* () { this .sql`DELETE FROM flarebot_browser_leases WHERE session_id=${sessionId}`; for (const scheduled of yield* agentCall(() => this.scheduler.listSchedules(), )) if ( scheduled.callback === "expireResearchBrowser" && scheduled.payload === sessionId ) yield* agentCall(() => this.scheduler.cancelSchedule(scheduled.id)); }); }
expire(sessionId: string) { return Effect.gen({ self: this }, function* () { const row = this.sql<{ attempts: number; }>`SELECT attempts FROM flarebot_browser_leases WHERE session_id=${sessionId}`[0]; if (!row || row.attempts >= 3) return; this .sql`UPDATE flarebot_browser_leases SET attempts=attempts+1 WHERE session_id=${sessionId}`; yield* closeResearchBrowser(this.browser, sessionId).pipe( Effect.andThen(this.release(sessionId)), Effect.catch(() => row.attempts < 2 ? agentCall(() => this.scheduler.schedule(5, "expireResearchBrowser", sessionId), ).pipe(Effect.asVoid) : Effect.void, ), ); }); }}