diff --git a/packages/ai/src/route/executor.ts b/packages/ai/src/route/executor.ts index ac30d230971..dd768a28460 100644 --- a/packages/ai/src/route/executor.ts +++ b/packages/ai/src/route/executor.ts @@ -297,7 +297,7 @@ export const classifyHttpFailure = (input: { }) } -const toHttpError = (redactedNames: ReadonlyArray) => (error: unknown) => { +export const mapHttpError = (error: unknown, redactedNames: ReadonlyArray) => { const transportError = (input: { readonly message: string readonly kind?: string | undefined @@ -314,23 +314,35 @@ const toHttpError = (redactedNames: ReadonlyArray) => (error: u }), }) - if (Cause.isTimeoutError(error)) { - return transportError({ message: error.message, kind: "Timeout" }) - } + const cause = + HttpClientError.isHttpClientError(error) && "cause" in error.reason + ? error.reason.cause + : error instanceof Error + ? error.cause + : undefined + const code = [cause, error] + .map((value) => (typeof value === "object" && value !== null ? Reflect.get(value, "code") : undefined)) + .find((value): value is string => typeof value === "string") + const request = HttpClientError.isHttpClientError(error) && "request" in error ? error.request : undefined + const raw = cause instanceof Error ? cause.message : error instanceof Error ? error.message : undefined + const detail = raw && request ? redactBody(raw, secretValues(request)) : raw + const message = code && detail && !detail.includes(code) ? `${code}: ${detail}` : detail + + if (Cause.isTimeoutError(error) || Cause.isTimeoutError(cause)) + return transportError({ message: message ?? "HTTP transport timed out", kind: code ?? "Timeout", request }) if (!HttpClientError.isHttpClientError(error)) { - return transportError({ message: error instanceof Error ? error.message : "HTTP transport failed" }) + return transportError({ message: message ?? "HTTP transport failed", kind: code, request }) } - const request = "request" in error ? error.request : undefined if (error.reason._tag === "TransportError") { return transportError({ - message: error.reason.description ?? "HTTP transport failed", - kind: error.reason._tag, + message: message ?? error.reason.description ?? "HTTP transport failed", + kind: code ?? error.reason._tag, request, }) } return transportError({ - message: `HTTP transport failed: ${error.reason._tag}`, - kind: error.reason._tag, + message: message ?? `HTTP transport failed: ${error.reason._tag}`, + kind: code ?? error.reason._tag, request, }) } @@ -343,15 +355,16 @@ export const layer: Layer.Layer = Layer.e Effect.gen(function* () { const redactedNames = yield* Headers.CurrentRedactedNames if (!middleware) - return yield* http - .execute(request) - .pipe(Effect.mapError(toHttpError(redactedNames)), Effect.flatMap(statusError(request, redactedNames))) + return yield* http.execute(request).pipe( + Effect.mapError((error) => mapHttpError(error, redactedNames)), + Effect.flatMap(statusError(request, redactedNames)), + ) const response = yield* middleware(request, (input) => http .execute(input) .pipe(Effect.mapError((cause) => (cause instanceof Error ? cause : new Error(String(cause))))), - ).pipe(Effect.mapError(toHttpError(redactedNames))) + ).pipe(Effect.mapError((error) => mapHttpError(error, redactedNames))) return yield* statusError(response.request, redactedNames)(response) }) return Service.of({ diff --git a/packages/ai/src/route/transport/http.ts b/packages/ai/src/route/transport/http.ts index 223afb353c1..add75a7c798 100644 --- a/packages/ai/src/route/transport/http.ts +++ b/packages/ai/src/route/transport/http.ts @@ -2,6 +2,7 @@ import { Effect, Stream } from "effect" import { Headers, HttpClientRequest } from "effect/unstable/http" import { Auth } from "../auth" import { render as renderEndpoint } from "../endpoint" +import { mapHttpError } from "../executor" import { Framing } from "../framing" import type { HttpMiddleware, Transport, TransportPrepareInput } from "./index" import * as ProviderShared from "../../protocols/shared" @@ -86,20 +87,16 @@ export const httpJson = (input: HttpJsonInput): HttpJs middleware: prepareInput.middleware, } }), - frames: (prepared, request, runtime) => + frames: (prepared, _request, runtime) => Stream.unwrap( runtime.http .execute(prepared.request, prepared.middleware) .pipe( Effect.map((response) => - prepared.framing.frame( - response.stream.pipe( - 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), - ), + Stream.unwrap( + Effect.map(Headers.CurrentRedactedNames, (redactedNames) => + prepared.framing.frame( + response.stream.pipe(Stream.mapError((error) => mapHttpError(error, redactedNames))), ), ), ), diff --git a/packages/ai/test/lib/http.ts b/packages/ai/test/lib/http.ts index f6c600555b9..fa69670d7ff 100644 --- a/packages/ai/test/lib/http.ts +++ b/packages/ai/test/lib/http.ts @@ -63,14 +63,14 @@ export const dynamicResponse = (handler: Handler) => runtimeLayer(handlerLayer(h * Layer that emits the supplied SSE chunks and then aborts mid-stream. Used to * exercise transport errors that surface during parsing. */ -export const truncatedStream = (chunks: ReadonlyArray) => +export const truncatedStream = (chunks: ReadonlyArray, error: Error = new Error("connection reset")) => dynamicResponse((input) => Effect.sync(() => { const encoder = new TextEncoder() const stream = new ReadableStream({ start(controller) { for (const chunk of chunks) controller.enqueue(encoder.encode(chunk)) - controller.error(new Error("connection reset")) + controller.error(error) }, }) return input.respond(stream, { headers: SSE_HEADERS }) diff --git a/packages/ai/test/provider/openai-chat.test.ts b/packages/ai/test/provider/openai-chat.test.ts index 6a0f7740ef6..abc3eaea83d 100644 --- a/packages/ai/test/provider/openai-chat.test.ts +++ b/packages/ai/test/provider/openai-chat.test.ts @@ -1221,12 +1221,18 @@ describe("OpenAI Chat route", () => { it.effect("surfaces transport errors that occur mid-stream", () => Effect.gen(function* () { - const layer = truncatedStream([ - `data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}\n\n`, - ]) + const layer = truncatedStream( + [`data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}\n\n`], + Object.assign(new Error("socket closed unexpectedly"), { code: "ECONNRESET" }), + ) const error = yield* LLMClient.generate(request).pipe(Effect.provide(layer), Effect.flip) - expect(error.message).toContain("Failed to read openai/openai-chat stream") + expect(error.reason).toMatchObject({ + _tag: "Transport", + message: "ECONNRESET: socket closed unexpectedly", + kind: "ECONNRESET", + url: "https://api.openai.test/v1/chat/completions", + }) }), )