From d8b1157f912539fe8b05e50b9239300bcc9e4672 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Mon, 28 Sep 2026 09:41:14 -0400 Subject: [PATCH] fix(server): RPC spans join the client's trace Signed-off-by: Yordis Prieto --- .../observability/RpcInstrumentation.test.ts | 74 ++++++++++++++++--- .../src/observability/RpcInstrumentation.ts | 28 +++---- apps/server/src/ws.ts | 3 +- 3 files changed, 78 insertions(+), 27 deletions(-) diff --git a/apps/server/src/observability/RpcInstrumentation.test.ts b/apps/server/src/observability/RpcInstrumentation.test.ts index c7fb7d12b095..5f9657927da0 100644 --- a/apps/server/src/observability/RpcInstrumentation.test.ts +++ b/apps/server/src/observability/RpcInstrumentation.test.ts @@ -5,14 +5,18 @@ import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; import * as Metric from "effect/Metric"; +import * as Option from "effect/Option"; +import * as Schema from "effect/Schema"; import * as Stream from "effect/Stream"; import * as Tracer from "effect/Tracer"; import * as TestClock from "effect/testing/TestClock"; +import { Rpc, RpcClient, RpcGroup, RpcServer } from "effect/unstable/rpc"; import { observeRpcEffect, observeRpcStream, observeRpcStreamEffect, + rpcServerTracingOptions, } from "./RpcInstrumentation.ts"; const hasMetricSnapshot = ( @@ -236,22 +240,68 @@ describe("RpcInstrumentation", () => { }), ); - it.effect("records spans for traced stream RPC handlers", () => + it.effect("runs traced RPC handlers inside the server span parented to the client span", () => Effect.gen(function* () { - const spanNames = yield* collectSpanNames( - Stream.runCollect( + const spans: Array = []; + const tracer = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + spans.push(span); + return span; + }, + }); + const group = RpcGroup.make( + Rpc.make("Unary", { success: Schema.String }), + Rpc.make("Streamed", { success: Schema.String, stream: true }), + ); + const handlers = group.toLayer({ + Unary: () => + observeRpcEffect("Unary", Effect.succeed("ok").pipe(Effect.withSpan("unary.child")), { + "rpc.aggregate": "test", + }), + Streamed: () => observeRpcStream( - "rpc.instrumentation.traced.stream", - Stream.fromEffect( - Effect.succeed("ok").pipe(Effect.withSpan("rpc.instrumentation.traced.stream.child")), - ), - { "rpc.aggregate": "test" }, + "Streamed", + Stream.fromEffect(Effect.succeed("ok").pipe(Effect.withSpan("streamed.child"))), ), - ), - ); + }); - assert.equal(spanNames.includes("ws.rpc.rpc.instrumentation.traced.stream"), true); - assert.equal(spanNames.includes("rpc.instrumentation.traced.stream.child"), true); + yield* Effect.gen(function* () { + let client!: Effect.Success< + ReturnType, never>> + >; + const server = yield* RpcServer.makeNoSerialization(group, { + ...rpcServerTracingOptions, + onFromServer: (response) => client.write(response), + }); + client = yield* RpcClient.makeNoSerialization(group, { + supportsAck: true, + onFromClient: ({ message }) => server.write(0, message), + }); + yield* client.client.Unary(); + yield* Stream.runCollect(client.client.Streamed()); + }).pipe(Effect.provide(handlers), Effect.withTracer(tracer), Effect.scoped); + + const spanNamed = (name: string) => spans.find((span) => span.name === name); + for (const [method, child] of [ + ["Unary", "unary.child"], + ["Streamed", "streamed.child"], + ] as const) { + const clientSpan = spanNamed(`RpcClient.${method}`); + const serverSpan = spanNamed(`ws.rpc.${method}`); + assert.ok(clientSpan && serverSpan); + assert.equal(serverSpan.traceId, clientSpan.traceId); + assert.equal(Option.getOrUndefined(serverSpan.parent)?.spanId, clientSpan.spanId); + assert.equal(serverSpan.attributes.get("rpc.method"), method); + assert.equal(serverSpan.attributes.get("rpc.system"), "effect-rpc"); + assert.equal(spans.filter((span) => span.name.endsWith(`.${method}`)).length, 2); + const childParent = Option.flatMap( + Option.fromNullishOr(spanNamed(child)), + (span) => span.parent, + ); + assert.equal(Option.getOrUndefined(childParent)?.spanId, serverSpan.spanId); + } + assert.equal(spanNamed("ws.rpc.Unary")?.attributes.get("rpc.aggregate"), "test"); }), ); diff --git a/apps/server/src/observability/RpcInstrumentation.ts b/apps/server/src/observability/RpcInstrumentation.ts index edbd705b3ee4..fb8cbff4317a 100644 --- a/apps/server/src/observability/RpcInstrumentation.ts +++ b/apps/server/src/observability/RpcInstrumentation.ts @@ -10,10 +10,17 @@ import * as Stream from "effect/Stream"; import { outcomeFromExit } from "./Attributes.ts"; import { metricAttributes, rpcRequestDuration, rpcRequestsTotal, withMetrics } from "./Metrics.ts"; -const RPC_SPAN_PREFIX = "ws.rpc"; -const DEFAULT_RPC_SPAN_ATTRIBUTES = { - "rpc.transport": "websocket", - "rpc.system": "effect-rpc", +/** + * Passed to `RpcServer.make` so the server opens each request's span as a child + * of the span the client sent with it. The observe helpers below only annotate + * that span. + */ +export const rpcServerTracingOptions = { + spanPrefix: "ws.rpc", + spanAttributes: { + "rpc.transport": "websocket", + "rpc.system": "effect-rpc", + }, } as const; const RPC_METHODS_WITH_TRACING_DISABLED: ReadonlySet = new Set([ WS_METHODS.serverGetTraceDiagnostics, @@ -30,7 +37,6 @@ const rpcSpanAttributes = ( method: string, traceAttributes?: Readonly>, ): Record => ({ - ...DEFAULT_RPC_SPAN_ATTRIBUTES, "rpc.method": method, ...traceAttributes, }); @@ -41,11 +47,7 @@ const withRpcEffectTracing = ( traceAttributes?: Readonly>, ): Effect.Effect => shouldTraceRpc(method) - ? effect.pipe( - Effect.withSpan(`${RPC_SPAN_PREFIX}.${method}`, { - attributes: rpcSpanAttributes(method, traceAttributes), - }), - ) + ? Effect.andThen(Effect.annotateCurrentSpan(rpcSpanAttributes(method, traceAttributes)), effect) : effect.pipe(Effect.provideService(References.TracerEnabled, false)); const withRpcStreamTracing = ( @@ -54,10 +56,8 @@ const withRpcStreamTracing = ( traceAttributes?: Readonly>, ): Stream.Stream => shouldTraceRpc(method) - ? stream.pipe( - Stream.withSpan(`${RPC_SPAN_PREFIX}.${method}`, { - attributes: rpcSpanAttributes(method, traceAttributes), - }), + ? Stream.unwrap( + Effect.as(Effect.annotateCurrentSpan(rpcSpanAttributes(method, traceAttributes)), stream), ) : stream.pipe(Stream.provideService(References.TracerEnabled, false)); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index ae09cd38ac45..db2e322a525d 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -100,6 +100,7 @@ import { observeRpcEffect as instrumentRpcEffect, observeRpcStream as instrumentRpcStream, observeRpcStreamEffect as instrumentRpcStreamEffect, + rpcServerTracingOptions, } from "./observability/RpcInstrumentation.ts"; import * as ProviderRegistry from "./provider/Services/ProviderRegistry.ts"; import * as ModelManifest from "./provider/ModelManifest.ts"; @@ -3011,7 +3012,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(WsRpcGroup, { disableTracing: true }).pipe( + yield* RpcServer.make(WsRpcGroup, rpcServerTracingOptions).pipe( Effect.provideService(RpcServer.Protocol, withTerminalOutputWindow(protocol)), Effect.forkScoped, );