Skip to content
46 changes: 19 additions & 27 deletions packages/effect-acp/src/_internal/shared.ts
Original file line number Diff line number Diff line change
@@ -1,30 +1,29 @@
import * as Effect from "effect/Effect";
import * as Schema from "effect/Schema";
import * as SchemaIssue from "effect/SchemaIssue";
import { RpcClientError } from "effect/unstable/rpc";

import * as AcpSchema from "../_generated/schema.gen.ts";
import * as AcpError from "../errors.ts";
const isError = Schema.is(AcpSchema.Error);
const isAcpRequestError = Schema.is(AcpError.AcpRequestError);

const formatSchemaIssue = SchemaIssue.makeFormatterDefault();

export const callRpc = <A>(
method: string,
effect: Effect.Effect<A, RpcClientError.RpcClientError | AcpSchema.Error>,
): Effect.Effect<A, AcpError.AcpError> =>
effect.pipe(
Effect.catchTag("RpcClientError", (error) =>
Effect.fail(
new AcpError.AcpTransportError({
detail: error.message,
cause: error,
}),
),
),
Effect.catchIf(isError, (error) =>
Effect.fail(AcpError.AcpRequestError.fromProtocolError(error)),
),
Effect.catchTags({
RpcClientError: (cause) =>
Effect.fail(
new AcpError.AcpTransportError({
operation: "call-rpc",
method,
cause,
}),
),
}),
);

