diff --git a/apps/server/src/observability/DefectReporter.test.ts b/apps/server/src/observability/DefectReporter.test.ts new file mode 100644 index 000000000000..28c81180ba67 --- /dev/null +++ b/apps/server/src/observability/DefectReporter.test.ts @@ -0,0 +1,102 @@ +import { assert, describe, it } from "@effect/vitest"; +import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as ErrorReporter from "effect/ErrorReporter"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Logger from "effect/Logger"; +import * as Queue from "effect/Queue"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import { Rpc, RpcClient, RpcGroup, RpcServer } from "effect/rpc"; + +import { WS_RPC_SERVER_OPTIONS } from "../ws.ts"; +import * as DefectReporter from "./DefectReporter.ts"; + +class TestRpcs extends RpcGroup.make( + Rpc.make("subscribe", { success: Schema.Number, stream: true }), + Rpc.make("boom", { success: Schema.Void }), +) {} +type TestRpc = RpcGroup.Rpcs; + +/** Runs `body` with every error log captured. */ +const withErrorLogs = ( + body: (logs: Queue.Queue>) => Effect.Effect, +) => + Effect.gen(function* () { + const logs = yield* Queue.unbounded>(); + const logger = Logger.make(({ cause, logLevel }) => { + if (logLevel === "Error") Queue.offerUnsafe(logs, cause); + }); + return yield* body(logs).pipe( + Effect.provide(Logger.layer([logger], { mergeWithExisting: false })), + ); + }); + +const errorMessage = (cause: Cause.Cause) => (Cause.squash(cause) as Error).message; + +describe("DefectReporter", () => { + it.effect("a dying RPC handler fails alone, and its defect is logged", () => + withErrorLogs((logs) => + Effect.gen(function* () { + const subscription = yield* Queue.unbounded(); + // The same client/server pairing as RpcTest.makeClient, with ws.ts's options. + let client!: Effect.Success< + ReturnType> + >; + const server = yield* RpcServer.makeNoSerialization(TestRpcs, { + ...WS_RPC_SERVER_OPTIONS, + onFromServer: (response) => client.write(response), + }).pipe( + // Provided to the handlers, as ws.ts does. + Effect.provide( + TestRpcs.toLayer({ + subscribe: () => Stream.fromQueue(subscription), + boom: () => Effect.die(new Error("handler bug")), + }).pipe(Layer.provide(DefectReporter.layer)), + ), + ); + client = yield* RpcClient.makeNoSerialization(TestRpcs, { + supportsAck: true, + onFromClient: ({ message }) => server.write(0, message), + }); + + const received = yield* Queue.unbounded(); + const sibling = yield* client.client.subscribe().pipe( + Stream.runForEach((value) => Queue.offer(received, value)), + Effect.forkScoped, + ); + yield* Queue.offer(subscription, 1); + assert.equal(yield* Queue.take(received), 1); + + const boom = yield* Effect.exit(client.client.boom()); + assert.isTrue(Exit.hasDies(boom)); + // The server reports before it answers, so the log is already written. + const logged = yield* Queue.clear(logs); + assert.deepEqual(logged.map(errorMessage), ["handler bug"]); + + // The sibling subscription on the same client keeps delivering. + yield* Queue.offer(subscription, 2); + const next = yield* Queue.take(received).pipe( + Effect.raceFirst(Fiber.await(sibling).pipe(Effect.as("subscription ended"))), + ); + assert.equal(next, 2); + assert.equal(yield* Queue.size(logs), 0); + }), + ).pipe(Effect.scoped), + ); + + it.effect("logs defects, not typed failures or interrupts", () => + withErrorLogs((logs) => + Effect.gen(function* () { + yield* ErrorReporter.report(Cause.fail(new Error("expected"))); + yield* ErrorReporter.report(Cause.interrupt()); + yield* ErrorReporter.report(Cause.die(new Error("bug"))); + + // Reporters log synchronously, so the logs are already written. + assert.deepEqual((yield* Queue.clear(logs)).map(errorMessage), ["bug"]); + }).pipe(Effect.provide(DefectReporter.layer)), + ), + ); +}); diff --git a/apps/server/src/observability/DefectReporter.ts b/apps/server/src/observability/DefectReporter.ts new file mode 100644 index 000000000000..220b2696d831 --- /dev/null +++ b/apps/server/src/observability/DefectReporter.ts @@ -0,0 +1,24 @@ +import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as ErrorReporter from "effect/ErrorReporter"; + +/** + * Logs the defects of WebSocket RPC handlers. RpcServer reports every failed + * handler exit; typed failures are expected responses, so only dies are logged. + */ +const reporter: ErrorReporter.ErrorReporter = { + [ErrorReporter.TypeId]: ErrorReporter.TypeId, + report: ({ cause, fiber }) => { + for (const reason of cause.reasons) { + if (reason._tag !== "Die" || ErrorReporter.isIgnored(reason.defect)) continue; + // Reporters are called synchronously from the failing fiber. Logging with + // its context keeps its loggers and annotations, and a fork cannot throw + // back into the server that reported. + Effect.runForkWith(fiber.context)( + Effect.logError("Unhandled defect", Cause.fromReasons([reason])), + ); + } + }, +}; + +export const layer = ErrorReporter.layer([reporter]); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 99a05b8e205b..dcf9fad56835 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -210,6 +210,7 @@ import * as WorktreeSetupTracker from "./project/WorktreeSetupTracker.ts"; import * as ServerEnvironment from "./environment/ServerEnvironment.ts"; import * as DirectEndpoints from "./environment/DirectEndpoints.ts"; import * as RemoteOpenTargets from "./environment/RemoteOpenTargets.ts"; +import * as DefectReporter from "./observability/DefectReporter.ts"; import * as BackgroundPolicy from "./background/BackgroundPolicy.ts"; import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts"; import { requiredScopeForDeviceList, rpcAuthorizationError } from "./auth/RpcAuthorization.ts"; @@ -3787,6 +3788,14 @@ const layerWsRpc = ( }), ); +// A defect in a handler's effect fails only its own request. RpcServer's default +// sends a socket-level Defect frame instead, and the client ends every pending +// request on the socket with it. DefectReporter logs these defects. +export const WS_RPC_SERVER_OPTIONS = { + disableTracing: true, + disableFatalDefects: true, +} as const; + export const layer = Layer.unwrap( Effect.gen(function* () { const previewAutomationBroker = yield* PreviewAutomationBroker.PreviewAutomationBroker; @@ -3830,7 +3839,7 @@ export const layer = Layer.unwrap( yield* analytics.record("client.connected", clientAnalyticsProps); const rpcWebSocketHttpEffect = yield* Effect.gen(function* () { const { protocol, httpEffect } = yield* RpcServer.makeProtocolWithHttpEffectWebsocket; - yield* RpcServer.make(ServerWsRpcGroup, { disableTracing: true }).pipe( + yield* RpcServer.make(ServerWsRpcGroup, WS_RPC_SERVER_OPTIONS).pipe( Effect.provideService(RpcServer.Protocol, withTerminalOutputWindow(protocol)), Effect.provide(RpcAuthorization.layer(session.scopes)), Effect.forkScoped, @@ -3847,6 +3856,9 @@ export const layer = Layer.unwrap( serverBrowser, ).pipe( Layer.provideMerge(RpcSerialization.layerJson), + // Request fibers run in the handlers' context, so this reporter sees + // their defects, not the rest of the server's. + Layer.provide(DefectReporter.layer), Layer.provide(Layer.succeed(SqlClient.SqlClient, sql)), Layer.provide(AgentSessionScanner.layer), Layer.provide(ProviderMaintenanceRunner.layer), diff --git a/packages/effect-acp/src/_internal/shared.ts b/packages/effect-acp/src/_internal/shared.ts index e4a98d51d8bf..527fb99304c6 100644 --- a/packages/effect-acp/src/_internal/shared.ts +++ b/packages/effect-acp/src/_internal/shared.ts @@ -1,3 +1,4 @@ +import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; import { RpcClientError } from "effect/rpc"; @@ -29,6 +30,19 @@ export const callRpc = ( }), ); +/** + * Runs a notification handler so it cannot stop the caller: a typed failure is + * dropped, as notifications have no reply, and a defect is logged. + */ +export const isolateNotificationHandler = (effect: Effect.Effect) => + effect.pipe( + Effect.catchCause((cause) => + Cause.hasDies(cause) + ? Effect.logError("ACP notification handler failed", cause) + : Effect.void, + ), + ); + export const runHandler = Effect.fnUntraced(function* >( handler: ((payload: A, ...args: Args) => Effect.Effect) | undefined, payload: A, @@ -39,6 +53,9 @@ export const runHandler = Effect.fnUntraced(function* + Effect.logError(`ACP request handler failed for '${method}'`, defect), + ), Effect.mapError((error) => AcpError.AcpRequestError.fromCoreHandlerError(error, method).toProtocolError(), ), diff --git a/packages/effect-acp/src/agent.test.ts b/packages/effect-acp/src/agent.test.ts index 9b69ecb5918d..416c4808a3ed 100644 --- a/packages/effect-acp/src/agent.test.ts +++ b/packages/effect-acp/src/agent.test.ts @@ -35,6 +35,14 @@ const SessionCancelNotification = jsonRpcNotification( const ExtPingNotification = jsonRpcNotification("x/ping", Schema.Struct({ count: Schema.Number })); const ExtRequest = jsonRpcRequest("x/test", Schema.Struct({ hello: Schema.String })); const ExtResponse = jsonRpcResponse(Schema.Struct({ ok: Schema.Boolean })); +/** A response whose cause is a handler's defect, as RpcServer encodes it. */ +const DieResponse = Schema.Struct({ + id: Schema.Number, + error: Schema.Struct({ + _tag: Schema.Literal("Cause"), + data: Schema.Tuple([Schema.Struct({ _tag: Schema.Literal("Die") })]), + }), +}); const decodeRequestPermissionRequest = Schema.decodeEffect( Schema.fromJsonString(RequestPermissionRequest), ); @@ -296,3 +304,33 @@ it.effect("effect-acp agent uses distinct ids for RPC calls and extension reques }).pipe(Effect.provide(context), Effect.ensuring(Scope.close(scope, Exit.void))); }), ); + +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(); + const scope = yield* Scope.make(); + const context = yield* Layer.buildWithScope(AcpAgent.layer(stdio), scope); + const agent = yield* Effect.service(AcpAgent.AcpAgent).pipe(Effect.provide(context)); + yield* agent.handleInitialize(() => Effect.die(new Error("handler bug"))); + + yield* Queue.offer( + input, + yield* encodeJsonl(InitializeRequest, { + jsonrpc: "2.0", + id: 7, + method: "initialize", + params: { + protocolVersion: 2, + capabilities: {}, + info: { name: "effect-acp-test", version: "0.0.0" }, + }, + headers: [], + }), + ); + const response = yield* Queue.take(output).pipe( + Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(DieResponse))), + ); + assert.equal(response.id, 7); + yield* Scope.close(scope, Exit.void); + }), +); diff --git a/packages/effect-acp/src/agent.ts b/packages/effect-acp/src/agent.ts index 153fea8428ee..2ddf5cc8cc0d 100644 --- a/packages/effect-acp/src/agent.ts +++ b/packages/effect-acp/src/agent.ts @@ -1,5 +1,6 @@ import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; +import * as ErrorReporter from "effect/ErrorReporter"; import * as Layer from "effect/Layer"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; @@ -18,6 +19,7 @@ import { callRpc, decodeExtNotificationRegistration, decodeExtRequestRegistration, + isolateNotificationHandler, runHandler, } from "./_internal/shared.ts"; @@ -280,7 +282,11 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* ( ), ), Effect.flatMap((decoded) => - Effect.forEach(cancelHandlers, (handler) => handler(decoded), { discard: true }), + Effect.forEach( + cancelHandlers, + (handler) => isolateNotificationHandler(handler(decoded)), + { discard: true }, + ), ), ); } @@ -421,7 +427,10 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* ( }), ); - yield* RpcServer.make(AcpRpcs.AgentRpcs).pipe( + yield* RpcServer.make(AcpRpcs.AgentRpcs, { disableFatalDefects: true }).pipe( + // runHandler logs handler defects with their method. A reporter inherited + // from the caller (a WebSocket request, say) would log them again. + Effect.provideService(ErrorReporter.CurrentErrorReporters, new Set()), Effect.provideService(RpcServer.Protocol, transport.serverProtocol), Effect.provide(layerAgentHandler), Effect.forkScoped, diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index d25df6ccf5e7..6209dcffa41f 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -5,11 +5,13 @@ import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; +import * as Logger from "effect/Logger"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Schema from "effect/Schema"; import * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; +import * as TestClock from "effect/testing/TestClock"; import { ChildProcess, ChildProcessSpawner } from "effect/process"; import * as NodeServices from "@effect/platform-node/NodeServices"; @@ -50,6 +52,16 @@ const PermissionRequest = jsonRpcRequest( const PermissionResponse = jsonRpcResponse(AcpSchema.RequestPermissionResponse); const ElicitationRequest = jsonRpcRequest("elicitation/create", AcpSchema.CreateElicitationRequest); const ElicitationResponse = jsonRpcResponse(AcpSchema.CreateElicitationResponse); +/** A JSON-RPC error response; only its id and the presence of an error matter here. */ +const ErrorResponse = Schema.Struct({ id: Schema.String, error: Schema.Unknown }); +/** A response whose cause is a handler's defect, as RpcServer encodes it. */ +const DieResponse = Schema.Struct({ + id: Schema.String, + error: Schema.Struct({ + _tag: Schema.Literal("Cause"), + data: Schema.Tuple([Schema.Struct({ _tag: Schema.Literal("Die") })]), + }), +}); const decodePromptRequestLine = Schema.decodeEffect(Schema.fromJsonString(PromptRequest)); const XAiPromptCompleteNotification = jsonRpcNotification( "_x.ai/session/prompt_complete", @@ -1130,6 +1142,117 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { }), ); + it.effect("answers each request whose handler dies, and keeps reading", () => + Effect.gen(function* () { + const errorLogs = yield* Queue.unbounded(); + const logger = Logger.make(({ logLevel, message }) => { + if (logLevel === "Error") Queue.offerUnsafe(errorLogs, String([message].flat()[0])); + }); + const { stdio, input, output } = yield* makeInMemoryStdio(); + const scope = yield* Scope.make(); + const acp = yield* AcpClient.make(stdio).pipe( + Effect.provideService(Scope.Scope, scope), + Effect.provide(Logger.layer([logger])), + ); + const bug = () => Effect.die(new Error("handler bug")); + yield* acp.handleRequestPermission(bug); + yield* acp.handleExtRequest("x/dies", Schema.Unknown, bug); + yield* acp.handleExtNotification("x/notification-dies", Schema.Unknown, bug); + yield* acp.handleSessionUpdate(bug); + const updates = yield* Queue.unbounded(); + yield* acp.handleSessionUpdate((notification) => + Queue.offer(updates, notification.sessionId).pipe(Effect.asVoid), + ); + yield* acp.handleExtRequest("x/test", Schema.Struct({ hello: Schema.String }), () => + Effect.succeed({ ok: true }), + ); + const send = (line: Effect.Effect) => + Effect.flatMap(line, (bytes) => Queue.offer(input, bytes)); + const decodeDie = Schema.decodeEffect(Schema.fromJsonString(DieResponse)); + // A stopped reader leaves these waiting; fail instead of hanging. + const next = (queue: Queue.Dequeue) => + TestClock.withLive(Queue.take(queue).pipe(Effect.timeout("2 seconds"))); + + yield* send( + encodeJsonl(PermissionRequest, { + jsonrpc: "2.0", + id: "permission-a", + method: "session/request_permission", + params: { + sessionId: "session-1", + title: "Tool", + subject: { + type: "tool_call" as const, + toolCall: { toolCallId: "tool-1", title: "Tool" }, + }, + options: [{ optionId: "allow", name: "Allow", kind: "allow_once" as const }], + }, + headers: [], + }), + ); + assert.equal((yield* next(output).pipe(Effect.flatMap(decodeDie))).id, "permission-a"); + + yield* send( + encodeJsonl(jsonRpcRequest("x/dies", Schema.Unknown), { + jsonrpc: "2.0", + id: "ext-a", + method: "x/dies", + params: {}, + headers: [], + }), + ); + const extDied = yield* next(output).pipe( + Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ErrorResponse))), + ); + assert.equal(extDied.id, "ext-a"); + + yield* send( + encodeJsonl(jsonRpcNotification("x/notification-dies", Schema.Unknown), { + jsonrpc: "2.0", + method: "x/notification-dies", + params: {}, + }), + ); + yield* send( + encodeJsonl(SessionUpdateNotification, { + jsonrpc: "2.0", + method: "session/update", + params: { + sessionId: "session-1", + update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text: "hi" } }, + }, + }), + ); + + // The next session handler still ran. + assert.equal(yield* next(updates), "session-1"); + + // The reader survived all three: a later request is still answered. + yield* send( + encodeJsonl(ExtRequest, { + jsonrpc: "2.0", + id: "ext-b", + method: "x/test", + params: { hello: "world" }, + headers: [], + }), + ); + const answered = yield* next(output).pipe( + Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ExtResponse))), + ); + assert.deepEqual([answered.id, answered.result], ["ext-b", { ok: true }]); + + // Every defect was logged at Error, once. + assert.deepEqual((yield* Queue.clear(errorLogs)).toSorted(), [ + "ACP extension request handler failed for 'x/dies'", + "ACP notification handler failed", + "ACP notification handler failed", + "ACP request handler failed for 'session/request_permission'", + ]); + yield* Scope.close(scope, Exit.void); + }), + ); + it.effect("answers elicitation/create with the flat action shape", () => Effect.gen(function* () { const { stdio, input, output } = yield* makeInMemoryStdio(); diff --git a/packages/effect-acp/src/client.ts b/packages/effect-acp/src/client.ts index 2df7f9409f4e..732e2d26650f 100644 --- a/packages/effect-acp/src/client.ts +++ b/packages/effect-acp/src/client.ts @@ -1,6 +1,7 @@ import * as Context from "effect/Context"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; +import * as ErrorReporter from "effect/ErrorReporter"; import * as Layer from "effect/Layer"; import * as Schema from "effect/Schema"; import * as Predicate from "effect/Predicate"; @@ -22,6 +23,7 @@ import { callRpc, decodeExtNotificationRegistration, decodeExtRequestRegistration, + isolateNotificationHandler, runHandler, } from "./_internal/shared.ts"; import { makeChildStdio, makeTerminationError } from "./_internal/stdio.ts"; @@ -838,9 +840,12 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* ( registration: BufferedNotificationHandler, notification: A, ) => - Effect.forEach(registration.handlers, (handler) => handler(notification).pipe(Effect.ignore), { - discard: true, - }); + // One handler failing or dying does not stop the others, or the reader. + Effect.forEach( + registration.handlers, + (handler) => isolateNotificationHandler(handler(notification)), + { discard: true }, + ); const flushBufferedNotifications = (registration: BufferedNotificationHandler) => Effect.suspend(() => { @@ -1099,7 +1104,10 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* ( }), ); - yield* RpcServer.make(AcpRpcs.CompatClientRpcs).pipe( + yield* RpcServer.make(AcpRpcs.CompatClientRpcs, { disableFatalDefects: true }).pipe( + // runHandler logs handler defects with their method. A reporter inherited + // from the caller (a WebSocket request, say) would log them again. + Effect.provideService(ErrorReporter.CurrentErrorReporters, new Set()), Effect.provideService(RpcServer.Protocol, transport.serverProtocol), Effect.provide(layerClientHandler), Effect.forkScoped, diff --git a/packages/effect-acp/src/protocol.test.ts b/packages/effect-acp/src/protocol.test.ts index 91677618e070..9c3f41e5f4c9 100644 --- a/packages/effect-acp/src/protocol.test.ts +++ b/packages/effect-acp/src/protocol.test.ts @@ -11,6 +11,7 @@ import * as Sink from "effect/Sink"; import * as Stdio from "effect/Stdio"; import * as Stream from "effect/Stream"; import * as Ref from "effect/Ref"; +import * as TestClock from "effect/testing/TestClock"; import { ChildProcess, ChildProcessSpawner } from "effect/process"; import { it, assert } from "@effect/vitest"; @@ -257,6 +258,42 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { }), ); + it.effect("terminates when a callback on the reader dies", () => + Effect.gen(function* () { + const { stdio, input } = yield* makeInMemoryStdio(); + const termination = yield* Deferred.make(); + yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + transformSessionUpdate: () => { + throw new Error("normalizer bug"); + }, + onTermination: (error) => Deferred.succeed(termination, error).pipe(Effect.asVoid), + }); + + yield* Queue.offer( + input, + encoder.encode( + `${encodeUnknownJsonString({ + jsonrpc: "2.0", + method: "session/update", + params: { + sessionId: "session-1", + update: { sessionUpdate: "plan", entries: [] }, + }, + })}\n`, + ), + ); + + // Pending requests are failed through termination instead of hanging. + const error = yield* TestClock.withLive( + Deferred.await(termination).pipe(Effect.timeout("2 seconds")), + ); + assert.instanceOf(error, AcpError.AcpTransportError); + assert.equal((error as AcpError.AcpTransportError).operation, "read-input-stream"); + }), + ); + it.effect("logs outgoing notifications when logOutgoing is enabled", () => Effect.gen(function* () { const { stdio } = yield* makeInMemoryStdio(); @@ -500,6 +537,38 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { }), ); + it.effect("answers an extension request whose handler also dies as an internal error", () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + yield* AcpProtocol.makeAcpPatchedProtocol({ + stdio, + serverRequestMethods: new Set(), + // The typed failure must not hide the defect from its cleanup. + onExtRequest: () => + Effect.fail(AcpError.AcpRequestError.invalidParams("bad params")).pipe( + Effect.ensuring(Effect.die(new Error("cleanup bug"))), + ), + }); + + yield* Queue.offer( + input, + yield* encodeJsonl(ExtRequest, { + jsonrpc: "2.0", + id: 9, + method: "x/test", + params: { hello: "world" }, + headers: [], + }), + ); + + const response = yield* Schema.decodeUnknownEffect( + Schema.fromJsonString(JsonRpcErrorResponse), + )(yield* Queue.take(output)); + assert.equal(response.id, 9); + assert.equal(response.error.code, -32603); + }), + ); + it.effect("preserves numeric ids for inbound extension requests", () => Effect.gen(function* () { const { stdio, input, output } = yield* makeInMemoryStdio(); diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index fb926cfca053..12e85a9e7542 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -3,6 +3,7 @@ import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import * as Deferred from "effect/Deferred"; 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"; @@ -20,6 +21,7 @@ import * as AcpSchema from "./schema.ts"; import * as AcpSchemaV1 from "./_generated/schema-v1.gen.ts"; import { CLIENT_METHODS } from "./_generated/meta.gen.ts"; import * as AcpError from "./errors.ts"; +import { isolateNotificationHandler } from "./_internal/shared.ts"; const isAcpError = Schema.is(AcpError.AcpError); export interface AcpProtocolLogEvent { @@ -385,8 +387,10 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi const dispatchNotification = (notification: AcpIncomingNotification) => Queue.offer(notificationQueue, notification).pipe( Effect.andThen( + // A failing or dying handler must not stop the reader, or every later + // message on the connection goes unanswered. options.onNotification - ? options.onNotification(notification).pipe(Effect.ignore) + ? isolateNotificationHandler(options.onNotification(notification)) : Effect.void, ), Effect.asVoid, @@ -458,12 +462,38 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi method: message.tag, }) .pipe( - Effect.matchEffect({ - onFailure: (error) => - respondWithError( - message.id, - AcpError.AcpRequestError.fromExtensionHandlerError(error, message.tag), - ), + Effect.matchCauseEffect({ + // A dying handler answers its own request, like a core handler, and + // leaves the reader running. A defect wins over a typed failure in + // the same cause, so it is never hidden behind an expected error. + onFailure: (cause) => { + const failure = Cause.hasDies(cause) ? Option.none() : Cause.findErrorOption(cause); + if (Option.isSome(failure)) { + return respondWithError( + message.id, + AcpError.AcpRequestError.fromExtensionHandlerError(failure.value, message.tag), + ); + } + return Effect.logError( + `ACP extension request handler failed for '${message.tag}'`, + cause, + ).pipe( + Effect.andThen( + respondWithError( + message.id, + AcpError.AcpRequestError.internalError( + `ACP extension request handler failed for method '${message.tag}'`, + undefined, + { + method: message.tag, + operation: "handle-extension-request", + cause: Cause.squash(cause), + }, + ), + ), + ), + ); + }, onSuccess: (value) => respondWithSuccess(message.id, value), }), ); @@ -685,8 +715,14 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi ), ), ), - Effect.matchEffect({ - onFailure: (error) => { + // Anything that ends the reader, including a defect in a callback it runs, + // terminates the connection so pending requests fail instead of hanging. + Effect.matchCauseEffect({ + onFailure: (cause) => { + // The reader's own interruption never gets here: an interrupted fiber + // skips failure handlers. An interrupt raised inside it ends it too. + const failure = Cause.findErrorOption(cause); + const error = Option.isSome(failure) ? failure.value : Cause.squash(cause); const normalized: AcpError.AcpError = isAcpError(error) ? error : new AcpError.AcpTransportError({