Compare commits

...

1 Commits

Author SHA1 Message Date
Aiden Cline ffb35ee936 fix(ai): preserve stream transport failures 2026-08-11 16:01:03 +00:00
4 changed files with 45 additions and 29 deletions
+27 -14
View File
@@ -297,7 +297,7 @@ export const classifyHttpFailure = (input: {
})
}
const toHttpError = (redactedNames: ReadonlyArray<string | RegExp>) => (error: unknown) => {
export const mapHttpError = (error: unknown, redactedNames: ReadonlyArray<string | RegExp>) => {
const transportError = (input: {
readonly message: string
readonly kind?: string | undefined
@@ -314,23 +314,35 @@ const toHttpError = (redactedNames: ReadonlyArray<string | RegExp>) => (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<Service, never, HttpClient.HttpClient> = 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({
+6 -9
View File
@@ -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 = <Body, Frame>(input: HttpJsonInput<Body, Frame>): 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))),
),
),
),
+2 -2
View File
@@ -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<string>) =>
export const truncatedStream = (chunks: ReadonlyArray<string>, 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 })
+10 -4
View File
@@ -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",
})
}),
)