diff --git a/packages/effect-acp/src/_internal/shared.ts b/packages/effect-acp/src/_internal/shared.ts index 937d931c404..7e43bbf8831 100644 --- a/packages/effect-acp/src/_internal/shared.ts +++ b/packages/effect-acp/src/_internal/shared.ts @@ -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 = ( + method: string, effect: Effect.Effect, ): Effect.Effect => 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* ( @@ -37,9 +36,7 @@ export const runHandler = Effect.fnUntraced(function* ( } return yield* handler(payload).pipe( Effect.mapError((error) => - isAcpRequestError(error) - ? error.toProtocolError() - : AcpError.AcpRequestError.internalError(error.message).toProtocolError(), + AcpError.AcpRequestError.fromCoreHandlerError(error, method).toProtocolError(), ), ); }); @@ -51,12 +48,7 @@ export function decodeExtRequestRegistration( ) { return (params: unknown): Effect.Effect => 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)), ); } @@ -68,12 +60,12 @@ export function decodeExtNotificationRegistration( ) { return (params: unknown): Effect.Effect => 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)), ); diff --git a/packages/effect-acp/src/_internal/stdio.ts b/packages/effect-acp/src/_internal/stdio.ts index 8ddb4d37d0f..393a1c591cb 100644 --- a/packages/effect-acp/src/_internal/stdio.ts +++ b/packages/effect-acp/src/_internal/stdio.ts @@ -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 }), diff --git a/packages/effect-acp/src/agent.ts b/packages/effect-acp/src/agent.ts index 5cad53c3d12..307028b0a80 100644 --- a/packages/effect-acp/src/agent.ts +++ b/packages/effect-acp/src/agent.ts @@ -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 }), @@ -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), diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index aca87d45c62..c732f80ef35 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -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" }); @@ -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"); }), ); diff --git a/packages/effect-acp/src/client.ts b/packages/effect-acp/src/client.ts index 6f3d6a0c9f8..61b3d71b49d 100644 --- a/packages/effect-acp/src/client.ts +++ b/packages/effect-acp/src/client.ts @@ -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) => diff --git a/packages/effect-acp/src/errors.test.ts b/packages/effect-acp/src/errors.test.ts new file mode 100644 index 00000000000..5187fabf5d2 --- /dev/null +++ b/packages/effect-acp/src/errors.test.ts @@ -0,0 +1,144 @@ +import { describe, expect, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Schema from "effect/Schema"; +import * as RpcClientError from "effect/unstable/rpc/RpcClientError"; + +import * as AcpSchema from "./_generated/schema.gen.ts"; +import { callRpc, runHandler } from "./_internal/shared.ts"; +import * as AcpError from "./errors.ts"; + +const decodeNestedNumberPayload = Schema.decodeUnknownEffect( + Schema.Struct({ profile: Schema.Struct({ token: Schema.Number }) }), +); +const encodeUnknownJson = Schema.encodeSync(Schema.UnknownFromJsonString); + +describe("effect-acp errors", () => { + it.effect("retains RPC method and cause without deriving the message from the cause", () => { + const rootCause = new Error("connection details that must not become the public message"); + const failure = new RpcClientError.RpcClientError({ + reason: new RpcClientError.RpcClientDefect({ + message: rootCause.message, + cause: rootCause, + }), + }); + + return Effect.gen(function* () { + const error = yield* callRpc("session/new", Effect.fail(failure)).pipe(Effect.flip); + + expect(error).toMatchObject({ + _tag: "AcpTransportError", + operation: "call-rpc", + method: "session/new", + cause: failure, + }); + expect(error.message).toBe("ACP transport operation call-rpc failed for method session/new."); + expect(error.message).not.toContain(rootCause.message); + }); + }); + + it.effect("preserves protocol request errors as request errors", () => { + const failure = AcpSchema.Error.make({ + code: -32602, + message: "Invalid params", + data: { field: "sessionId" }, + }); + + return Effect.gen(function* () { + const error = yield* callRpc("session/load", Effect.fail(failure)).pipe(Effect.flip); + + expect(error).toMatchObject({ + _tag: "AcpRequestError", + code: -32602, + errorMessage: "Invalid params", + data: { field: "sessionId" }, + }); + }); + }); + + it("does not expose legacy diagnostic detail as the transport message", () => { + const cause = new Error("connection refused at a private endpoint"); + const error = new AcpError.AcpTransportError({ + detail: cause.message, + cause, + }); + + expect(error.message).toBe("ACP transport operation failed."); + expect(error.cause).toBe(cause); + }); + + it("preserves structured extension handler failures behind stable request errors", () => { + const cause = new AcpError.AcpTransportError({ + operation: "read-input-stream", + cause: new Error("private transport diagnostics"), + }); + const error = AcpError.AcpRequestError.fromExtensionHandlerError(cause, "x/test"); + + expect(error).toMatchObject({ + code: -32603, + method: "x/test", + operation: "handle-extension-request", + cause, + }); + expect(error.message).toBe("ACP extension request handler failed for method 'x/test'"); + expect(error.message).not.toContain(cause.message); + }); + + it.effect("uses the structured mapper for core handler failures", () => { + const cause = new AcpError.AcpTransportError({ + operation: "read-input-stream", + cause: new Error("private transport diagnostics"), + }); + + return Effect.gen(function* () { + const error = yield* runHandler(() => Effect.fail(cause), {}, "fs/read_text_file").pipe( + Effect.flip, + ); + + expect(error).toMatchObject({ + code: -32603, + message: "ACP request handler failed for method 'fs/read_text_file'", + }); + expect(error.message).not.toContain(cause.message); + }); + }); + + it.effect("keeps invalid extension payload values only in the exact schema cause", () => + Effect.gen(function* () { + const secret = "acp-schema-payload-secret"; + const cause = yield* decodeNestedNumberPayload({ profile: { token: secret } }).pipe( + Effect.flip, + ); + const error = AcpError.AcpRequestError.invalidExtensionPayload("x/private", cause); + const { cause: directCause, ...directDiagnostics } = error; + + expect(directCause).toBe(cause); + expect(error).toMatchObject({ + method: "x/private", + operation: "decode-extension-request-payload", + maximumPathDepth: 2, + }); + expect(error.issueCount).toBeGreaterThan(0); + expect(error.issueKinds).toContain("Pointer"); + expect(error.message).toBe("Invalid payload for ACP extension method 'x/private'."); + expect(error.message).not.toContain(secret); + expect(encodeUnknownJson(directDiagnostics)).not.toContain(secret); + expect(encodeUnknownJson(error.toProtocolError())).not.toContain(secret); + + const protocolError = AcpError.AcpProtocolParseError.fromSchemaError( + "decode-notification-payload", + "x/private", + cause, + ); + const { cause: protocolCause, ...protocolDiagnostics } = protocolError; + expect(protocolCause).toBe(cause); + expect(protocolError).toMatchObject({ + method: "x/private", + operation: "decode-notification-payload", + maximumPathDepth: 2, + }); + expect(protocolError.message).not.toContain(secret); + expect(encodeUnknownJson(protocolDiagnostics)).not.toContain(secret); + expect("detail" in protocolError).toBe(false); + }), + ); +}); diff --git a/packages/effect-acp/src/errors.ts b/packages/effect-acp/src/errors.ts index 91668f841f9..b3c0dee6294 100644 --- a/packages/effect-acp/src/errors.ts +++ b/packages/effect-acp/src/errors.ts @@ -1,7 +1,77 @@ import * as Schema from "effect/Schema"; +import type * as SchemaIssue from "effect/SchemaIssue"; import * as AcpSchema from "./_generated/schema.gen.ts"; +export const AcpRequestOperation = Schema.Literals([ + "decode-extension-request-payload", + "handle-request", + "handle-extension-request", +]); +export type AcpRequestOperation = typeof AcpRequestOperation.Type; + +export const AcpSchemaIssueKind = Schema.Literals([ + "Filter", + "Encoding", + "Pointer", + "Composite", + "AnyOf", + "InvalidType", + "InvalidValue", + "MissingKey", + "UnexpectedKey", + "Forbidden", + "OneOf", +]); +export type AcpSchemaIssueKind = typeof AcpSchemaIssueKind.Type; + +export interface AcpSchemaIssueDiagnostics { + readonly issueCount: number; + readonly issueKinds: ReadonlyArray; + readonly maximumPathDepth: number; +} + +const schemaIssueDiagnostics = (root: SchemaIssue.Issue): AcpSchemaIssueDiagnostics => { + let issueCount = 0; + let maximumPathDepth = 0; + const issueKinds = new Set(); + + const visit = (issue: SchemaIssue.Issue, pathDepth: number): void => { + issueCount += 1; + issueKinds.add(issue._tag); + maximumPathDepth = Math.max(maximumPathDepth, pathDepth); + switch (issue._tag) { + case "Filter": + case "Encoding": + visit(issue.issue, pathDepth); + break; + case "Pointer": + visit(issue.issue, pathDepth + issue.path.length); + break; + case "Composite": + case "AnyOf": + for (const child of issue.issues) visit(child, pathDepth); + break; + } + }; + + visit(root, 0); + return { + issueCount, + issueKinds: [...issueKinds], + maximumPathDepth, + }; +}; + +export interface AcpRequestDiagnostics { + readonly method?: string; + readonly operation?: AcpRequestOperation; + readonly cause?: unknown; + readonly issueCount?: number; + readonly issueKinds?: ReadonlyArray; + readonly maximumPathDepth?: number; +} + export class AcpSpawnError extends Schema.TaggedErrorClass()("AcpSpawnError", { command: Schema.optional(Schema.String), cause: Schema.Defect(), @@ -27,27 +97,68 @@ export class AcpProcessExitedError extends Schema.TaggedErrorClass()( "AcpProtocolParseError", { - detail: Schema.String, - cause: Schema.optional(Schema.Defect()), + operation: AcpProtocolParseOperation, + method: Schema.optionalKey(Schema.String), + issueCount: Schema.optionalKey(Schema.Number), + issueKinds: Schema.optionalKey(Schema.Array(AcpSchemaIssueKind)), + maximumPathDepth: Schema.optionalKey(Schema.Number), + cause: Schema.Defect(), }, ) { override get message() { - return `Failed to parse ACP protocol message: ${this.detail}`; + const method = this.method === undefined ? "" : ` for method '${this.method}'`; + return `ACP protocol operation '${this.operation}' failed${method}.`; + } + + static fromSchemaError( + operation: AcpProtocolParseOperation, + method: string, + cause: Schema.SchemaError, + ) { + return new AcpProtocolParseError({ + operation, + method, + ...schemaIssueDiagnostics(cause.issue), + cause, + }); } } export class AcpTransportError extends Schema.TaggedErrorClass()( "AcpTransportError", { - detail: Schema.String, + operation: Schema.optional( + Schema.Literals(["call-rpc", "read-input-stream", "read-process-exit-status"]), + ), + method: Schema.optional(Schema.String), + detail: Schema.optional(Schema.String), cause: Schema.Defect(), }, ) { override get message() { - return this.detail; + const method = this.method ? ` for method ${this.method}` : ""; + return this.operation + ? `ACP transport operation ${this.operation} failed${method}.` + : "ACP transport operation failed."; + } +} + +export class AcpInputStreamEndedError extends Schema.TaggedErrorClass()( + "AcpInputStreamEndedError", + {}, +) { + override get message() { + return "ACP input stream ended."; } } @@ -55,6 +166,12 @@ export class AcpRequestError extends Schema.TaggedErrorClass()( code: AcpSchema.ErrorCode, errorMessage: Schema.String, data: Schema.optional(Schema.Unknown), + method: Schema.optionalKey(Schema.String), + operation: Schema.optionalKey(AcpRequestOperation), + issueCount: Schema.optionalKey(Schema.Number), + issueKinds: Schema.optionalKey(Schema.Array(AcpSchemaIssueKind)), + maximumPathDepth: Schema.optionalKey(Schema.Number), + cause: Schema.optionalKey(Schema.Defect()), }) { override get message() { return this.errorMessage; @@ -68,6 +185,36 @@ export class AcpRequestError extends Schema.TaggedErrorClass()( }); } + static fromCoreHandlerError(error: AcpError, method: string) { + if (error._tag === "AcpRequestError") { + return error; + } + return AcpRequestError.internalError( + `ACP request handler failed for method '${method}'`, + undefined, + { + method, + operation: "handle-request", + cause: error, + }, + ); + } + + static fromExtensionHandlerError(error: AcpError, method: string) { + if (error._tag === "AcpRequestError") { + return error; + } + return AcpRequestError.internalError( + `ACP extension request handler failed for method '${method}'`, + undefined, + { + method, + operation: "handle-extension-request", + cause: error, + }, + ); + } + static parseError(message = "Parse error", data?: unknown) { return new AcpRequestError({ code: -32700, @@ -99,11 +246,29 @@ export class AcpRequestError extends Schema.TaggedErrorClass()( }); } - static internalError(message = "Internal error", data?: unknown) { + static invalidExtensionPayload(method: string, cause: Schema.SchemaError) { + const diagnostics = schemaIssueDiagnostics(cause.issue); + return new AcpRequestError({ + code: -32602, + errorMessage: `Invalid payload for ACP extension method '${method}'.`, + data: diagnostics, + method, + operation: "decode-extension-request-payload", + ...diagnostics, + cause, + }); + } + + static internalError( + message = "Internal error", + data?: unknown, + diagnostics: AcpRequestDiagnostics = {}, + ) { return new AcpRequestError({ code: -32603, errorMessage: message, ...(data !== undefined ? { data } : {}), + ...diagnostics, }); } @@ -138,6 +303,7 @@ export const AcpError = Schema.Union([ AcpProcessExitedError, AcpProtocolParseError, AcpTransportError, + AcpInputStreamEndedError, ]); 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 093d4acfcfa..c8e03dd7235 100644 --- a/packages/effect-acp/src/protocol.test.ts +++ b/packages/effect-acp/src/protocol.test.ts @@ -48,6 +48,8 @@ const decodeExtRequest = Schema.decodeEffect(Schema.fromJsonString(ExtRequest)); const decodeRequestPermissionResponse = Schema.decodeEffect( Schema.fromJsonString(RequestPermissionResponse), ); +const encodeUnknownJsonString = Schema.encodeUnknownSync(Schema.UnknownFromJsonString); +const encoder = new TextEncoder(); const mockPeerPath = Effect.map(Effect.service(Path.Path), (path) => path.join(import.meta.dirname, "../test/fixtures/acp-mock-peer.ts"), ); @@ -132,6 +134,49 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { }), ); + it.effect("keeps invalid core notification values only in the schema cause", () => + Effect.gen(function* () { + const secret = "acp-core-notification-secret-sentinel"; + const { stdio, input } = yield* makeInMemoryStdio(); + const termination = yield* Deferred.make(); + yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + yield* Queue.offer( + input, + encoder.encode( + `${encodeUnknownJsonString({ + jsonrpc: "2.0", + method: "session/update", + params: { + sessionId: { secret }, + update: { + sessionUpdate: "plan", + entries: [], + }, + }, + })}\n`, + ), + ); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, AcpError.AcpProtocolParseError); + const parseError = error as AcpError.AcpProtocolParseError; + const { cause, ...directDiagnostics } = parseError; + assert.equal(parseError.operation, "decode-notification-payload"); + assert.equal(parseError.method, "session/update"); + assert.isAbove(parseError.issueCount ?? 0, 0); + assert.include(parseError.issueKinds ?? [], "Pointer"); + assert.isAbove(parseError.maximumPathDepth ?? 0, 0); + assert.isTrue(Schema.isSchemaError(cause)); + assert.notInclude(parseError.message, secret); + assert.notInclude(encodeUnknownJsonString(directDiagnostics), secret); + }), + ); + it.effect("logs outgoing notifications when logOutgoing is enabled", () => Effect.gen(function* () { const { stdio } = yield* makeInMemoryStdio(); @@ -172,6 +217,38 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { }), ); + it.effect("logs decode failures without copying the cause or wire payload", () => + Effect.gen(function* () { + const secret = "acp-wire-secret-sentinel"; + const { stdio, input } = yield* makeInMemoryStdio(); + const events: Array = []; + const termination = yield* Deferred.make(); + yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + logIncoming: true, + logger: (event) => + Effect.sync(() => { + events.push(event); + }), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + yield* Queue.offer(input, encoder.encode(`{"secret":"${secret}"\n`)); + yield* Deferred.await(termination); + + const event = events.find(({ stage }) => stage === "decode_failed"); + assert.deepEqual(event, { + direction: "incoming", + stage: "decode_failed", + payload: { + operation: "decode-wire-message", + }, + }); + assert.notInclude(encodeUnknownJsonString(event), secret); + }), + ); + it.effect("fails notification encoding through the declared ACP error channel", () => Effect.gen(function* () { const { stdio } = yield* makeInMemoryStdio(); @@ -182,13 +259,16 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { const bigintError = yield* transport.notify("x/test", 1n).pipe(Effect.flip); assert.instanceOf(bigintError, AcpError.AcpProtocolParseError); - assert.equal(bigintError.detail, "Failed to encode ACP message"); + assert.equal(bigintError.operation, "encode-message"); + assert.instanceOf(bigintError.cause, TypeError); + assert.equal(bigintError.message, "ACP protocol operation 'encode-message' failed."); const circular: Record = {}; circular.self = circular; const circularError = yield* transport.notify("x/test", circular).pipe(Effect.flip); assert.instanceOf(circularError, AcpError.AcpProtocolParseError); - assert.equal(circularError.detail, "Failed to encode ACP message"); + assert.equal(circularError.operation, "encode-message"); + assert.instanceOf(circularError.cause, TypeError); }), ); @@ -381,14 +461,35 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { assert.equal((message as { readonly _tag?: string })._tag, "ClientProtocolError"); const defect = (message as { readonly error: { readonly reason: unknown } }).error.reason as { readonly _tag: string; + readonly message: string; readonly cause: unknown; }; assert.equal(defect._tag, "RpcClientDefect"); + assert.equal(defect.message, "ACP protocol terminated."); assert.instanceOf(defect.cause, AcpError.AcpProcessExitedError); assert.equal((defect.cause as AcpError.AcpProcessExitedError).code, 7); }), ); + it.effect("classifies an input stream ending without inventing a cause", () => + Effect.gen(function* () { + const { stdio, input } = yield* makeInMemoryStdio(); + const termination = yield* Deferred.make(); + yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + yield* Queue.end(input); + + const error = yield* Deferred.await(termination); + assert.instanceOf(error, AcpError.AcpInputStreamEndedError); + assert.equal(error.message, "ACP input stream ended."); + assert.equal("cause" in error, false); + }), + ); + it.effect("does not emit a second process-exit error after a decode failure", () => Effect.gen(function* () { const handle = yield* makeHandle({ @@ -413,9 +514,40 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { assert.equal((message as { readonly _tag?: string })._tag, "ClientProtocolError"); const defect = (message as { readonly error: { readonly reason: unknown } }).error.reason as { readonly _tag: string; + readonly message: string; readonly cause: unknown; }; assert.equal(defect._tag, "RpcClientDefect"); + assert.equal(defect.message, "ACP protocol terminated."); + assert.instanceOf(defect.cause, AcpError.AcpProtocolParseError); + }), + ); + + it.effect("keeps client send failure messages independent from the cause", () => + Effect.gen(function* () { + const { stdio } = yield* makeInMemoryStdio(); + const transport = yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + }); + + const failure = yield* transport.clientProtocol + .send(0, { + _tag: "Request", + id: "request-1", + tag: "x/test", + payload: 1n, + headers: [], + }) + .pipe(Effect.flip); + const defect = failure.reason as { + readonly _tag: string; + readonly message: string; + readonly cause: unknown; + }; + + assert.equal(defect._tag, "RpcClientDefect"); + assert.equal(defect.message, "Failed to send ACP protocol message."); assert.instanceOf(defect.cause, AcpError.AcpProtocolParseError); }), ); diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index 56a7ce81ab8..6c3bd399028 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -17,7 +17,6 @@ import * as AcpSchema from "./_generated/schema.gen.ts"; import { CLIENT_METHODS } from "./_generated/meta.gen.ts"; import * as AcpError from "./errors.ts"; const isAcpError = Schema.is(AcpError.AcpError); -const isAcpRequestError = Schema.is(AcpError.AcpRequestError); export interface AcpProtocolLogEvent { readonly direction: "incoming" | "outgoing"; @@ -114,7 +113,7 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi try: () => parser.encode(message), catch: (cause) => new AcpError.AcpProtocolParseError({ - detail: "Failed to encode ACP message", + operation: "encode-message", cause, }), }); @@ -184,7 +183,7 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi _tag: "ClientProtocolError", error: new RpcClientError.RpcClientError({ reason: new RpcClientError.RpcClientDefect({ - message: error.message, + message: "ACP protocol terminated.", cause: error, }), }), @@ -243,7 +242,11 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi } return options.onExtRequest(message.tag, message.payload).pipe( Effect.matchEffect({ - onFailure: (error) => respondWithError(message.id, normalizeToRequestError(error)), + onFailure: (error) => + respondWithError( + message.id, + AcpError.AcpRequestError.fromExtensionHandlerError(error, message.tag), + ), onSuccess: (value) => respondWithSuccess(message.id, value), }), ); @@ -261,12 +264,12 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi params, }) satisfies AcpIncomingNotification, ), - Effect.mapError( - (cause) => - new AcpError.AcpProtocolParseError({ - detail: `Invalid ${CLIENT_METHODS.session_update} notification payload`, - cause, - }), + Effect.mapError((cause) => + AcpError.AcpProtocolParseError.fromSchemaError( + "decode-notification-payload", + CLIENT_METHODS.session_update, + cause, + ), ), Effect.flatMap(dispatchNotification), ); @@ -281,12 +284,12 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi params, }) satisfies AcpIncomingNotification, ), - Effect.mapError( - (cause) => - new AcpError.AcpProtocolParseError({ - detail: `Invalid ${CLIENT_METHODS.session_elicitation_complete} notification payload`, - cause, - }), + Effect.mapError((cause) => + AcpError.AcpProtocolParseError.fromSchemaError( + "decode-notification-payload", + CLIENT_METHODS.session_elicitation_complete, + cause, + ), ), Effect.flatMap(dispatchNotification), ); @@ -379,7 +382,7 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi >, catch: (cause) => new AcpError.AcpProtocolParseError({ - detail: "Failed to decode ACP wire message", + operation: "decode-wire-message", cause, }), }), @@ -396,8 +399,13 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi direction: "incoming", stage: "decode_failed", payload: { - detail: error.detail, - cause: error.cause, + operation: error.operation, + ...(error.method === undefined ? {} : { method: error.method }), + ...(error.issueCount === undefined ? {} : { issueCount: error.issueCount }), + ...(error.issueKinds === undefined ? {} : { issueKinds: error.issueKinds }), + ...(error.maximumPathDepth === undefined + ? {} + : { maximumPathDepth: error.maximumPathDepth }), }, }), ), @@ -413,7 +421,7 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi const normalized: AcpError.AcpError = isAcpError(error) ? error : new AcpError.AcpTransportError({ - detail: error instanceof Error ? error.message : String(error), + operation: "read-input-stream", cause: error, }); return handleTermination(() => Effect.succeed(normalized)); @@ -421,13 +429,7 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi onSuccess: () => handleTermination( () => - options.terminationError ?? - Effect.succeed( - new AcpError.AcpTransportError({ - detail: "ACP input stream ended", - cause: new Error("ACP input stream ended"), - }), - ), + options.terminationError ?? Effect.succeed(new AcpError.AcpInputStreamEndedError({})), ), }), Effect.forkScoped, @@ -441,7 +443,18 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi Stream.runForEach((message) => f(message)), Effect.forever, ), - send: (_clientId, request) => offerOutgoing(request).pipe(Effect.mapError(toRpcClientError)), + send: (_clientId, request) => + offerOutgoing(request).pipe( + Effect.mapError( + (error) => + new RpcClientError.RpcClientError({ + reason: new RpcClientError.RpcClientDefect({ + message: "Failed to send ACP protocol message.", + cause: error, + }), + }), + ), + ), supportsAck: true, supportsTransferables: false, }); @@ -521,16 +534,3 @@ function isProtocolError( typeof value.message === "string" ); } - -function normalizeToRequestError(error: AcpError.AcpError): AcpError.AcpRequestError { - return isAcpRequestError(error) ? error : AcpError.AcpRequestError.internalError(error.message); -} - -function toRpcClientError(error: AcpError.AcpError): RpcClientError.RpcClientError { - return new RpcClientError.RpcClientError({ - reason: new RpcClientError.RpcClientDefect({ - message: error.message, - cause: error, - }), - }); -} diff --git a/packages/effect-acp/src/terminal.ts b/packages/effect-acp/src/terminal.ts index 088ff863738..b892f040436 100644 --- a/packages/effect-acp/src/terminal.ts +++ b/packages/effect-acp/src/terminal.ts @@ -23,23 +23,3 @@ export interface AcpTerminal { */ readonly release: Effect.Effect; } - -export interface MakeTerminalOptions { - readonly sessionId: string; - readonly terminalId: string; - readonly output: Effect.Effect; - readonly waitForExit: Effect.Effect; - readonly kill: Effect.Effect; - readonly release: Effect.Effect; -} - -export function makeTerminal(options: MakeTerminalOptions): AcpTerminal { - return { - sessionId: options.sessionId, - terminalId: options.terminalId, - output: options.output, - waitForExit: options.waitForExit, - kill: options.kill, - release: options.release, - }; -}