Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254import 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<ReadableStream<Uint8Array>>; closeTemporary(): Promise<boolean>;}
export class ShellLeases { private readonly launches = new Map<string, { stop: () => void }>(); private readonly closing = new Map<string, Promise<boolean>>(); private readonly timers = new Map<string, ReturnType<typeof setTimeout>>(); 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<unknown>) => 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<never>(); 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<boolean, AgentFailure> { 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); }); }}