export const runHandler = Effect.fnUntraced(function* <A, B>(
Expand All @@ -37,9 +36,7 @@ export const runHandler = Effect.fnUntraced(function* <A, B>(
}
return yield* handler(payload).pipe(
Effect.mapError((error) =>
isAcpRequestError(error)
Comment thread
cursor[bot] marked this conversation as resolved.
? error.toProtocolError()
: AcpError.AcpRequestError.internalError(error.message).toProtocolError(),
AcpError.AcpRequestError.fromCoreHandlerError(error, method).toProtocolError(),
),
);
});
Expand All @@ -51,12 +48,7 @@ export function decodeExtRequestRegistration<A, I>(
) {
return (params: unknown): Effect.Effect<unknown, AcpError.AcpError> =>
Schema.decodeUnknownEffect(payload)(params).pipe(
Effect.mapError((error) =>
AcpError.AcpRequestError.invalidParams(
`Invalid ${method} payload: ${formatSchemaIssue(error.issue)}`,
{ issue: error.issue },
),
),
Effect.mapError((error) => AcpError.AcpRequestError.invalidExtensionPayload(method, error)),
Effect.flatMap((decoded) => handler(decoded)),
);
}
Expand All @@ -68,12 +60,12 @@ export function decodeExtNotificationRegistration<A, I>(
) {
return (params: unknown): Effect.Effect<void, AcpError.AcpError> =>
Schema.decodeUnknownEffect(payload)(params).pipe(
Effect.mapError(
(error) =>
new AcpError.AcpProtocolParseError({
detail: `Invalid ${method} notification payload: ${formatSchemaIssue(error.issue)}`,
cause: error,
}),
Effect.mapError((error) =>
AcpError.AcpProtocolParseError.fromSchemaError(
"decode-notification-payload",
method,
error,
),
),
Effect.flatMap((decoded) => handler(decoded)),
);
Expand Down
2 changes: 1 addition & 1 deletion packages/effect-acp/src/_internal/stdio.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ export const makeTerminationError = (
Effect.match(handle.exitCode, {
onFailure: (cause) =>
new AcpError.AcpTransportError({
detail: "Failed to determine ACP process exit status",
operation: "read-process-exit-status",
cause,
}),
onSuccess: (code) => new AcpError.AcpProcessExitedError({ code }),
Expand Down
94 changes: 54 additions & 40 deletions packages/effect-acp/src/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -288,12 +288,12 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* (
notification.method === AGENT_METHODS.session_cancel
) {
return decodeCancelNotification(notification.params).pipe(
Effect.mapError(
(error) =>
new AcpError.AcpProtocolParseError({
detail: `Invalid ${AGENT_METHODS.session_cancel} notification payload`,
cause: error,
}),
Effect.mapError((error) =>
AcpError.AcpProtocolParseError.fromSchemaError(
"decode-notification-payload",
AGENT_METHODS.session_cancel,
error,
),
),
Effect.flatMap((decoded) =>
Effect.forEach(cancelHandlers, (handler) => handler(decoded), { discard: true }),
Expand Down Expand Up @@ -376,41 +376,55 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* (
},
client: {
requestPermission: (payload) =>
callRpc(rpc[CLIENT_METHODS.session_request_permission](payload)),
elicit: (payload) => callRpc(rpc[CLIENT_METHODS.session_elicitation](payload)),
readTextFile: (payload) => callRpc(rpc[CLIENT_METHODS.fs_read_text_file](payload)),
writeTextFile: (payload) => callRpc(rpc[CLIENT_METHODS.fs_write_text_file](payload)),
callRpc(
CLIENT_METHODS.session_request_permission,
rpc[CLIENT_METHODS.session_request_permission](payload),
),
elicit: (payload) =>
callRpc(
CLIENT_METHODS.session_elicitation,
rpc[CLIENT_METHODS.session_elicitation](payload),
),
readTextFile: (payload) =>
callRpc(CLIENT_METHODS.fs_read_text_file, rpc[CLIENT_METHODS.fs_read_text_file](payload)),
writeTextFile: (payload) =>
callRpc(CLIENT_METHODS.fs_write_text_file, rpc[CLIENT_METHODS.fs_write_text_file](payload)),
createTerminal: (payload) =>
callRpc(rpc[CLIENT_METHODS.terminal_create](payload)).pipe(
Effect.map((response) =>
AcpTerminal.makeTerminal({
sessionId: payload.sessionId,
terminalId: response.terminalId,
output: callRpc(
rpc[CLIENT_METHODS.terminal_output]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
waitForExit: callRpc(
rpc[CLIENT_METHODS.terminal_wait_for_exit]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
kill: callRpc(
rpc[CLIENT_METHODS.terminal_kill]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
release: callRpc(
rpc[CLIENT_METHODS.terminal_release]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
}),
callRpc(CLIENT_METHODS.terminal_create, rpc[CLIENT_METHODS.terminal_create](payload)).pipe(
Effect.map(
(response) =>
({
sessionId: payload.sessionId,
terminalId: response.terminalId,
output: callRpc(
CLIENT_METHODS.terminal_output,
rpc[CLIENT_METHODS.terminal_output]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
waitForExit: callRpc(
CLIENT_METHODS.terminal_wait_for_exit,
rpc[CLIENT_METHODS.terminal_wait_for_exit]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
kill: callRpc(
CLIENT_METHODS.terminal_kill,
rpc[CLIENT_METHODS.terminal_kill]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
release: callRpc(
CLIENT_METHODS.terminal_release,
rpc[CLIENT_METHODS.terminal_release]({
sessionId: payload.sessionId,
terminalId: response.terminalId,
}),
),
}) satisfies AcpTerminal.AcpTerminal,
),
),
sessionUpdate: (payload) => transport.notify(CLIENT_METHODS.session_update, payload),
Expand Down
6 changes: 3 additions & 3 deletions packages/effect-acp/src/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,7 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => {
);

it.effect(
"returns formatted invalid params when a typed extension request payload is wrong",
"returns structured invalid params without exposing values from typed extension request payloads",
() =>
Effect.gen(function* () {
const handle = yield* makeHandle({ ACP_MOCK_BAD_TYPED_REQUEST: "1" });
Expand Down Expand Up @@ -213,8 +213,8 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => {
assert.fail("Expected prompt to fail for invalid typed extension payload");
}
const rendered = Cause.pretty(result.cause);
assert.include(rendered, "Invalid x/typed_request payload:");
assert.include(rendered, "Expected string, got 123");
assert.include(rendered, "Invalid payload for ACP extension method 'x/typed_request'.");
assert.notInclude(rendered, "Expected string, got 123");
}),
);

Expand Down
37 changes: 25 additions & 12 deletions packages/effect-acp/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -462,19 +462,32 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* (
notify: transport.notify,
},
agent: {
initialize: (payload) => callRpc(rpc[AGENT_METHODS.initialize](payload)),
authenticate: (payload) => callRpc(rpc[AGENT_METHODS.authenticate](payload)),
logout: (payload) => callRpc(rpc[AGENT_METHODS.logout](payload)),
createSession: (payload) => callRpc(rpc[AGENT_METHODS.session_new](payload)),
loadSession: (payload) => callRpc(rpc[AGENT_METHODS.session_load](payload)),
listSessions: (payload) => callRpc(rpc[AGENT_METHODS.session_list](payload)),
forkSession: (payload) => callRpc(rpc[AGENT_METHODS.session_fork](payload)),
resumeSession: (payload) => callRpc(rpc[AGENT_METHODS.session_resume](payload)),
closeSession: (payload) => callRpc(rpc[AGENT_METHODS.session_close](payload)),
setSessionModel: (payload) => callRpc(rpc[AGENT_METHODS.session_set_model](payload)),
initialize: (payload) =>
callRpc(AGENT_METHODS.initialize, rpc[AGENT_METHODS.initialize](payload)),
authenticate: (payload) =>
callRpc(AGENT_METHODS.authenticate, rpc[AGENT_METHODS.authenticate](payload)),
logout: (payload) => callRpc(AGENT_METHODS.logout, rpc[AGENT_METHODS.logout](payload)),
createSession: (payload) =>
callRpc(AGENT_METHODS.session_new, rpc[AGENT_METHODS.session_new](payload)),
loadSession: (payload) =>
callRpc(AGENT_METHODS.session_load, rpc[AGENT_METHODS.session_load](payload)),
listSessions: (payload) =>
callRpc(AGENT_METHODS.session_list, rpc[AGENT_METHODS.session_list](payload)),
forkSession: (payload) =>
callRpc(AGENT_METHODS.session_fork, rpc[AGENT_METHODS.session_fork](payload)),
resumeSession: (payload) =>
callRpc(AGENT_METHODS.session_resume, rpc[AGENT_METHODS.session_resume](payload)),
closeSession: (payload) =>
callRpc(AGENT_METHODS.session_close, rpc[AGENT_METHODS.session_close](payload)),
setSessionModel: (payload) =>
callRpc(AGENT_METHODS.session_set_model, rpc[AGENT_METHODS.session_set_model](payload)),
setSessionConfigOption: (payload) =>
callRpc(rpc[AGENT_METHODS.session_set_config_option](payload)),
prompt: (payload) => callRpc(rpc[AGENT_METHODS.session_prompt](payload)),
callRpc(
AGENT_METHODS.session_set_config_option,
rpc[AGENT_METHODS.session_set_config_option](payload),
),
prompt: (payload) =>
callRpc(AGENT_METHODS.session_prompt, rpc[AGENT_METHODS.session_prompt](payload)),
cancel: (payload) => transport.notify(AGENT_METHODS.session_cancel, payload),
},
handleRequestPermission: (handler) =>
Expand Down
Loading
Loading