From fcf1604453ae8850cda0734c4392b59cb42df33a Mon Sep 17 00:00:00 2001 From: Julius Marminge Date: Mon, 5 Oct 2026 23:05:01 -0700 Subject: [PATCH] fix(relay): export traces through one tracer, one request span each The relay Worker exported traces twice. Alongside Alchemy's Axiom.Telemetry it built its own OTLP tracer per request, and Alchemy's runtime opened its own http.server span around ours, so every request produced two server spans in Axiom: ours with the route and redaction, and an empty sibling. Axiom.Telemetry is now the only exporter. The runtime's request span is off for every path, leaving ours, and the schema error attributes our tracer added are applied to the current tracer instead. Co-Authored-By: Claude Opus 5.5 (1M context) --- infra/relay/src/hooks/HookForwarder.test.ts | 10 +-- infra/relay/src/http/Api.test.ts | 33 ++++++++ infra/relay/src/http/Api.ts | 10 ++- infra/relay/src/observability.test.ts | 84 ++++----------------- infra/relay/src/observability.ts | 38 +++------- infra/relay/src/worker.ts | 45 +++++------ 6 files changed, 87 insertions(+), 133 deletions(-) diff --git a/infra/relay/src/hooks/HookForwarder.test.ts b/infra/relay/src/hooks/HookForwarder.test.ts index dad1431a133f..f9df51a46974 100644 --- a/infra/relay/src/hooks/HookForwarder.test.ts +++ b/infra/relay/src/hooks/HookForwarder.test.ts @@ -615,15 +615,13 @@ describe("HookForwarder", () => { }); const harness = makeHarness(); const handler = yield* harness.httpEffect; - // As the worker runtime runs it: its own tracer around ours, off for hook - // paths. This checks the predicate; whether alchemy applies it per event - // is only visible on a deployed worker (see worker.ts). + // As the worker runtime runs it: its own tracer around ours, turned off + // (see worker.ts). Whether alchemy applies the predicate per event is + // only visible on a deployed worker. yield* HttpMiddleware.tracer( traceRelayHttpRequestWith(handler, Layer.succeed(Tracer.Tracer, tracer)), ).pipe( - Effect.provideService(HttpMiddleware.TracerDisabledWhen, (request) => - HookForwarder.isRelayHookPath(request.url), - ), + Effect.provideService(HttpMiddleware.TracerDisabledWhen, () => true), Effect.withTracer(tracer), Effect.provideService( HttpServerRequest.HttpServerRequest, diff --git a/infra/relay/src/http/Api.test.ts b/infra/relay/src/http/Api.test.ts index 4df0a27181ef..453d7d867210 100644 --- a/infra/relay/src/http/Api.test.ts +++ b/infra/relay/src/http/Api.test.ts @@ -30,6 +30,7 @@ import * as Tracer from "effect/Tracer"; import * as Etag from "effect/http/Etag"; import * as HttpEffect from "effect/http/HttpEffect"; import * as HttpRouter from "effect/http/HttpRouter"; +import * as HttpMiddleware from "effect/http/HttpMiddleware"; import * as HttpServerRequest from "effect/http/HttpServerRequest"; import * as HttpServerResponse from "effect/http/HttpServerResponse"; import * as HttpApi from "effect/http-api/HttpApi"; @@ -1195,6 +1196,38 @@ describe("relay request tracing", () => { }), ); + it.effect("records one server span inside the worker's disabled HTTP tracer", () => + Effect.gen(function* () { + const spans: Array = []; + const tracer = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + spans.push(span); + return span; + }, + }); + const request = HttpServerRequest.fromWeb( + new Request("https://relay.test/v1/mobile/devices", { method: "POST" }), + ); + + // As the worker runtime runs it: its own tracer around ours, turned off. + yield* HttpMiddleware.tracer( + traceRelayHttpRequestWith( + Effect.succeed(HttpServerResponse.empty({ status: 204 })), + Layer.empty, + ), + ).pipe( + Effect.provideService(HttpMiddleware.TracerDisabledWhen, () => true), + Effect.provideService(HttpServerRequest.HttpServerRequest, request), + Effect.withTracer(tracer), + ); + yield* Effect.yieldNow; + + expect(spans.filter((span) => span.kind === "server")).toHaveLength(1); + expect(spans[0]?.attributes.get("url.path")).toBe("/v1/mobile/devices"); + }), + ); + it.effect("fails hung requests with a 504 before the client's 10s abort", () => Effect.gen(function* () { const spans: Array = []; diff --git a/infra/relay/src/http/Api.ts b/infra/relay/src/http/Api.ts index 7e146022b8c5..69a48f8f6b81 100644 --- a/infra/relay/src/http/Api.ts +++ b/infra/relay/src/http/Api.ts @@ -228,8 +228,12 @@ export const traceRelayHttpRequest = ( Effect.andThen(relayRequestDeadline(httpEffect)), ); if (!isRelayHookPath(request.url)) { - // HttpMiddleware finalizes its span on the dispatcher; do not close a request-scoped exporter first. - return yield* HttpMiddleware.tracer(traced).pipe(Effect.ensuring(Effect.yieldNow)); + return yield* HttpMiddleware.tracer(traced).pipe( + // The worker turns its own request span off; this one is ours. + Effect.provideService(HttpMiddleware.TracerDisabledWhen, () => false), + // HttpMiddleware finalizes its span on the dispatcher; do not close a request-scoped exporter first. + Effect.ensuring(Effect.yieldNow), + ); } // Hook URLs carry a secret token: the tracer and deadline log see a redacted // request, while the route itself still receives the original. A webhook @@ -253,7 +257,7 @@ export const traceRelayHttpRequest = ( ), ).pipe( Effect.provideService(HttpServerRequest.HttpServerRequest, redacted), - // The worker disables its own span for hook paths; this one is ours. + // The worker turns its own request span off; this one is ours. Effect.provideService(HttpMiddleware.TracerDisabledWhen, () => false), Effect.ensuring(Effect.yieldNow), ); diff --git a/infra/relay/src/observability.test.ts b/infra/relay/src/observability.test.ts index 52f2cd1e7434..922f9fd7d439 100644 --- a/infra/relay/src/observability.test.ts +++ b/infra/relay/src/observability.test.ts @@ -1,46 +1,20 @@ -import * as NodeHttpServer from "@effect/platform-node/NodeHttpServer"; import { expect, it } from "@effect/vitest"; -import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; -import * as Redacted from "effect/Redacted"; -import * as Schema from "effect/Schema"; -import * as HttpServer from "effect/http/HttpServer"; -import * as HttpServerRequest from "effect/http/HttpServerRequest"; -import * as HttpServerResponse from "effect/http/HttpServerResponse"; -import type { OtlpTracer } from "effect/observability"; +import * as Tracer from "effect/Tracer"; import * as EnvironmentConnector from "./environments/EnvironmentConnector.ts"; import * as Observability from "./observability.ts"; -interface ExportedRequest { - readonly authorization: string | undefined; - readonly body: string; - readonly dataset: string | undefined; -} - -const otlpAttributeValue = (value: { - readonly stringValue?: string | null; - readonly boolValue?: boolean | null; - readonly intValue?: string | number | null; - readonly doubleValue?: number | null; -}) => value.stringValue ?? value.boolValue ?? value.intValue ?? value.doubleValue; - -const decodeJson = Schema.decodeUnknownEffect(Schema.fromJsonString(Schema.Unknown)); - -it.effect("exports schema error fields as span attributes", () => +it.effect("adds schema error fields to spans on the current tracer", () => Effect.gen(function* () { - const exportedRequest = yield* Deferred.make(); - yield* HttpServer.serveEffect( - Effect.gen(function* () { - const request = yield* HttpServerRequest.HttpServerRequest; - yield* Deferred.succeed(exportedRequest, { - authorization: request.headers.authorization, - body: yield* request.text, - dataset: request.headers["x-axiom-dataset"], - }); - return HttpServerResponse.empty({ status: 204 }); - }), - ); + const spans: Array = []; + const tracer = Tracer.make({ + span: (options) => { + const span = new Tracer.NativeSpan(options); + spans.push(span); + return span; + }, + }); yield* Effect.fail( new EnvironmentConnector.EnvironmentConnectNotAuthorized({ @@ -51,44 +25,16 @@ it.effect("exports schema error fields as span attributes", () => ).pipe( Effect.withSpan("relay.test.schema_error"), Effect.exit, - Effect.provide( - Observability.layer({ - tracesEndpoint: "/v1/traces", - tracesDatasetName: "relay-test-traces", - ingestToken: Redacted.make("test-token"), - }), - ), + Observability.withSchemaErrorSpanAttributes, + Effect.withTracer(tracer), ); - const request = yield* Deferred.await(exportedRequest).pipe(Effect.timeout("1 second")); - const payload = (yield* decodeJson(request.body)) as OtlpTracer.TraceData; - const resourceAttributes = Object.fromEntries( - payload.resourceSpans - .flatMap((resourceSpan) => resourceSpan.resource.attributes) - .map((attribute) => [attribute.key, otlpAttributeValue(attribute.value)]), - ); - const span = payload.resourceSpans - .flatMap((resourceSpan) => resourceSpan.scopeSpans) - .flatMap((scopeSpan) => scopeSpan.spans) - .find((candidate) => candidate.name === "relay.test.schema_error"); - const attributes = Object.fromEntries( - (span?.attributes ?? []).map((attribute) => [ - attribute.key, - otlpAttributeValue(attribute.value), - ]), - ); - - expect(request.authorization).toBe("Bearer test-token"); - expect(request.dataset).toBe("relay-test-traces"); - expect(resourceAttributes).toMatchObject({ - "service.name": "t3code-relay", - "service.namespace": "t3code", - }); - expect(attributes).toMatchObject({ + expect(spans.map((span) => span.name)).toEqual(["relay.test.schema_error"]); + expect(Object.fromEntries(spans[0]!.attributes)).toMatchObject({ "error.type": "EnvironmentConnectNotAuthorized", "error.environmentId": "environment-1", "error.operation": "connect", "error.reason": "managed_endpoint_allocation_not_ready", }); - }).pipe(Effect.provide(NodeHttpServer.layerTest), Effect.scoped), + }), ); diff --git a/infra/relay/src/observability.ts b/infra/relay/src/observability.ts index c0f377093a1f..cb469c9166d3 100644 --- a/infra/relay/src/observability.ts +++ b/infra/relay/src/observability.ts @@ -4,12 +4,9 @@ import * as Output from "alchemy/Output"; import * as Cause from "effect/Cause"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; -import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; -import * as Redacted from "effect/Redacted"; import * as Schema from "effect/Schema"; import * as Tracer from "effect/Tracer"; -import { OtlpExporter, OtlpSerialization, OtlpTracer } from "effect/observability"; import { relayResourceNameForStage } from "./deploymentConfig.ts"; @@ -212,27 +209,14 @@ const withSchemaErrorAttributes = (delegate: Tracer.Tracer): Tracer.Tracer => ...(delegate.context ? { context: delegate.context } : {}), }); -export const layer = (input: { - readonly tracesEndpoint: string; - readonly tracesDatasetName: string; - readonly ingestToken: Redacted.Redacted; -}) => - Layer.effect( - Tracer.Tracer, - OtlpTracer.make({ - url: input.tracesEndpoint, - resource: { - serviceName: "t3code-relay", - attributes: { - "service.namespace": "t3code", - "service.runtime": "cloudflare-worker", - "service.component": "relay", - }, - }, - headers: { - Authorization: `Bearer ${Redacted.value(input.ingestToken)}`, - "X-Axiom-Dataset": input.tracesDatasetName, - }, - exportInterval: "1 second", - }).pipe(Effect.map(withSchemaErrorAttributes)), - ).pipe(Layer.provideMerge(OtlpExporter.layerFlusher), Layer.provide(OtlpSerialization.layerJson)); +/** + * Adds a failed span's schema error fields (`error.type`, `error.`) to + * the span, on whichever tracer is current. Provide it around the work whose + * spans should carry them; the Worker's telemetry still owns export. + */ +export const withSchemaErrorSpanAttributes = ( + effect: Effect.Effect, +): Effect.Effect => + Effect.flatMap(Effect.tracer, (tracer) => + effect.pipe(Effect.withTracer(withSchemaErrorAttributes(tracer))), + ); diff --git a/infra/relay/src/worker.ts b/infra/relay/src/worker.ts index 71d0b4a0d787..adb77288057c 100644 --- a/infra/relay/src/worker.ts +++ b/infra/relay/src/worker.ts @@ -134,7 +134,7 @@ export const layer = Api.make( const relayApiZone = yield* RelayApiZone; const managedEndpointZone = yield* ManagedEndpointZone; const randomApnsDeliveryJobSigningSecret = yield* ApnsDeliveryJobSigningSecret; - const observability = yield* RelayObservability; + yield* RelayObservability; // // 2. Create bindings @@ -159,10 +159,6 @@ export const layer = Api.make( const apnsDeliveryQueueSender = yield* Cloudflare.Queues.WriteQueue(apnsDeliveryQueue); const fcmDeliveryQueueSender = yield* Cloudflare.Queues.WriteQueue(fcmDeliveryQueue); - const axiomDatasetName = yield* observability.traces.name; - const axiomIngestToken = yield* observability.workerIngestToken.token; - const axiomTracesEndpoint = yield* observability.traces.otelTracesEndpoint; - const clerkSecretKey = yield* Config.Redacted("CLERK_SECRET_KEY"); const clerkPublishableKey = yield* Config.String("CLERK_PUBLISHABLE_KEY"); const clerkJwtAudience = yield* Config.String("CLERK_JWT_AUDIENCE"); @@ -231,14 +227,6 @@ export const layer = Api.make( }); }); - const layerRelayTrace = Layer.unwrap( - Effect.all({ - tracesDatasetName: axiomDatasetName, - tracesEndpoint: axiomTracesEndpoint, - ingestToken: axiomIngestToken, - }).pipe(Effect.map(Observability.layer)), - ); - // Each managed endpoint's held webhook requests live in its own Durable Object. const inboxCall = (operation: HookInbox.HookInboxError["operation"], endpointKey: string) => @@ -442,8 +430,8 @@ export const layer = Api.make( { concurrency: 2, discard: true }, ).pipe( Effect.withSpan("relay.cron.prune_expired_state"), - // Export cron spans to Axiom like HTTP spans; the scope flushes them before the run ends. - Effect.provide(Layer.merge(layerRuntime, layerRelayTrace)), + Observability.withSchemaErrorSpanAttributes, + Effect.provide(layerRuntime), ), ); @@ -462,7 +450,11 @@ export const layer = Api.make( HttpRouter.toHttpEffect, Effect.provideService(HttpRouter.RouterConfig, RELAY_HTTP_ROUTER_CONFIG), withoutCapturedParentSpan, - Effect.flatMap((httpEffect) => traceRelayHttpRequestWith(httpEffect, layerRelayTrace)), + Effect.map((httpEffect) => + traceRelayHttpRequestWith(httpEffect, Layer.empty).pipe( + Observability.withSchemaErrorSpanAttributes, + ), + ), ); return { fetch }; @@ -477,20 +469,17 @@ export const layer = Api.make( Layer.provideMerge(Cloudflare.DNS.ReadWriteDnsHttp), Layer.provideMerge(Cloudflare.Workers.RateLimitBinding), Layer.provideMerge(HookInboxObject.layer), - // The worker runtime opens its own HTTP span around ours. For webhook - // paths it would record the raw URL, token included, and adopt the - // sender's traceparent, so only our redacted span covers those. - // Registered as telemetry: request-time context is assembled per - // event, and only these layers are built into it. + // The worker runtime opens its own HTTP span around ours. Ours carries + // the route, header redaction and webhook URL redaction, and drops a + // webhook sender's traceparent, so the runtime's span is off for every + // request rather than duplicating it. Registered as telemetry: + // request-time context is assembled per event, and only these layers + // are built into it. Layer.provideMerge( - Alchemy.Telemetry.layer( - Layer.succeed(HttpMiddleware.TracerDisabledWhen)((request) => - HookForwarder.isRelayHookPath(request.url), - ), - ), + Alchemy.Telemetry.layer(Layer.succeed(HttpMiddleware.TracerDisabledWhen)(() => true)), ), - // Exports spans from events the HTTP tracer does not wrap, notably - // HookInboxObject calls and alarms, to the same Axiom dataset. + // The Worker's only trace exporter: every event (fetch, queue, cron, + // HookInboxObject calls and alarms) exports through it to Axiom. Layer.provideMerge( Layer.unwrap( Effect.map(RelayObservability, (observability) =>