diff --git a/apps/server/src/provider/Layers/CursorProvider.ts b/apps/server/src/provider/Layers/CursorProvider.ts index dfaa8ffe338e..c8622fee29e6 100644 --- a/apps/server/src/provider/Layers/CursorProvider.ts +++ b/apps/server/src/provider/Layers/CursorProvider.ts @@ -653,7 +653,15 @@ const discoverCursorModelsViaListAvailableModels = ( (acp) => Effect.gen(function* () { yield* acp.start(); - const response = yield* acp.request("cursor/list_available_models", {}); + // A Cursor agent CLI that stays alive but stops answering extension + // requests must not hang model discovery forever. + const response = yield* acp.request( + "cursor/list_available_models", + {}, + { + timeout: "30 seconds", + }, + ); const decoded = yield* decodeCursorListAvailableModelsResponse(response); return buildCursorDiscoveredModelsFromAvailableModelsResponse(decoded); }), diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 0d0b81e910a3..ec70a5f1c0bb 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -285,6 +285,7 @@ export class AcpSessionRuntime extends Context.Service< readonly request: ( method: string, payload: unknown, + options?: EffectAcpProtocol.AcpPatchedRequestOptions, ) => Effect.Effect; /** * Sends a generic ACP extension notification. @@ -1124,9 +1125,11 @@ export const make = ( ); }), ), - request: (method, payload) => + request: (method, payload, options) => ensureConnected.pipe( - Effect.andThen(runLoggedRequest(method, payload, acp.raw.request(method, payload))), + Effect.andThen( + runLoggedRequest(method, payload, acp.raw.request(method, payload, options)), + ), ), notify: (method, payload) => ensureConnected.pipe(Effect.andThen(acp.raw.notify(method, payload))), diff --git a/packages/effect-acp/src/errors.ts b/packages/effect-acp/src/errors.ts index 01f296bb8cb6..38432a7d2805 100644 --- a/packages/effect-acp/src/errors.ts +++ b/packages/effect-acp/src/errors.ts @@ -368,6 +368,18 @@ export class AcpRequestError extends Schema.TaggedError()("AcpR } } +export class AcpRequestTimeoutError extends Schema.TaggedError()( + "AcpRequestTimeoutError", + { + method: Schema.String, + requestId: AcpRequestId, + }, +) { + override get message() { + return `ACP request '${this.method}' (id ${this.requestId}) timed out waiting for a response.`; + } +} + export const AcpError = Schema.Union([ AcpRequestError, AcpSpawnError, @@ -375,6 +387,7 @@ export const AcpError = Schema.Union([ AcpProtocolParseError, AcpTransportError, AcpInputStreamEndedError, + AcpRequestTimeoutError, ]); export type AcpError = typeof AcpError.Type; diff --git a/packages/effect-acp/src/protocol.test.ts b/packages/effect-acp/src/protocol.test.ts index cef04ffcd498..22a5b8d2bab5 100644 --- a/packages/effect-acp/src/protocol.test.ts +++ b/packages/effect-acp/src/protocol.test.ts @@ -11,6 +11,7 @@ import { ChildProcess, ChildProcessSpawner } from "effect/unstable/process"; import { it, assert } from "@effect/vitest"; import * as NodeServices from "@effect/platform-node/NodeServices"; +import * as TestClock from "effect/testing/TestClock"; import * as AcpSchema from "./_generated/schema.gen.ts"; import * as AcpProtocol from "./protocol.ts"; @@ -766,5 +767,107 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { assert.equal(yield* Queue.size(output), 0); }), ); + it.effect("fails ext requests that never receive a response when a timeout is set", () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const transport = yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + }); + + const fiber = yield* Effect.forkScoped( + transport + .request("x/test", { hello: "world" }, { timeout: "50 millis" }) + .pipe(Effect.catch((error) => Effect.succeed(error))), + ); + const outbound = yield* Queue.take(output); + assert.deepEqual(yield* decodeExtRequest(outbound), { + jsonrpc: "2.0", + id: 1, + method: "x/test", + params: { + hello: "world", + }, + headers: [], + }); + + // Advance the test clock past the deadline: the request itself fails + // with a timeout error instead of hanging on the wedged peer. + yield* TestClock.adjust("100 millis"); + const error = yield* Fiber.join(fiber); + assert.instanceOf(error, AcpError.AcpRequestTimeoutError); + assert.strictEqual(error.method, "x/test"); + assert.strictEqual(error.requestId, "1"); + assert.match(error.message, /timed out/); + + // A late response for the timed-out request is ignored, and the next + // request still routes by id — the pending entry was dropped. + yield* Queue.offer( + input, + yield* encodeJsonl(ExtResponse, { + jsonrpc: "2.0", + id: 1, + result: { ok: false }, + }), + ); + const pending = yield* transport + .request("x/test", { hello: "world" }) + .pipe(Effect.forkScoped); + const outbound2 = yield* Queue.take(output); + assert.deepEqual(yield* decodeExtRequest(outbound2), { + jsonrpc: "2.0", + id: 2, + method: "x/test", + params: { + hello: "world", + }, + headers: [], + }); + yield* Queue.offer( + input, + yield* encodeJsonl(ExtResponse, { + jsonrpc: "2.0", + id: 2, + result: { ok: true }, + }), + ); + assert.deepEqual(yield* Fiber.join(pending), { ok: true }); + }), + ); + + it.effect("keeps ext requests unbounded by default", () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const transport = yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + }); + + // No timeout option: the existing behavior — waits until the peer + // answers, however long that takes. + const pending = yield* transport + .request("x/test", { hello: "world" }) + .pipe(Effect.forkScoped); + const outbound = yield* Queue.take(output); + assert.deepEqual(yield* decodeExtRequest(outbound), { + jsonrpc: "2.0", + id: 1, + method: "x/test", + params: { + hello: "world", + }, + headers: [], + }); + yield* Queue.offer( + input, + yield* encodeJsonl(ExtResponse, { + jsonrpc: "2.0", + id: 1, + result: { ok: true }, + }), + ); + assert.deepEqual(yield* Fiber.join(pending), { ok: true }); + }), + ); } }); diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index a84c2f5a57b4..bddc50e2766c 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -1,7 +1,9 @@ import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; import * as Exit from "effect/Exit"; +import * as Option from "effect/Option"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import type * as PlatformError from "effect/PlatformError"; @@ -68,11 +70,21 @@ export interface AcpPatchedProtocolOptions { readonly onTermination?: (error: AcpError.AcpError) => Effect.Effect; } +export interface AcpPatchedRequestOptions { + // Bounds how long the request waits for a response. A wedged peer must not + // hold the caller forever. Defaults to unbounded (existing behavior). + readonly timeout?: Duration.Input; +} + export interface AcpPatchedProtocol { readonly clientProtocol: RpcClient.Protocol["Service"]; readonly serverProtocol: RpcServer.Protocol["Service"]; readonly incoming: Stream.Stream; - readonly request: (method: string, payload: unknown) => Effect.Effect; + readonly request: ( + method: string, + payload: unknown, + options?: AcpPatchedRequestOptions, + ) => Effect.Effect; readonly notify: (method: string, payload: unknown) => Effect.Effect; } @@ -585,7 +597,11 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi yield* Queue.offer(outgoing, encoded); }); - const sendRequest = Effect.fn("sendRequest")(function* (method: string, payload: unknown) { + const sendRequest = Effect.fn("sendRequest")(function* ( + method: string, + payload: unknown, + options?: AcpPatchedRequestOptions, + ) { yield* ensureActive; const requestId = yield* Ref.modify( nextRequestId, @@ -602,9 +618,25 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi payload, headers: [], }).pipe(Effect.tapError(() => removeExtPending(requestId))); - return yield* Deferred.await(deferred).pipe( + if (options?.timeout === undefined) { + return yield* Deferred.await(deferred).pipe( + Effect.onInterrupt(() => removeExtPending(requestId)), + ); + } + const response = yield* Deferred.await(deferred).pipe( Effect.onInterrupt(() => removeExtPending(requestId)), + Effect.timeoutOption(options.timeout), ); + if (Option.isNone(response)) { + // The peer never answered. Drop the pending entry (a late response + // is ignored by resolveExtPending) and fail instead of hanging. + yield* removeExtPending(requestId); + return yield* new AcpError.AcpRequestTimeoutError({ + method, + requestId: String(requestId), + }); + } + return response.value; }); return {