diff --git a/packages/llm/src/route/transport/http.ts b/packages/llm/src/route/transport/http.ts index a70c4bf5382..6556174bcc6 100644 --- a/packages/llm/src/route/transport/http.ts +++ b/packages/llm/src/route/transport/http.ts @@ -131,31 +131,32 @@ export const httpJson =
(input: HttpJsonInput): HttpJs frames: (prepared, request, runtime) => LLMHttpTelemetry.stream( prepared.request, - Stream.unwrap( - runtime.http.execute(prepared.request).pipe( - Effect.map((response) => - Stream.unwrap( - Effect.gen(function* () { - const received = yield* LLMHttpTelemetry.ResponseReceived - if (received) yield* received(response.status) - const firstChunk = yield* LLMHttpTelemetry.ResponseChunkReceived - return prepared.framing.frame( - response.stream.pipe( - Stream.tap(() => firstChunk ?? Effect.void), - Stream.mapError((error) => - ProviderShared.eventError( - `${request.model.provider}/${request.model.route.id}`, - `Failed to read ${request.model.provider}/${request.model.route.id} stream`, - ProviderShared.errorText(error), + (providerRequest) => + Stream.unwrap( + runtime.http.execute(providerRequest).pipe( + Effect.map((response) => + Stream.unwrap( + Effect.gen(function* () { + const received = yield* LLMHttpTelemetry.ResponseReceived + if (received) yield* received(response.status) + const firstChunk = yield* LLMHttpTelemetry.ResponseChunkReceived + return prepared.framing.frame( + response.stream.pipe( + Stream.tap(() => firstChunk ?? Effect.void), + Stream.mapError((error) => + ProviderShared.eventError( + `${request.model.provider}/${request.model.route.id}`, + `Failed to read ${request.model.provider}/${request.model.route.id} stream`, + ProviderShared.errorText(error), + ), ), ), - ), - ) - }), + ) + }), + ), ), ), ), - ), ), }) diff --git a/packages/llm/src/telemetry/http.ts b/packages/llm/src/telemetry/http.ts index 2e9be90f3f0..19cecca2736 100644 --- a/packages/llm/src/telemetry/http.ts +++ b/packages/llm/src/telemetry/http.ts @@ -2,7 +2,7 @@ export * as LLMHttpTelemetry from "./http" import { Cause, Clock, Context, Effect, Exit, Option, References, Stream } from "effect" import { ParentSpan, type Span } from "effect/Tracer" -import { HttpClient, HttpClientRequest } from "effect/unstable/http" +import { HttpClient, HttpClientRequest, HttpTraceContext } from "effect/unstable/http" import { ATTR_ERROR_TYPE, ATTR_HTTP_REQUEST_METHOD, @@ -40,11 +40,14 @@ const observe = (effect: Effect.Effect) => () => Effect.void, ) -export const stream = (request: HttpClientRequest.HttpClientRequest, source: Stream.Stream) => +export const stream = ( + request: HttpClientRequest.HttpClientRequest, + source: (request: HttpClientRequest.HttpClientRequest) => Stream.Stream, +) => Stream.unwrap( Effect.gen(function* () { const parent = yield* CurrentModelSpan - if (!parent) return source + if (!parent) return source(request) const url = URL.canParse(request.url) ? new URL(request.url) : undefined const port = url?.port ? Number(url.port) @@ -71,12 +74,16 @@ export const stream = (request: HttpClientRequest.HttpClientRequest, sourc }) const state: State = { responseReceived: false } yield* Effect.addFinalizer((exit) => observe(finalize(span, state, exit))) - return observeStream(span, state, source) + return observeStream( + span, + state, + source(HttpClientRequest.setHeaders(request, HttpTraceContext.toHeaders(span))), + ) }).pipe( Effect.withTracerEnabled(true), Effect.catchCauseIf( (cause) => !Cause.hasInterrupts(cause), - () => Effect.succeed(source), + () => Effect.succeed(source(request)), ), ), ) diff --git a/packages/llm/test/telemetry.test.ts b/packages/llm/test/telemetry.test.ts index af0a5243a4c..eb7c4b5c584 100644 --- a/packages/llm/test/telemetry.test.ts +++ b/packages/llm/test/telemetry.test.ts @@ -184,7 +184,8 @@ describe("GenAI telemetry", () => { ? span.status.endTime >= http.status.endTime : false, ).toBeTrue() - expect(traceparent).toBeUndefined() + expect(traceparent?.split("-")[1]).toBe(http?.traceId) + expect(traceparent?.split("-")[2]).toBe(http?.spanId) }), ) @@ -303,7 +304,7 @@ describe("GenAI telemetry", () => { yield* Effect.useSpan("ambient", () => Effect.all( [ - LLMHttpTelemetry.stream(HttpClientRequest.post("https://example.test/path"), Stream.empty).pipe( + LLMHttpTelemetry.stream(HttpClientRequest.post("https://example.test/path"), () => Stream.empty).pipe( Stream.runDrain, ), LLMWebSocketTelemetry.stream("wss://example.test/path", Stream.empty).pipe(Stream.runDrain), @@ -524,7 +525,7 @@ describe("GenAI telemetry", () => { }), ) - it.live("does not mutate provider request headers", () => + it.live("propagates the transport span to provider requests", () => Effect.acquireUseRelease( Effect.sync(() => { const received: { traceparent?: string; b3?: string } = {} @@ -564,9 +565,10 @@ describe("GenAI telemetry", () => { ) const http = spans.find((span) => span.attributes.get(ATTR_HTTP_REQUEST_METHOD) === "POST") - expect(http).toBeDefined() - expect(received.traceparent).toBeUndefined() - expect(received.b3).toBeUndefined() + const traceparent = received.traceparent?.split("-") + expect(traceparent?.[1]).toBe(http?.traceId) + expect(traceparent?.[2]).toBe(http?.spanId) + expect(received.b3).toStartWith(`${http?.traceId}-${http?.spanId}-`) }), ({ server }) => Effect.promise(() => server.stop(true)), ),