import type { Agent } from "agents"; import * as Effect from "effect/Effect"; import { SHELL_MAX_ACTIVE, SHELL_MAX_MS } from "../shared/shell"; import { AgentFailure, agentCall, agentValidation } from "./agent-io"; import type { PersonalAgent } from "./personal-agent"; import { shellInput } from "./shell-tool"; type ShellScheduler = Pick< PersonalAgent, "schedule" | "listSchedules" | "cancelSchedule" >; interface ShellSandbox { execTemporary( command: string, expiresAt: number, id: string, ): Promise>; closeTemporary(): Promise; } export class ShellLeases { private readonly launches = new Map void }>(); private readonly closing = new Map>(); private readonly timers = new Map>(); constructor( private readonly sql: Agent["sql"], private readonly scheduler: ShellScheduler, private readonly sandbox: (id: string) => ShellSandbox, private readonly active: (id: string) => boolean, private readonly waitUntil: (work: Promise) => void, ) {} initialize() { this.sql`CREATE TABLE IF NOT EXISTS flarebot_shell_leases ( id TEXT PRIMARY KEY, conversation_id TEXT NOT NULL, expires_at INTEGER NOT NULL, status TEXT NOT NULL )`; } closeConversation(id: string) { return Effect.suspend(() => Effect.forEach( this.sql<{ id: string; }>`SELECT id FROM flarebot_shell_leases WHERE conversation_id=${id}`, (lease) => this.close(lease.id), { concurrency: "unbounded", discard: true }, ), ); } recover() { return Effect.suspend(() => Effect.forEach( this.sql<{ id: string }>`SELECT id FROM flarebot_shell_leases`, (row) => this.close(row.id), { discard: true }, ), ); } reserve(conversationId: string, timeoutMs: number) { return Effect.gen({ self: this }, function* () { if (!this.active(conversationId)) return yield* Effect.fail( new AgentFailure({ message: "Conversation not found" }), ); if ( !Number.isInteger(timeoutMs) || timeoutMs < 1000 || timeoutMs > SHELL_MAX_MS ) return yield* Effect.fail( new AgentFailure({ message: "Invalid shell lifetime" }), ); if ( this.sql`SELECT id FROM flarebot_shell_leases`.length >= SHELL_MAX_ACTIVE ) return yield* Effect.fail( new AgentFailure({ message: "Shell capacity unavailable" }), ); const id = `shell-${crypto.randomUUID()}`, expiresAt = Date.now() + timeoutMs; this .sql`INSERT INTO flarebot_shell_leases VALUES (${id},${conversationId},${expiresAt},'reserved')`; yield* agentCall(() => this.scheduler.schedule( new Date(Math.ceil(expiresAt / 1000) * 1000), "expireShellWorkspace", id, { idempotent: true }, ), ); this.timers.set( id, setTimeout( () => this.waitUntil(Effect.runPromise(this.close(id))), Math.max(0, expiresAt - Date.now()), ), ); return { id, expiresAt }; }); } launch(id: string, command: string) { return Effect.gen({ self: this }, function* () { yield* agentValidation(() => shellInput.parse({ command })); const row = this.sql<{ conversation_id: string; expires_at: number; status: string; }>`SELECT * FROM flarebot_shell_leases WHERE id=${id}`[0]; if ( !row || row.status !== "reserved" || !this.active(row.conversation_id) || Date.now() >= row.expires_at ) return yield* Effect.fail( new AgentFailure({ message: "Shell workspace unavailable" }), ); this .sql`UPDATE flarebot_shell_leases SET status='launching' WHERE id=${id}`; const interrupted = Promise.withResolvers(); this.launches.set(id, { stop: () => interrupted.reject( new AgentFailure({ message: "Shell interrupted" }), ), }); // A late native launch must finish under the parent and destroy again after // an earlier close. Interruption of the caller alone cannot revoke an RPC. const launch = Effect.runPromise( Effect.gen({ self: this }, function* () { const stream = yield* agentCall( () => this.sandbox(id).execTemporary(command, row.expires_at, id), "Shell launch unavailable", ); const current = this.sql<{ status: string; }>`SELECT status FROM flarebot_shell_leases WHERE id=${id}`[0]; if ( !current || current.status === "closing" || !this.active(row.conversation_id) || Date.now() >= row.expires_at ) { yield* agentCall(() => stream.cancel()).pipe( Effect.catchTag("AgentFailure", () => Effect.void), ); return yield* Effect.fail( new AgentFailure({ message: "Shell interrupted" }), ); } this .sql`UPDATE flarebot_shell_leases SET status='running' WHERE id=${id}`; return stream; }).pipe( Effect.catchTag("AgentFailure", () => Effect.gen({ self: this }, function* () { this .sql`UPDATE flarebot_shell_leases SET status='closing' WHERE id=${id}`; return yield* Effect.fail( new AgentFailure({ message: "Shell launch unavailable" }), ); }), ), Effect.ensuring( Effect.gen({ self: this }, function* () { this.launches.delete(id); if ( this .sql`SELECT id FROM flarebot_shell_leases WHERE id=${id} AND status='closing'` .length ) { const earlier = this.closing.get(id); if (earlier) yield* agentCall(() => earlier); yield* this.close(id); } }).pipe(Effect.orDie), ), ), ); this.waitUntil( launch.then( () => {}, () => {}, ), ); return yield* Effect.raceFirst( agentCall(() => launch, "Shell launch unavailable"), agentCall(() => interrupted.promise, "Shell interrupted"), ); }); } close(id: string): Effect.Effect { return Effect.suspend(() => { if (!this.sql`SELECT id FROM flarebot_shell_leases WHERE id=${id}`.length) return Effect.succeed(true); this .sql`UPDATE flarebot_shell_leases SET status='closing' WHERE id=${id}`; this.launches.get(id)?.stop(); clearTimeout(this.timers.get(id) ?? null); this.timers.delete(id); const existing = this.closing.get(id); if (existing) return agentCall(() => existing); const launchPending = this.launches.has(id); const closing = Effect.runPromise( Effect.gen({ self: this }, function* () { const destroy = this.sandbox(id).closeTemporary(); this.waitUntil(destroy.catch(() => {})); const destroyed = yield* agentCall(() => destroy).pipe( Effect.timeoutOrElse({ duration: 2000, orElse: () => Effect.succeed(false), }), ); if (!destroyed || launchPending || this.launches.has(id)) return false; this.sql`DELETE FROM flarebot_shell_leases WHERE id=${id}`; for (const scheduled of yield* agentCall(() => this.scheduler.listSchedules(), )) if ( scheduled.callback === "expireShellWorkspace" && scheduled.payload === id ) yield* agentCall(() => this.scheduler.cancelSchedule(scheduled.id), ); return true; }).pipe( Effect.catchTag("AgentFailure", () => Effect.succeed(false)), Effect.ensuring( Effect.gen({ self: this }, function* () { if ( this.sql`SELECT id FROM flarebot_shell_leases WHERE id=${id}` .length ) yield* agentCall(() => this.scheduler.schedule(5, "expireShellWorkspace", id), ); }).pipe(Effect.orDie), ), Effect.ensuring(Effect.sync(() => this.closing.delete(id))), ), ); this.closing.set(id, closing); this.waitUntil(closing); return agentCall(() => closing); }); } }