feat(llm): propagate trace context to HTTP requests
This commit is contained in:
parent
18da3deee7
commit
aacec8ea06
3 changed files with 41 additions and 31 deletions
|
|
@ -131,31 +131,32 @@ export const httpJson = <Body, Frame>(input: HttpJsonInput<Body, Frame>): HttpJs
|
||||||
frames: (prepared, request, runtime) =>
|
frames: (prepared, request, runtime) =>
|
||||||
LLMHttpTelemetry.stream(
|
LLMHttpTelemetry.stream(
|
||||||
prepared.request,
|
prepared.request,
|
||||||
Stream.unwrap(
|
(providerRequest) =>
|
||||||
runtime.http.execute(prepared.request).pipe(
|
Stream.unwrap(
|
||||||
Effect.map((response) =>
|
runtime.http.execute(providerRequest).pipe(
|
||||||
Stream.unwrap(
|
Effect.map((response) =>
|
||||||
Effect.gen(function* () {
|
Stream.unwrap(
|
||||||
const received = yield* LLMHttpTelemetry.ResponseReceived
|
Effect.gen(function* () {
|
||||||
if (received) yield* received(response.status)
|
const received = yield* LLMHttpTelemetry.ResponseReceived
|
||||||
const firstChunk = yield* LLMHttpTelemetry.ResponseChunkReceived
|
if (received) yield* received(response.status)
|
||||||
return prepared.framing.frame(
|
const firstChunk = yield* LLMHttpTelemetry.ResponseChunkReceived
|
||||||
response.stream.pipe(
|
return prepared.framing.frame(
|
||||||
Stream.tap(() => firstChunk ?? Effect.void),
|
response.stream.pipe(
|
||||||
Stream.mapError((error) =>
|
Stream.tap(() => firstChunk ?? Effect.void),
|
||||||
ProviderShared.eventError(
|
Stream.mapError((error) =>
|
||||||
`${request.model.provider}/${request.model.route.id}`,
|
ProviderShared.eventError(
|
||||||
`Failed to read ${request.model.provider}/${request.model.route.id} stream`,
|
`${request.model.provider}/${request.model.route.id}`,
|
||||||
ProviderShared.errorText(error),
|
`Failed to read ${request.model.provider}/${request.model.route.id} stream`,
|
||||||
|
ProviderShared.errorText(error),
|
||||||
|
),
|
||||||
),
|
),
|
||||||
),
|
),
|
||||||
),
|
)
|
||||||
)
|
}),
|
||||||
}),
|
),
|
||||||
),
|
),
|
||||||
),
|
),
|
||||||
),
|
),
|
||||||
),
|
|
||||||
),
|
),
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,7 @@ export * as LLMHttpTelemetry from "./http"
|
||||||
|
|
||||||
import { Cause, Clock, Context, Effect, Exit, Option, References, Stream } from "effect"
|
import { Cause, Clock, Context, Effect, Exit, Option, References, Stream } from "effect"
|
||||||
import { ParentSpan, type Span } from "effect/Tracer"
|
import { ParentSpan, type Span } from "effect/Tracer"
|
||||||
import { HttpClient, HttpClientRequest } from "effect/unstable/http"
|
import { HttpClient, HttpClientRequest, HttpTraceContext } from "effect/unstable/http"
|
||||||
import {
|
import {
|
||||||
ATTR_ERROR_TYPE,
|
ATTR_ERROR_TYPE,
|
||||||
ATTR_HTTP_REQUEST_METHOD,
|
ATTR_HTTP_REQUEST_METHOD,
|
||||||
|
|
@ -40,11 +40,14 @@ const observe = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
|
||||||
() => Effect.void,
|
() => Effect.void,
|
||||||
)
|
)
|
||||||
|
|
||||||
export const stream = <A, R>(request: HttpClientRequest.HttpClientRequest, source: Stream.Stream<A, LLMError, R>) =>
|
export const stream = <A, R>(
|
||||||
|
request: HttpClientRequest.HttpClientRequest,
|
||||||
|
source: (request: HttpClientRequest.HttpClientRequest) => Stream.Stream<A, LLMError, R>,
|
||||||
|
) =>
|
||||||
Stream.unwrap(
|
Stream.unwrap(
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const parent = yield* CurrentModelSpan
|
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 url = URL.canParse(request.url) ? new URL(request.url) : undefined
|
||||||
const port = url?.port
|
const port = url?.port
|
||||||
? Number(url.port)
|
? Number(url.port)
|
||||||
|
|
@ -71,12 +74,16 @@ export const stream = <A, R>(request: HttpClientRequest.HttpClientRequest, sourc
|
||||||
})
|
})
|
||||||
const state: State = { responseReceived: false }
|
const state: State = { responseReceived: false }
|
||||||
yield* Effect.addFinalizer((exit) => observe(finalize(span, state, exit)))
|
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(
|
}).pipe(
|
||||||
Effect.withTracerEnabled(true),
|
Effect.withTracerEnabled(true),
|
||||||
Effect.catchCauseIf(
|
Effect.catchCauseIf(
|
||||||
(cause) => !Cause.hasInterrupts(cause),
|
(cause) => !Cause.hasInterrupts(cause),
|
||||||
() => Effect.succeed(source),
|
() => Effect.succeed(source(request)),
|
||||||
),
|
),
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -184,7 +184,8 @@ describe("GenAI telemetry", () => {
|
||||||
? span.status.endTime >= http.status.endTime
|
? span.status.endTime >= http.status.endTime
|
||||||
: false,
|
: false,
|
||||||
).toBeTrue()
|
).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", () =>
|
yield* Effect.useSpan("ambient", () =>
|
||||||
Effect.all(
|
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,
|
Stream.runDrain,
|
||||||
),
|
),
|
||||||
LLMWebSocketTelemetry.stream("wss://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.acquireUseRelease(
|
||||||
Effect.sync(() => {
|
Effect.sync(() => {
|
||||||
const received: { traceparent?: string; b3?: string } = {}
|
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")
|
const http = spans.find((span) => span.attributes.get(ATTR_HTTP_REQUEST_METHOD) === "POST")
|
||||||
expect(http).toBeDefined()
|
const traceparent = received.traceparent?.split("-")
|
||||||
expect(received.traceparent).toBeUndefined()
|
expect(traceparent?.[1]).toBe(http?.traceId)
|
||||||
expect(received.b3).toBeUndefined()
|
expect(traceparent?.[2]).toBe(http?.spanId)
|
||||||
|
expect(received.b3).toStartWith(`${http?.traceId}-${http?.spanId}-`)
|
||||||
}),
|
}),
|
||||||
({ server }) => Effect.promise(() => server.stop(true)),
|
({ server }) => Effect.promise(() => server.stop(true)),
|
||||||
),
|
),
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue