From 5cfffb8ed9be94367f1353a801cbb428311ad7c1 Mon Sep 17 00:00:00 2001 From: Kevin Rajan Date: Mon, 28 Sep 2026 02:01:33 -0500 Subject: [PATCH] [muse] fix(acp): bound extension requests with per-request timeout option The ACP patched protocol's sendRequest awaited its response Deferred with no deadline. A wedged ACP agent that stays alive but never answers an extension request (e.g. cursor/list_available_models during Cursor model discovery) hung the caller forever; termination only rescued the dead-process case. Mirrors the codex app-server transport fix: request() now accepts an optional timeout (default unbounded, existing behavior). On expiry the pending entry is dropped (late responses are ignored by resolveExtPending) and the request fails with AcpRequestTimeoutError. The runtime's request() passes the option through, and Cursor model discovery is now bounded at 30s. Tests: protocol.test.ts gains a TestClock test (timeout fires, late response ignored, next request routes by id) plus a default-unbounded control. Red-on-base verified (hangs, 15s vitest timeout); 29/29 green with the fix. Authored with AI assistance (Muse, Meta's Muse Spark) under the contributor's direction. --- .../src/provider/Layers/CursorProvider.ts | 10 +- .../src/provider/acp/AcpSessionRuntime.ts | 7 +- packages/effect-acp/src/errors.ts | 13 +++ packages/effect-acp/src/protocol.test.ts | 103 ++++++++++++++++++ packages/effect-acp/src/protocol.ts | 38 ++++++- 5 files changed, 165 insertions(+), 6 deletions(-) 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 {