Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
92 changes: 91 additions & 1 deletion packages/effect-acp/src/agent.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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,
Expand All @@ -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<Array<AcpProtocol.AcpRequestContext>>([]);
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();
Expand Down
2 changes: 1 addition & 1 deletion packages/effect-acp/src/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
24 changes: 23 additions & 1 deletion packages/effect-acp/src/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
}),
);
Expand Down
2 changes: 1 addition & 1 deletion packages/effect-acp/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
7 changes: 7 additions & 0 deletions packages/effect-acp/src/protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
Loading