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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
102 changes: 102 additions & 0 deletions apps/server/src/observability/DefectReporter.test.ts
Original file line number Diff line number Diff line change
@@ -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<typeof TestRpcs>;

/** Runs `body` with every error log captured. */
const withErrorLogs = <A, E, R>(
body: (logs: Queue.Queue<Cause.Cause<unknown>>) => Effect.Effect<A, E, R>,
) =>
Effect.gen(function* () {
const logs = yield* Queue.unbounded<Cause.Cause<unknown>>();
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<unknown>) => (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<number>();
// The same client/server pairing as RpcTest.makeClient, with ws.ts's options.
let client!: Effect.Success<
ReturnType<typeof RpcClient.makeNoSerialization<TestRpc, never>>
>;
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<number>();
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)),
),
);
});
24 changes: 24 additions & 0 deletions apps/server/src/observability/DefectReporter.ts
Original file line number Diff line number Diff line change
@@ -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]);
14 changes: 13 additions & 1 deletion apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand All @@ -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),
Expand Down
17 changes: 17 additions & 0 deletions packages/effect-acp/src/_internal/shared.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -29,6 +30,19 @@ export const callRpc = <A>(
}),
);

/**
* 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 = <E, R>(effect: Effect.Effect<void, E, R>) =>
effect.pipe(
Effect.catchCause((cause) =>
Cause.hasDies(cause)
? Effect.logError("ACP notification handler failed", cause)
: Effect.void,
),
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
);

export const runHandler = Effect.fnUntraced(function* <A, B, Args extends ReadonlyArray<unknown>>(
handler: ((payload: A, ...args: Args) => Effect.Effect<B, AcpError.AcpError>) | undefined,
payload: A,
Expand All @@ -39,6 +53,9 @@ export const runHandler = Effect.fnUntraced(function* <A, B, Args extends Readon
return yield* Effect.fail(AcpError.AcpRequestError.methodNotFound(method).toProtocolError());
}
return yield* handler(payload, ...args).pipe(
Effect.tapDefect((defect) =>
Effect.logError(`ACP request handler failed for '${method}'`, defect),
),
Effect.mapError((error) =>
AcpError.AcpRequestError.fromCoreHandlerError(error, method).toProtocolError(),
),
Expand Down
38 changes: 38 additions & 0 deletions packages/effect-acp/src/agent.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
);
Expand Down Expand Up @@ -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);
}),
);
13 changes: 11 additions & 2 deletions packages/effect-acp/src/agent.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -18,6 +19,7 @@ import {
callRpc,
decodeExtNotificationRegistration,
decodeExtRequestRegistration,
isolateNotificationHandler,
runHandler,
} from "./_internal/shared.ts";

Expand Down Expand Up @@ -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 },
),
),
);
}
Expand Down Expand Up @@ -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,
Expand Down
Loading
Loading