From 3fd5091bf03e8c0d101d84359187a2f309c7d657 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sat, 3 Oct 2026 23:20:09 -0700 Subject: [PATCH 1/8] fix(server): one failing RPC handler no longer ends the client's other requests A defect in any WebSocket RPC handler was sent as a socket-level Defect frame, and the client ends every pending request on the socket when it gets one: shell, thread and terminal subscriptions included. No ErrorReporter was installed, so the defect was not logged on the server. Set disableFatalDefects so a defect fails only its own request, and install a defect-only ErrorReporter in the server's observability layer so defects from RPC, HTTP and HttpApi handlers are logged. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/observability/DefectReporter.test.ts | 123 ++++++++++++++++++ .../src/observability/DefectReporter.ts | 41 ++++++ .../src/observability/Layers/Observability.ts | 9 +- apps/server/src/ws.ts | 10 +- 4 files changed, 181 insertions(+), 2 deletions(-) create mode 100644 apps/server/src/observability/DefectReporter.test.ts create mode 100644 apps/server/src/observability/DefectReporter.ts diff --git a/apps/server/src/observability/DefectReporter.test.ts b/apps/server/src/observability/DefectReporter.test.ts new file mode 100644 index 000000000000..0a7ea4a92369 --- /dev/null +++ b/apps/server/src/observability/DefectReporter.test.ts @@ -0,0 +1,123 @@ +import { assert, describe, it } from "@effect/vitest"; +import * as Cause from "effect/Cause"; +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 Logger from "effect/Logger"; +import * as Queue from "effect/Queue"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import { Rpc, RpcGroup, RpcMessage, RpcSerialization, RpcServer } from "effect/unstable/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 }), +) {} + +/** Runs `body` with the defect reporter installed and 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( + Layer.merge(DefectReporter.layer, 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(); + const responses = yield* Queue.unbounded(); + const receive = yield* Deferred.make[0]>(); + const protocol = yield* RpcServer.Protocol.make((write) => + Effect.gen(function* () { + yield* Deferred.succeed(receive, write); + const serialization = yield* RpcSerialization.RpcSerialization; + return { + disconnects: yield* Queue.unbounded(), + send: (_clientId, response) => Queue.offer(responses, response), + end: () => Effect.void, + clientIds: Effect.succeed(new Set([0])), + initialMessage: Effect.succeedNone, + supportsAck: true, + supportsTransferables: false, + supportsSpanPropagation: false, + supportsNotifications: true, + codecFor: serialization.codecFor, + }; + }), + ); + yield* RpcServer.make(TestRpcs, WS_RPC_SERVER_OPTIONS).pipe( + Effect.provide( + TestRpcs.toLayer({ + subscribe: () => Stream.fromQueue(subscription), + boom: () => Effect.die(new Error("handler bug")), + }), + ), + Effect.provideService(RpcServer.Protocol, protocol), + Effect.forkScoped, + ); + const write = yield* Deferred.await(receive); + + yield* write(0, { _tag: "Request", id: "1", tag: "subscribe", payload: null, headers: [] }); + yield* Queue.offer(subscription, 1); + assert.deepEqual(yield* Queue.take(responses), { + _tag: "Chunk", + requestId: "1", + values: [1], + }); + yield* write(0, { _tag: "Ack", requestId: "1" }); + + yield* write(0, { _tag: "Request", id: "2", tag: "boom", payload: null, headers: [] }); + const boom = yield* Queue.take(responses); + assert.equal(boom._tag, "Exit"); + if (boom._tag === "Exit") { + assert.equal(boom.requestId, "2"); + assert.equal(boom.exit._tag, "Failure"); + } + const logged = yield* Queue.take(logs); + assert.isTrue(Cause.hasDies(logged)); + assert.equal(errorMessage(logged), "handler bug"); + + // The sibling subscription on the same client keeps delivering. + yield* Queue.offer(subscription, 2); + assert.deepEqual(yield* Queue.take(responses), { + _tag: "Chunk", + requestId: "1", + values: [2], + }); + assert.equal(yield* Queue.size(logs), 0); + }), + ).pipe(Effect.provide(RpcSerialization.layerJson), Effect.scoped), + ); + + it.effect("does not log typed failures, or a typed failure rethrown as a defect", () => + withErrorLogs((logs) => + Effect.gen(function* () { + // HttpApi reports a schema failure, then rethrows the same error as a defect. + const badRequest = new Error("bad request"); + yield* ErrorReporter.report(Cause.fail(badRequest)); + yield* ErrorReporter.report(Cause.die(badRequest)); + yield* ErrorReporter.report(Cause.fail(new Error("expected"))); + yield* ErrorReporter.report(Cause.die(new Error("bug"))); + + assert.equal(errorMessage(yield* Queue.take(logs)), "bug"); + assert.equal(yield* Queue.size(logs), 0); + }), + ), + ); +}); diff --git a/apps/server/src/observability/DefectReporter.ts b/apps/server/src/observability/DefectReporter.ts new file mode 100644 index 000000000000..e5d6efc87ccc --- /dev/null +++ b/apps/server/src/observability/DefectReporter.ts @@ -0,0 +1,41 @@ +import * as Cause from "effect/Cause"; +import * as Effect from "effect/Effect"; +import * as ErrorReporter from "effect/ErrorReporter"; + +/** + * Logs defects that reach a reporting boundary: RPC handlers, HTTP routes and + * HttpApi endpoints. Typed failures are expected responses, and HttpApi reports + * every one of them, so they are not logged. + * + * Errors first seen as typed failures are remembered. HttpApi reports a bad + * request as a failure and then rethrows it as a defect so the HTTP layer answers + * 400, and that second report must not be logged. Ignored errors, such as the + * response HttpEffect attaches to every failed request, are skipped. + */ +export const make = (): ErrorReporter.ErrorReporter => { + const seen = new WeakSet(); + return { + [ErrorReporter.TypeId]: ErrorReporter.TypeId, + report: ({ cause, fiber }) => { + if (seen.has(cause)) return; + seen.add(cause); + for (const reason of cause.reasons) { + if (reason._tag === "Interrupt") continue; + const value = reason._tag === "Fail" ? reason.error : reason.defect; + if (typeof value === "object" && value !== null) { + if (seen.has(value)) continue; + seen.add(value); + } + if (reason._tag === "Fail" || ErrorReporter.isIgnored(value)) 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 boundary that reported. + Effect.runForkWith(fiber.context)( + Effect.logError("Unhandled defect", Cause.fromReasons([reason])), + ); + } + }, + }; +}; + +export const layer = ErrorReporter.layer([Effect.sync(make)]); diff --git a/apps/server/src/observability/Layers/Observability.ts b/apps/server/src/observability/Layers/Observability.ts index 77ebe8410a91..fdb92176681d 100644 --- a/apps/server/src/observability/Layers/Observability.ts +++ b/apps/server/src/observability/Layers/Observability.ts @@ -17,6 +17,7 @@ import * as ServerConfig from "../../config.ts"; import * as ResourceAttribution from "../../resourceTelemetry/ResourceAttribution.ts"; import { ServerLoggerLive } from "../../serverLogger.ts"; import * as BrowserTraceCollector from "../BrowserTraceCollector.ts"; +import * as DefectReporter from "../DefectReporter.ts"; export const ObservabilityLive = Layer.unwrap( Effect.gen(function* () { @@ -95,7 +96,13 @@ export const ObservabilityLive = Layer.unwrap( return otelWarningsLayer.pipe( Layer.provideMerge( - Layer.mergeAll(ServerLoggerLive, traceReferencesLayer, tracerLayer, metricsLayer), + Layer.mergeAll( + ServerLoggerLive, + traceReferencesLayer, + tracerLayer, + metricsLayer, + DefectReporter.layer, + ), ), Layer.provide( OtelEnvironment.layerResourceAttributes(config.otelEnvironment.resourceAttributes), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index ae331995d201..77b7ca48505f 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -3763,6 +3763,14 @@ const makeWsRpcLayer = ( }), ); +// A handler defect 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 the defect either way. +export const WS_RPC_SERVER_OPTIONS = { + disableTracing: true, + disableFatalDefects: true, +} as const; + export const websocketRpcRouteLayer = Layer.unwrap( Effect.gen(function* () { const previewAutomationBroker = yield* PreviewAutomationBroker.PreviewAutomationBroker; @@ -3805,7 +3813,7 @@ export const websocketRpcRouteLayer = 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(rpcScopeAuthorizationLayer(session.scopes)), Effect.forkScoped, From 3d268331af9049d7ae6f826034031d558117c074 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 00:32:19 -0700 Subject: [PATCH 2/8] chore(server): keep the defect reporter constructor module-private Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/observability/DefectReporter.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/server/src/observability/DefectReporter.ts b/apps/server/src/observability/DefectReporter.ts index e5d6efc87ccc..10db53a0b208 100644 --- a/apps/server/src/observability/DefectReporter.ts +++ b/apps/server/src/observability/DefectReporter.ts @@ -12,7 +12,7 @@ import * as ErrorReporter from "effect/ErrorReporter"; * 400, and that second report must not be logged. Ignored errors, such as the * response HttpEffect attaches to every failed request, are skipped. */ -export const make = (): ErrorReporter.ErrorReporter => { +const make = (): ErrorReporter.ErrorReporter => { const seen = new WeakSet(); return { [ErrorReporter.TypeId]: ErrorReporter.TypeId, From 06944d707f06fc80bfcd326e5025dc1c8fff5778 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 23:28:49 -0700 Subject: [PATCH 3/8] fix(server): ACP requests whose handler dies get their own error The ACP client and agent RPC servers also sent a handler defect as a socket-level error with id -32603, so the peer's request never got an answer. They now use disableFatalDefects too. DefectReporter moves from the global observability layer to the WS RPC handlers. Installed globally it logged MCP tool defects a second time, since McpServer already logs them. It now logs only dies, so the HttpApi fail-then-die bookkeeping is gone. The RPC test pairs a real client with the server: with fatal defects the client ended the sibling subscription, which the server-only test could not see. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/observability/DefectReporter.test.ts | 100 +++++++----------- .../src/observability/DefectReporter.ts | 49 +++------ .../src/observability/Layers/Observability.ts | 9 +- apps/server/src/ws.ts | 10 +- packages/effect-acp/src/agent.test.ts | 33 ++++++ packages/effect-acp/src/agent.ts | 2 +- packages/effect-acp/src/client.test.ts | 36 +++++++ packages/effect-acp/src/client.ts | 2 +- 8 files changed, 135 insertions(+), 106 deletions(-) diff --git a/apps/server/src/observability/DefectReporter.test.ts b/apps/server/src/observability/DefectReporter.test.ts index 0a7ea4a92369..66ad2f32e570 100644 --- a/apps/server/src/observability/DefectReporter.test.ts +++ b/apps/server/src/observability/DefectReporter.test.ts @@ -1,14 +1,15 @@ import { assert, describe, it } from "@effect/vitest"; import * as Cause from "effect/Cause"; -import * as Deferred from "effect/Deferred"; 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, RpcGroup, RpcMessage, RpcSerialization, RpcServer } from "effect/unstable/rpc"; +import { Rpc, RpcClient, RpcGroup, RpcServer } from "effect/unstable/rpc"; import { WS_RPC_SERVER_OPTIONS } from "../ws.ts"; import * as DefectReporter from "./DefectReporter.ts"; @@ -17,8 +18,9 @@ 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 the defect reporter installed and every error log captured. */ +/** Runs `body` with every error log captured. */ const withErrorLogs = ( body: (logs: Queue.Queue>) => Effect.Effect, ) => @@ -28,9 +30,7 @@ const withErrorLogs = ( if (logLevel === "Error") Queue.offerUnsafe(logs, cause); }); return yield* body(logs).pipe( - Effect.provide( - Layer.merge(DefectReporter.layer, Logger.layer([logger], { mergeWithExisting: false })), - ), + Effect.provide(Logger.layer([logger], { mergeWithExisting: false })), ); }); @@ -41,83 +41,63 @@ describe("DefectReporter", () => { withErrorLogs((logs) => Effect.gen(function* () { const subscription = yield* Queue.unbounded(); - const responses = yield* Queue.unbounded(); - const receive = yield* Deferred.make[0]>(); - const protocol = yield* RpcServer.Protocol.make((write) => - Effect.gen(function* () { - yield* Deferred.succeed(receive, write); - const serialization = yield* RpcSerialization.RpcSerialization; - return { - disconnects: yield* Queue.unbounded(), - send: (_clientId, response) => Queue.offer(responses, response), - end: () => Effect.void, - clientIds: Effect.succeed(new Set([0])), - initialMessage: Effect.succeedNone, - supportsAck: true, - supportsTransferables: false, - supportsSpanPropagation: false, - supportsNotifications: true, - codecFor: serialization.codecFor, - }; - }), - ); - yield* RpcServer.make(TestRpcs, WS_RPC_SERVER_OPTIONS).pipe( + // The same client/server pairing as RpcTest.makeClient, with ws.ts's options. + // oxlint-disable-next-line prefer-const + 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)), ), - Effect.provideService(RpcServer.Protocol, protocol), - Effect.forkScoped, ); - const write = yield* Deferred.await(receive); + client = yield* RpcClient.makeNoSerialization(TestRpcs, { + supportsAck: true, + onFromClient: ({ message }) => server.write(0, message), + }); - yield* write(0, { _tag: "Request", id: "1", tag: "subscribe", payload: null, headers: [] }); + 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.deepEqual(yield* Queue.take(responses), { - _tag: "Chunk", - requestId: "1", - values: [1], - }); - yield* write(0, { _tag: "Ack", requestId: "1" }); + assert.equal(yield* Queue.take(received), 1); - yield* write(0, { _tag: "Request", id: "2", tag: "boom", payload: null, headers: [] }); - const boom = yield* Queue.take(responses); - assert.equal(boom._tag, "Exit"); - if (boom._tag === "Exit") { - assert.equal(boom.requestId, "2"); - assert.equal(boom.exit._tag, "Failure"); - } - const logged = yield* Queue.take(logs); - assert.isTrue(Cause.hasDies(logged)); - assert.equal(errorMessage(logged), "handler bug"); + 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); - assert.deepEqual(yield* Queue.take(responses), { - _tag: "Chunk", - requestId: "1", - values: [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.provide(RpcSerialization.layerJson), Effect.scoped), + ).pipe(Effect.scoped), ); - it.effect("does not log typed failures, or a typed failure rethrown as a defect", () => + it.effect("logs defects, not typed failures or interrupts", () => withErrorLogs((logs) => Effect.gen(function* () { - // HttpApi reports a schema failure, then rethrows the same error as a defect. - const badRequest = new Error("bad request"); - yield* ErrorReporter.report(Cause.fail(badRequest)); - yield* ErrorReporter.report(Cause.die(badRequest)); yield* ErrorReporter.report(Cause.fail(new Error("expected"))); + yield* ErrorReporter.report(Cause.interrupt()); yield* ErrorReporter.report(Cause.die(new Error("bug"))); assert.equal(errorMessage(yield* Queue.take(logs)), "bug"); assert.equal(yield* Queue.size(logs), 0); - }), + }).pipe(Effect.provide(DefectReporter.layer)), ), ); }); diff --git a/apps/server/src/observability/DefectReporter.ts b/apps/server/src/observability/DefectReporter.ts index 10db53a0b208..220b2696d831 100644 --- a/apps/server/src/observability/DefectReporter.ts +++ b/apps/server/src/observability/DefectReporter.ts @@ -3,39 +3,22 @@ import * as Effect from "effect/Effect"; import * as ErrorReporter from "effect/ErrorReporter"; /** - * Logs defects that reach a reporting boundary: RPC handlers, HTTP routes and - * HttpApi endpoints. Typed failures are expected responses, and HttpApi reports - * every one of them, so they are not logged. - * - * Errors first seen as typed failures are remembered. HttpApi reports a bad - * request as a failure and then rethrows it as a defect so the HTTP layer answers - * 400, and that second report must not be logged. Ignored errors, such as the - * response HttpEffect attaches to every failed request, are skipped. + * Logs the defects of WebSocket RPC handlers. RpcServer reports every failed + * handler exit; typed failures are expected responses, so only dies are logged. */ -const make = (): ErrorReporter.ErrorReporter => { - const seen = new WeakSet(); - return { - [ErrorReporter.TypeId]: ErrorReporter.TypeId, - report: ({ cause, fiber }) => { - if (seen.has(cause)) return; - seen.add(cause); - for (const reason of cause.reasons) { - if (reason._tag === "Interrupt") continue; - const value = reason._tag === "Fail" ? reason.error : reason.defect; - if (typeof value === "object" && value !== null) { - if (seen.has(value)) continue; - seen.add(value); - } - if (reason._tag === "Fail" || ErrorReporter.isIgnored(value)) 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 boundary that reported. - Effect.runForkWith(fiber.context)( - Effect.logError("Unhandled defect", Cause.fromReasons([reason])), - ); - } - }, - }; +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([Effect.sync(make)]); +export const layer = ErrorReporter.layer([reporter]); diff --git a/apps/server/src/observability/Layers/Observability.ts b/apps/server/src/observability/Layers/Observability.ts index fdb92176681d..77ebe8410a91 100644 --- a/apps/server/src/observability/Layers/Observability.ts +++ b/apps/server/src/observability/Layers/Observability.ts @@ -17,7 +17,6 @@ import * as ServerConfig from "../../config.ts"; import * as ResourceAttribution from "../../resourceTelemetry/ResourceAttribution.ts"; import { ServerLoggerLive } from "../../serverLogger.ts"; import * as BrowserTraceCollector from "../BrowserTraceCollector.ts"; -import * as DefectReporter from "../DefectReporter.ts"; export const ObservabilityLive = Layer.unwrap( Effect.gen(function* () { @@ -96,13 +95,7 @@ export const ObservabilityLive = Layer.unwrap( return otelWarningsLayer.pipe( Layer.provideMerge( - Layer.mergeAll( - ServerLoggerLive, - traceReferencesLayer, - tracerLayer, - metricsLayer, - DefectReporter.layer, - ), + Layer.mergeAll(ServerLoggerLive, traceReferencesLayer, tracerLayer, metricsLayer), ), Layer.provide( OtelEnvironment.layerResourceAttributes(config.otelEnvironment.resourceAttributes), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 77b7ca48505f..16a9114b68e5 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -205,6 +205,7 @@ import * as RepositoryIdentityResolver from "./project/RepositoryIdentityResolve import * as WorktreeSetupTracker from "./project/WorktreeSetupTracker.ts"; import * as ServerEnvironment from "./environment/ServerEnvironment.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 { @@ -3763,9 +3764,9 @@ const makeWsRpcLayer = ( }), ); -// A handler defect 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 the defect either way. +// 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, @@ -3829,6 +3830,9 @@ export const websocketRpcRouteLayer = Layer.unwrap( previewAutomationBroker, ).pipe( Layer.provideMerge(RpcSerialization.layerJson), + // Request fibers run in the handlers' context, so this reporter sees + // their defects and nothing else on the server. + 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/agent.test.ts b/packages/effect-acp/src/agent.test.ts index 9b69ecb5918d..76ceeffe59a6 100644 --- a/packages/effect-acp/src/agent.test.ts +++ b/packages/effect-acp/src/agent.test.ts @@ -35,6 +35,8 @@ 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 JSON-RPC error response; only its id and the presence of an error matter here. */ +const ErrorResponse = Schema.Struct({ id: Schema.Number, error: Schema.Unknown }); const decodeRequestPermissionRequest = Schema.decodeEffect( Schema.fromJsonString(RequestPermissionRequest), ); @@ -296,3 +298,34 @@ 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(ErrorResponse))), + ); + assert.equal(response.id, 7); + assert.isNotNull(response.error); + yield* Scope.close(scope, Exit.void); + }), +); diff --git a/packages/effect-acp/src/agent.ts b/packages/effect-acp/src/agent.ts index 2cd9a3ef8b4c..e680a12cb744 100644 --- a/packages/effect-acp/src/agent.ts +++ b/packages/effect-acp/src/agent.ts @@ -421,7 +421,7 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* ( }), ); - yield* RpcServer.make(AcpRpcs.AgentRpcs).pipe( + yield* RpcServer.make(AcpRpcs.AgentRpcs, { disableFatalDefects: true }).pipe( Effect.provideService(RpcServer.Protocol, transport.serverProtocol), Effect.provide(agentHandlerLayer), Effect.forkScoped, diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index 94a1cccdd879..6e3416149ea6 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -50,6 +50,8 @@ 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 }); const decodePromptRequestLine = Schema.decodeEffect(Schema.fromJsonString(PromptRequest)); const XAiPromptCompleteNotification = jsonRpcNotification( "_x.ai/session/prompt_complete", @@ -1130,6 +1132,40 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { }), ); + it.effect("answers a request whose handler dies with an error for that request", () => + Effect.gen(function* () { + const { stdio, input, output } = yield* makeInMemoryStdio(); + const scope = yield* Scope.make(); + const acp = yield* AcpClient.make(stdio).pipe(Effect.provideService(Scope.Scope, scope)); + yield* acp.handleRequestPermission(() => Effect.die(new Error("handler bug"))); + + yield* Queue.offer( + input, + yield* 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: [], + }), + ); + const response = yield* Queue.take(output).pipe( + Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ErrorResponse))), + ); + assert.equal(response.id, "permission-a"); + assert.isNotNull(response.error); + 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 593fdd443095..8b4da1647807 100644 --- a/packages/effect-acp/src/client.ts +++ b/packages/effect-acp/src/client.ts @@ -1099,7 +1099,7 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* ( }), ); - yield* RpcServer.make(AcpRpcs.CompatClientRpcs).pipe( + yield* RpcServer.make(AcpRpcs.CompatClientRpcs, { disableFatalDefects: true }).pipe( Effect.provideService(RpcServer.Protocol, transport.serverProtocol), Effect.provide(clientHandlerLayer), Effect.forkScoped, From f963d366bdd60f1c15e5b958278da53354d23a86 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Sun, 4 Oct 2026 23:33:36 -0700 Subject: [PATCH 4/8] test(server): drop an unused lint directive Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/observability/DefectReporter.test.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/apps/server/src/observability/DefectReporter.test.ts b/apps/server/src/observability/DefectReporter.test.ts index 66ad2f32e570..4b556d96fcee 100644 --- a/apps/server/src/observability/DefectReporter.test.ts +++ b/apps/server/src/observability/DefectReporter.test.ts @@ -42,7 +42,6 @@ describe("DefectReporter", () => { Effect.gen(function* () { const subscription = yield* Queue.unbounded(); // The same client/server pairing as RpcTest.makeClient, with ws.ts's options. - // oxlint-disable-next-line prefer-const let client!: Effect.Success< ReturnType> >; From 20fddd2a259a748abcf88406c8ef3fac16c4c150 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 02:13:08 -0700 Subject: [PATCH 5/8] fix(effect-acp): a dying extension or notification handler no longer stops the reader disableFatalDefects only covers core methods, which go through RpcServer. Extension requests and notifications are handled in the protocol layer, where Effect.ignore and matchEffect let a defect through and end the stdin reader: that request and every later one on the connection went unanswered, and termination never fired. A dying extension handler now answers its own request with an internal error, and notification handlers are guarded on their whole cause. The ACP tests now assert the defect itself, and the client test checks that a later request is still answered after each kind of handler dies. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../src/observability/DefectReporter.test.ts | 4 +- packages/effect-acp/src/agent.test.ts | 13 +++- packages/effect-acp/src/client.test.ts | 75 +++++++++++++++++-- packages/effect-acp/src/client.ts | 12 ++- packages/effect-acp/src/protocol.ts | 31 ++++++-- 5 files changed, 112 insertions(+), 23 deletions(-) diff --git a/apps/server/src/observability/DefectReporter.test.ts b/apps/server/src/observability/DefectReporter.test.ts index 4b556d96fcee..cce2e294a406 100644 --- a/apps/server/src/observability/DefectReporter.test.ts +++ b/apps/server/src/observability/DefectReporter.test.ts @@ -94,8 +94,8 @@ describe("DefectReporter", () => { yield* ErrorReporter.report(Cause.interrupt()); yield* ErrorReporter.report(Cause.die(new Error("bug"))); - assert.equal(errorMessage(yield* Queue.take(logs)), "bug"); - assert.equal(yield* Queue.size(logs), 0); + // 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/packages/effect-acp/src/agent.test.ts b/packages/effect-acp/src/agent.test.ts index 76ceeffe59a6..416c4808a3ed 100644 --- a/packages/effect-acp/src/agent.test.ts +++ b/packages/effect-acp/src/agent.test.ts @@ -35,8 +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 JSON-RPC error response; only its id and the presence of an error matter here. */ -const ErrorResponse = Schema.Struct({ id: Schema.Number, error: Schema.Unknown }); +/** 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), ); @@ -322,10 +328,9 @@ it.effect("effect-acp agent answers a request whose handler dies with an error f }), ); const response = yield* Queue.take(output).pipe( - Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ErrorResponse))), + Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(DieResponse))), ); assert.equal(response.id, 7); - assert.isNotNull(response.error); yield* Scope.close(scope, Exit.void); }), ); diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index 6e3416149ea6..639f401405b8 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -52,6 +52,14 @@ const ElicitationRequest = jsonRpcRequest("elicitation/create", AcpSchema.Create 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", @@ -1132,16 +1140,28 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { }), ); - it.effect("answers a request whose handler dies with an error for that request", () => + it.effect("answers each request whose handler dies, and keeps reading", () => Effect.gen(function* () { const { stdio, input, output } = yield* makeInMemoryStdio(); const scope = yield* Scope.make(); const acp = yield* AcpClient.make(stdio).pipe(Effect.provideService(Scope.Scope, scope)); - yield* acp.handleRequestPermission(() => Effect.die(new Error("handler bug"))); + const bug = () => Effect.die(new Error("handler bug")); + yield* acp.handleRequestPermission(bug); + yield* acp.handleExtRequest("x/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)); - yield* Queue.offer( - input, - yield* encodeJsonl(PermissionRequest, { + yield* send( + encodeJsonl(PermissionRequest, { jsonrpc: "2.0", id: "permission-a", method: "session/request_permission", @@ -1157,11 +1177,50 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { headers: [], }), ); - const response = yield* Queue.take(output).pipe( + assert.equal((yield* Queue.take(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* Queue.take(output).pipe( Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ErrorResponse))), ); - assert.equal(response.id, "permission-a"); - assert.isNotNull(response.error); + assert.equal(extDied.id, "ext-a"); + + 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* Queue.take(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* Queue.take(output).pipe( + Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ExtResponse))), + ); + assert.deepEqual([answered.id, answered.result], ["ext-b", { ok: true }]); yield* Scope.close(scope, Exit.void); }), ); diff --git a/packages/effect-acp/src/client.ts b/packages/effect-acp/src/client.ts index 8b4da1647807..cc7584100628 100644 --- a/packages/effect-acp/src/client.ts +++ b/packages/effect-acp/src/client.ts @@ -838,9 +838,15 @@ 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) => + handler(notification).pipe( + Effect.ignoreCause({ log: true, message: "ACP notification handler failed" }), + ), + { discard: true }, + ); const flushBufferedNotifications = (registration: BufferedNotificationHandler) => Effect.suspend(() => { diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index 354de5d656a6..72a6138ad0da 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"; @@ -366,8 +367,12 @@ 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) + ? options + .onNotification(notification) + .pipe(Effect.ignoreCause({ log: true, message: "ACP notification handler failed" })) : Effect.void, ), Effect.asVoid, @@ -439,12 +444,26 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi method: message.tag, }) .pipe( - Effect.matchEffect({ - onFailure: (error) => - respondWithError( + Effect.matchCauseEffect({ + // A dying handler answers its own request, like a core handler, and + // leaves the reader running. + onFailure: (cause) => { + const failure = Cause.findErrorOption(cause); + return respondWithError( message.id, - AcpError.AcpRequestError.fromExtensionHandlerError(error, message.tag), - ), + Option.isSome(failure) + ? AcpError.AcpRequestError.fromExtensionHandlerError(failure.value, message.tag) + : 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), }), ); From 3afd7649eed9393cb2dbe31ad383a6732c49191f Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 04:30:58 -0700 Subject: [PATCH 6/8] fix(effect-acp): log handler defects at Error and terminate on reader defects - A defect in a core or extension request handler is logged before the agent gets its error; before, it reached no log. - Notification handlers drop typed failures silently, as before, and log defects at Error. ignoreCause logged both at Info. - A defect that ends the reader outside any handler (a session update normalizer, a logger) now terminates the connection, so pending requests fail instead of hanging. Its own interrupt still doesn't. - The agent's cancel handlers are isolated from each other. The client test also covers a dying extension notification and checks each defect is logged once at Error; a protocol test covers a dying normalizer. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/effect-acp/src/_internal/shared.ts | 17 ++++++++++ packages/effect-acp/src/agent.ts | 7 ++++- packages/effect-acp/src/client.test.ts | 26 ++++++++++++++- packages/effect-acp/src/client.ts | 6 ++-- packages/effect-acp/src/protocol.test.ts | 34 ++++++++++++++++++++ packages/effect-acp/src/protocol.ts | 35 +++++++++++++++------ 6 files changed, 109 insertions(+), 16 deletions(-) diff --git a/packages/effect-acp/src/_internal/shared.ts b/packages/effect-acp/src/_internal/shared.ts index 96dd8f7a9a48..eb3fc4638bd0 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/unstable/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.ts b/packages/effect-acp/src/agent.ts index e680a12cb744..e51dfd87877d 100644 --- a/packages/effect-acp/src/agent.ts +++ b/packages/effect-acp/src/agent.ts @@ -18,6 +18,7 @@ import { callRpc, decodeExtNotificationRegistration, decodeExtRequestRegistration, + isolateNotificationHandler, runHandler, } from "./_internal/shared.ts"; @@ -280,7 +281,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 }, + ), ), ); } diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index 639f401405b8..a35343483785 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -5,6 +5,7 @@ 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"; @@ -1142,12 +1143,20 @@ 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)); + 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) => @@ -1193,6 +1202,13 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { ); 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", @@ -1221,6 +1237,14 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { 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); }), ); diff --git a/packages/effect-acp/src/client.ts b/packages/effect-acp/src/client.ts index cc7584100628..aafbd8990c73 100644 --- a/packages/effect-acp/src/client.ts +++ b/packages/effect-acp/src/client.ts @@ -22,6 +22,7 @@ import { callRpc, decodeExtNotificationRegistration, decodeExtRequestRegistration, + isolateNotificationHandler, runHandler, } from "./_internal/shared.ts"; import { makeChildStdio, makeTerminationError } from "./_internal/stdio.ts"; @@ -841,10 +842,7 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* ( // One handler failing or dying does not stop the others, or the reader. Effect.forEach( registration.handlers, - (handler) => - handler(notification).pipe( - Effect.ignoreCause({ log: true, message: "ACP notification handler failed" }), - ), + (handler) => isolateNotificationHandler(handler(notification)), { discard: true }, ); diff --git a/packages/effect-acp/src/protocol.test.ts b/packages/effect-acp/src/protocol.test.ts index 051f03627f03..5db018e6c429 100644 --- a/packages/effect-acp/src/protocol.test.ts +++ b/packages/effect-acp/src/protocol.test.ts @@ -257,6 +257,40 @@ 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* Deferred.await(termination); + 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(); diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index 72a6138ad0da..66910b98e395 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -21,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 { @@ -370,9 +371,7 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi // 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.ignoreCause({ log: true, message: "ACP notification handler failed" })) + ? isolateNotificationHandler(options.onNotification(notification)) : Effect.void, ), Effect.asVoid, @@ -449,11 +448,20 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi // leaves the reader running. onFailure: (cause) => { const failure = Cause.findErrorOption(cause); - return respondWithError( - message.id, - Option.isSome(failure) - ? AcpError.AcpRequestError.fromExtensionHandlerError(failure.value, message.tag) - : AcpError.AcpRequestError.internalError( + 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, { @@ -462,6 +470,8 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi cause: Cause.squash(cause), }, ), + ), + ), ); }, onSuccess: (value) => respondWithSuccess(message.id, value), @@ -685,8 +695,13 @@ 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) => { + const failure = Cause.findErrorOption(cause); + if (Option.isNone(failure) && Cause.hasInterruptsOnly(cause)) return Effect.void; + const error = Option.isSome(failure) ? failure.value : Cause.squash(cause); const normalized: AcpError.AcpError = isAcpError(error) ? error : new AcpError.AcpTransportError({ From fc9ee1649b2451fb42ff43c0f78a5b6d6b108d3d Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 07:47:53 -0700 Subject: [PATCH 7/8] fix(effect-acp): an interrupt inside the reader terminates the connection The reader's own interruption never reaches its failure handler, so the interrupt-only shortcut only caught interrupts raised inside it, and ended the reader without terminating. Drop it. ACP RpcServers no longer inherit an ErrorReporter from the fiber that built them. Built inside a WebSocket request they picked up DefectReporter and logged handler defects a second time; runHandler already logs them with their method. The new waits in the ACP tests are bounded, so a regression fails in seconds instead of at the test timeout. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/ws.ts | 2 +- packages/effect-acp/src/agent.ts | 4 ++++ packages/effect-acp/src/client.test.ts | 12 ++++++++---- packages/effect-acp/src/client.ts | 4 ++++ packages/effect-acp/src/protocol.test.ts | 5 ++++- packages/effect-acp/src/protocol.ts | 3 ++- 6 files changed, 23 insertions(+), 7 deletions(-) diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 16a9114b68e5..e2f6364891bf 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -3831,7 +3831,7 @@ export const websocketRpcRouteLayer = Layer.unwrap( ).pipe( Layer.provideMerge(RpcSerialization.layerJson), // Request fibers run in the handlers' context, so this reporter sees - // their defects and nothing else on the server. + // their defects, not the rest of the server's. Layer.provide(DefectReporter.layer), Layer.provide(Layer.succeed(SqlClient.SqlClient, sql)), Layer.provide(AgentSessionScanner.layer), diff --git a/packages/effect-acp/src/agent.ts b/packages/effect-acp/src/agent.ts index e51dfd87877d..30b8182bcf1b 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"; @@ -427,6 +428,9 @@ export const make = Effect.fn("effect-acp/AcpAgent.make")(function* ( ); 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(agentHandlerLayer), Effect.forkScoped, diff --git a/packages/effect-acp/src/client.test.ts b/packages/effect-acp/src/client.test.ts index a35343483785..228fd41a5af2 100644 --- a/packages/effect-acp/src/client.test.ts +++ b/packages/effect-acp/src/client.test.ts @@ -11,6 +11,7 @@ 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/unstable/process"; import * as NodeServices from "@effect/platform-node/NodeServices"; @@ -1168,6 +1169,9 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { 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, { @@ -1186,7 +1190,7 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { headers: [], }), ); - assert.equal((yield* Queue.take(output).pipe(Effect.flatMap(decodeDie))).id, "permission-a"); + assert.equal((yield* next(output).pipe(Effect.flatMap(decodeDie))).id, "permission-a"); yield* send( encodeJsonl(jsonRpcRequest("x/dies", Schema.Unknown), { @@ -1197,7 +1201,7 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { headers: [], }), ); - const extDied = yield* Queue.take(output).pipe( + const extDied = yield* next(output).pipe( Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ErrorResponse))), ); assert.equal(extDied.id, "ext-a"); @@ -1221,7 +1225,7 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { ); // The next session handler still ran. - assert.equal(yield* Queue.take(updates), "session-1"); + assert.equal(yield* next(updates), "session-1"); // The reader survived all three: a later request is still answered. yield* send( @@ -1233,7 +1237,7 @@ it.layer(NodeServices.layer)("effect-acp client", (it) => { headers: [], }), ); - const answered = yield* Queue.take(output).pipe( + const answered = yield* next(output).pipe( Effect.flatMap(Schema.decodeEffect(Schema.fromJsonString(ExtResponse))), ); assert.deepEqual([answered.id, answered.result], ["ext-b", { ok: true }]); diff --git a/packages/effect-acp/src/client.ts b/packages/effect-acp/src/client.ts index aafbd8990c73..98f3ef549b19 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"; @@ -1104,6 +1105,9 @@ export const make = Effect.fn("effect-acp/AcpClient.make")(function* ( ); 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(clientHandlerLayer), Effect.forkScoped, diff --git a/packages/effect-acp/src/protocol.test.ts b/packages/effect-acp/src/protocol.test.ts index 5db018e6c429..389601bf6d89 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/unstable/process"; import { it, assert } from "@effect/vitest"; @@ -285,7 +286,9 @@ it.layer(NodeServices.layer)("effect-acp protocol", (it) => { ); // Pending requests are failed through termination instead of hanging. - const error = yield* Deferred.await(termination); + 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"); }), diff --git a/packages/effect-acp/src/protocol.ts b/packages/effect-acp/src/protocol.ts index 66910b98e395..5129b242b738 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -699,8 +699,9 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi // 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); - if (Option.isNone(failure) && Cause.hasInterruptsOnly(cause)) return Effect.void; const error = Option.isSome(failure) ? failure.value : Cause.squash(cause); const normalized: AcpError.AcpError = isAcpError(error) ? error From 6d36ef769dfe110056f95f2f96f88291eafcabd2 Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Tue, 6 Oct 2026 11:54:17 -0700 Subject: [PATCH 8/8] fix(effect-acp): a dying extension handler is never reported as its typed error When an extension request handler failed with an ACP error and its cleanup then died, the response reported the typed error and the defect was neither answered nor logged. A defect now wins: the request gets an internal error and the defect is logged. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/effect-acp/src/protocol.test.ts | 32 ++++++++++++++++++++++++ packages/effect-acp/src/protocol.ts | 5 ++-- 2 files changed, 35 insertions(+), 2 deletions(-) diff --git a/packages/effect-acp/src/protocol.test.ts b/packages/effect-acp/src/protocol.test.ts index c8ee1d892385..9c3f41e5f4c9 100644 --- a/packages/effect-acp/src/protocol.test.ts +++ b/packages/effect-acp/src/protocol.test.ts @@ -537,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 006df1149df6..12e85a9e7542 100644 --- a/packages/effect-acp/src/protocol.ts +++ b/packages/effect-acp/src/protocol.ts @@ -464,9 +464,10 @@ export const makeAcpPatchedProtocol = Effect.fn("makeAcpPatchedProtocol")(functi .pipe( Effect.matchCauseEffect({ // A dying handler answers its own request, like a core handler, and - // leaves the reader running. + // 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.findErrorOption(cause); + const failure = Cause.hasDies(cause) ? Option.none() : Cause.findErrorOption(cause); if (Option.isSome(failure)) { return respondWithError( message.id,