diff --git a/packages/effect-acp/src/agent.test.ts b/packages/effect-acp/src/agent.test.ts index 416c4808a3ed..3ba861e70292 100644 --- a/packages/effect-acp/src/agent.test.ts +++ b/packages/effect-acp/src/agent.test.ts @@ -226,7 +226,7 @@ it.effect("effect-acp agent handles core agent requests and outbound client requ }), ); -it.effect("effect-acp agent uses distinct ids for RPC calls and extension requests", () => +it.effect("effect-acp agent keeps RPC ids within int32 and apart from extension ids", () => Effect.gen(function* () { const { stdio, input, output } = yield* makeInMemoryStdio(); const scope = yield* Scope.make(); @@ -275,6 +275,11 @@ it.effect("effect-acp agent uses distinct ids for RPC calls and extension reques : yield* decodedExt(firstOutbound); assert.notEqual(permissionRequest.id, extRequest.id); + // Keep locally generated ids compatible with SDKs that decode signed int32. + assert.typeOf(permissionRequest.id, "number"); + assert.equal(permissionRequest.id, 2 ** 30); + assert.isAtMost(Number(permissionRequest.id), 2 ** 31 - 1); + assert.equal(extRequest.id, 1); yield* Queue.offer( input, @@ -301,10 +306,95 @@ it.effect("effect-acp agent uses distinct ids for RPC calls and extension reques const permission = yield* Fiber.join(permissionFiber); assert.equal(permission.outcome.outcome, "selected"); assert.deepEqual(yield* Fiber.join(extFiber), { ok: true }); + + const nextPermissionFiber = yield* agent.client + .requestPermission(permissionRequest.params) + .pipe(Effect.forkScoped); + const nextPermissionRequest = yield* decodedPermission(yield* Queue.take(output)); + assert.equal(nextPermissionRequest.id, 2 ** 30 + 1); + yield* Queue.offer( + input, + yield* encodeJsonl(RequestPermissionResponse, { + jsonrpc: "2.0", + id: nextPermissionRequest.id, + result: { outcome: { outcome: "cancelled" } }, + }), + ); + assert.equal((yield* Fiber.join(nextPermissionFiber)).outcome.outcome, "cancelled"); }).pipe(Effect.provide(context), Effect.ensuring(Scope.close(scope, Exit.void))); }), ); +it.effect.each([0, 2 ** 30, 2 ** 32])( + "effect-acp agent answers incoming id %s independently of outbound ids", + (requestId) => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const agent = yield* AcpAgent.make(stdio); + const contexts = yield* Ref.make>([]); + yield* agent.handleInitialize((_request, context) => + Ref.update(contexts, (current) => [...current, context]).pipe( + Effect.as({ + protocolVersion: 2, + capabilities: {}, + info: { name: "mock-agent", version: "0.0.0" }, + }), + ), + ); + + // Opposite directions may reuse an id even while a request is pending. + const outboundFiber = yield* agent.client + .requestPermission({ + sessionId: "session-1", + title: "Allow mock action", + subject: { + type: "tool_call", + toolCall: { toolCallId: "tool-1", title: "Allow mock action" }, + }, + options: [{ optionId: "allow", name: "Allow", kind: "allow_once" }], + }) + .pipe(Effect.forkScoped); + const outbound = yield* decodeRequestPermissionRequest(yield* Queue.take(output)); + + yield* Queue.offer( + input, + yield* encodeJsonl(InitializeRequest, { + jsonrpc: "2.0", + id: requestId, + method: "initialize", + params: { + protocolVersion: 2, + capabilities: {}, + info: { name: "test-client", version: "0.0.0" }, + }, + headers: [], + }), + ); + assert.deepEqual(yield* decodeInitializeResponse(yield* Queue.take(output)), { + jsonrpc: "2.0", + id: requestId, + result: { + protocolVersion: 2, + capabilities: {}, + info: { name: "mock-agent", version: "0.0.0" }, + }, + }); + assert.deepEqual(yield* Ref.get(contexts), [ + { requestId: `$t3:jsonrpc:number:${requestId}`, method: "initialize" }, + ]); + + yield* Queue.offer( + input, + yield* encodeJsonl(RequestPermissionResponse, { + jsonrpc: "2.0", + id: outbound.id, + result: { outcome: { outcome: "cancelled" } }, + }), + ); + assert.equal((yield* Fiber.join(outboundFiber)).outcome.outcome, "cancelled"); + }), +); + it.effect("effect-acp agent answers a request whose handler dies with an error for it", () => Effect.gen(function* () { const { stdio, input, output } = yield* makeInMemoryStdio(); diff --git a/packages/effect-acp/src/agent.ts b/packages/effect-acp/src/agent.ts index 2ddf5cc8cc0d..e60059937b7b 100644 --- a/packages/effect-acp/src/agent.ts +++ b/packages/effect-acp/src/agent.ts @@ -436,7 +436,7 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* ( Effect.forkScoped, ); - let nextRpcRequestId = 2 ** 32; + let nextRpcRequestId = AcpProtocol.RPC_REQUEST_ID_START; const rpc = yield* RpcClient.make(AcpRpcs.ClientRpcs, { generateRequestId: () => RpcMessage.RequestId(nextRpcRequestId++), }).pipe(Effect.provideService(RpcClient.Protocol, transport.clientProtocol)); diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index 6209dcffa41f..362c2118e2aa 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -755,7 +755,7 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { }), ); - it.effect("uses distinct ids for RPC calls and extension requests", () => + it.effect("keeps RPC ids within int32 and apart from extension ids", () => Effect.gen(function* () { const { stdio, input, output } = yield* makeInMemoryStdio(); const scope = yield* Scope.make(); @@ -796,6 +796,11 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { : yield* decodedExt(firstOutbound); assert.notEqual(initializeRequest.id, extRequest.id); + // Keep locally generated ids compatible with SDKs that decode signed int32. + assert.typeOf(initializeRequest.id, "number"); + assert.equal(initializeRequest.id, 2 ** 30); + assert.isAtMost(Number(initializeRequest.id), 2 ** 31 - 1); + assert.equal(extRequest.id, 1); yield* Queue.offer( input, @@ -823,6 +828,23 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { yield* Fiber.join(initializeFiber); assert.deepEqual(yield* Fiber.join(extFiber), { ok: true }); + + const sessionFiber = yield* acp.agent + .createSession({ cwd: "/tmp", mcpServers: [] }) + .pipe(Effect.forkScoped); + const sessionRequest = yield* Schema.decodeEffect( + Schema.fromJsonString(jsonRpcRequest("session/new", AcpSchema.NewSessionRequest)), + )(yield* Queue.take(output)); + assert.equal(sessionRequest.id, 2 ** 30 + 1); + yield* Queue.offer( + input, + yield* encodeJsonl(jsonRpcResponse(AcpSchema.NewSessionResponse), { + jsonrpc: "2.0", + id: sessionRequest.id, + result: { sessionId: "session-1" }, + }), + ); + assert.equal((yield* Fiber.join(sessionFiber)).sessionId, "session-1"); yield* Scope.close(scope, Exit.void); }), ); diff --git a/packages/effect-acp/src/client.ts b/packages/effect-acp/src/client.ts index 732e2d26650f..ae158186fca5 100644 --- a/packages/effect-acp/src/client.ts +++ b/packages/effect-acp/src/client.ts @@ -1113,7 +1113,7 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* ( Effect.forkScoped, ); - let nextRpcRequestId = 2 ** 32; + let nextRpcRequestId = AcpProtocol.RPC_REQUEST_ID_START; const rpc = yield* RpcClient.make(AcpRpcs.CompatAgentRpcs, { generateRequestId: () => RpcMessage.RequestId(nextRpcRequestId++), }).pipe(Effect.provideService(RpcClient.Protocol, transport.clientProtocol)); diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index 12e85a9e7542..f026211fd5d4 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -177,6 +177,13 @@ const encodeJsonRpcNotification = Schema.encodeUnknownExit( ), ); +/** + * First id for typed Effect RPC requests. Starts within signed int32 for SDKs that reject + * larger numeric ids (e.g. the Kotlin ACP SDK decodes ids as `Int`), while staying far + * above extension ids, which count up from 1. + */ +export const RPC_REQUEST_ID_START = 2 ** 30; + const isEffectRpcRequestId = (requestId: AcpError.AcpRequestId): boolean => typeof requestId === "number" && Number.isSafeInteger(requestId);