diff --git a/app/server/modules/agents/controller/session.ts b/app/server/modules/agents/controller/session.ts index cd719efd..0187d112 100644 --- a/app/server/modules/agents/controller/session.ts +++ b/app/server/modules/agents/controller/session.ts @@ -1,4 +1,4 @@ -import { Effect, Queue, Ref, type Scope } from "effect"; +import { Deferred, Effect, Queue, Ref, type Scope } from "effect"; import type { AgentKind } from "../../../db/schema"; import { createControllerMessage, @@ -29,10 +29,8 @@ type SessionState = { lastPongAt: number | null; }; -type PendingVolumeCommand = { - resolve: (payload: VolumeCommandResponsePayload) => void; - reject: (error: Error) => void; - timeout: ReturnType; +type InFlightVolumeCommand = { + deferred: Deferred.Deferred; }; export type ControllerAgentSessionEvent = @@ -58,7 +56,7 @@ export const createControllerAgentSession = ( Effect.gen(function* () { let isClosed = false; const outboundQueue = yield* Queue.bounded(64); - const pendingVolumeCommands = new Map(); + const inFlightVolumeCommands = yield* Ref.make(new Map()); const state = yield* Ref.make({ isReady: false, lastSeenAt: null, @@ -77,18 +75,33 @@ export const createControllerAgentSession = ( const updateState = (update: (current: SessionState) => SessionState) => Ref.update(state, update); - const rejectPendingVolumeCommands = () => { - for (const [commandId, pending] of pendingVolumeCommands) { - clearTimeout(pending.timeout); - pending.reject(new Error(`Agent session closed before volume command ${commandId} completed`)); + const setInFlightVolumeCommand = (commandId: string, inFlight: InFlightVolumeCommand) => + Ref.update(inFlightVolumeCommands, (current) => new Map(current).set(commandId, inFlight)); + + const removeInFlightVolumeCommand = (commandId: string) => + Ref.modify(inFlightVolumeCommands, (current) => { + const inFlight = current.get(commandId) ?? null; + const next = new Map(current); + next.delete(commandId); + return [inFlight, next]; + }); + + const rejectInFlightVolumeCommands = Effect.gen(function* () { + const inFlightCommands = yield* Ref.get(inFlightVolumeCommands); + yield* Ref.set(inFlightVolumeCommands, new Map()); + + for (const [commandId, inFlight] of inFlightCommands) { + yield* Deferred.fail( + inFlight.deferred, + new Error(`Agent session closed before volume command ${commandId} completed`), + ); } - pendingVolumeCommands.clear(); - }; + }); const releaseSession = Effect.gen(function* () { const disconnectedAt = Date.now(); yield* updateState((current) => ({ ...current, isReady: false, lastSeenAt: disconnectedAt })); - yield* Effect.sync(rejectPendingVolumeCommands); + yield* rejectInFlightVolumeCommands; yield* onEvent({ type: "agent.disconnected" }); yield* Queue.shutdown(outboundQueue); @@ -152,17 +165,16 @@ export const createControllerAgentSession = ( return yield* Effect.never; }); - const handleVolumeCommandResult = (payload: VolumeCommandResponsePayload) => { - const pending = pendingVolumeCommands.get(payload.commandId); - if (!pending) { - logger.warn(`Received response for unknown volume command ${payload.commandId}`); - return; - } + const handleVolumeCommandResult = (payload: VolumeCommandResponsePayload) => + Effect.gen(function* () { + const inFlight = yield* removeInFlightVolumeCommand(payload.commandId); + if (!inFlight) { + yield* logger.effect.warn(`Received response for unknown volume command ${payload.commandId}`); + return; + } - pendingVolumeCommands.delete(payload.commandId); - clearTimeout(pending.timeout); - pending.resolve(payload); - }; + yield* Deferred.succeed(inFlight.deferred, payload); + }); const handleAgentMessage = (message: AgentMessage) => Effect.gen(function* () { @@ -181,7 +193,7 @@ export const createControllerAgentSession = ( break; } case "volume.commandResult": { - yield* Effect.sync(() => handleVolumeCommandResult(message.payload)); + yield* handleVolumeCommandResult(message.payload); break; } default: { @@ -212,37 +224,26 @@ export const createControllerAgentSession = ( }, sendBackup: (payload) => offerOutbound(createControllerMessage("backup.run", payload)), sendBackupCancel: (payload) => offerOutbound(createControllerMessage("backup.cancel", payload)), - runVolumeCommand: (command) => { - return Effect.tryPromise({ - try: async () => { - const commandId = Bun.randomUUIDv7(); - const response = new Promise((resolve, reject) => { - const timeout = setTimeout(() => { - pendingVolumeCommands.delete(commandId); - reject(new Error(`Volume command ${command.name} timed out`)); - }, 60_000); + runVolumeCommand: (command) => + Effect.gen(function* () { + const commandId = Bun.randomUUIDv7(); + const deferred = yield* Deferred.make(); + yield* setInFlightVolumeCommand(commandId, { deferred }); - pendingVolumeCommands.set(commandId, { resolve, reject, timeout }); - }); + const queued = yield* offerOutbound(createControllerMessage("volume.command", { commandId, command })); + if (!queued) { + yield* removeInFlightVolumeCommand(commandId); + return yield* Effect.fail(new Error(`Failed to queue volume command ${command.name}`)); + } - const queued = await Effect.runPromise( - offerOutbound(createControllerMessage("volume.command", { commandId, command })), - ); - - if (!queued) { - const pending = pendingVolumeCommands.get(commandId); - if (pending) { - clearTimeout(pending.timeout); - pendingVolumeCommands.delete(commandId); - } - throw new Error(`Failed to queue volume command ${command.name}`); - } - - return response; - }, - catch: (error) => (error instanceof Error ? error : new Error(toMessage(error))), - }); - }, + return yield* Deferred.await(deferred).pipe( + Effect.timeoutFail({ + duration: "60 seconds", + onTimeout: () => new Error(`Volume command ${command.name} timed out`), + }), + Effect.ensuring(removeInFlightVolumeCommand(commandId)), + ); + }), isReady: () => Ref.get(state).pipe(Effect.map((current) => current.isReady)), run, };