229 lines
7.9 KiB
TypeScript
229 lines
7.9 KiB
TypeScript
import { describe, expect } from "bun:test"
|
|
import { Effect, Fiber, Layer, Random, Ref } from "effect"
|
|
import * as TestClock from "effect/testing/TestClock"
|
|
import { Headers, HttpClient, HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
|
|
import { LLM, ProviderChunkError, ProviderRequestError } from "../src"
|
|
import { LLMClient, RequestExecutor } from "../src/adapter"
|
|
import * as OpenAIChat from "../src/protocols/openai-chat"
|
|
import { dynamicResponse } from "./lib/http"
|
|
import { deltaChunk } from "./lib/openai-chunks"
|
|
import { sseRaw } from "./lib/sse"
|
|
import { it } from "./lib/effect"
|
|
|
|
const request = HttpClientRequest.post("https://provider.test/v1/chat?api_key=secret&key=secret&debug=1").pipe(
|
|
HttpClientRequest.setHeaders(Headers.fromInput({ authorization: "Bearer secret", "x-safe": "visible" })),
|
|
)
|
|
|
|
const responsesLayer = (responses: ReadonlyArray<Response>) =>
|
|
RequestExecutor.layer.pipe(
|
|
Layer.provide(
|
|
Layer.unwrap(
|
|
Effect.gen(function* () {
|
|
const cursor = yield* Ref.make(0)
|
|
return Layer.succeed(
|
|
HttpClient.HttpClient,
|
|
HttpClient.make((request) =>
|
|
Effect.gen(function* () {
|
|
const index = yield* Ref.getAndUpdate(cursor, (value) => value + 1)
|
|
return HttpClientResponse.fromWeb(request, responses[index] ?? responses[responses.length - 1])
|
|
}),
|
|
),
|
|
)
|
|
}),
|
|
),
|
|
),
|
|
)
|
|
|
|
const countedResponsesLayer = (attempts: Ref.Ref<number>, responses: ReadonlyArray<Response>) =>
|
|
RequestExecutor.layer.pipe(
|
|
Layer.provide(
|
|
Layer.unwrap(
|
|
Effect.gen(function* () {
|
|
const cursor = yield* Ref.make(0)
|
|
return Layer.succeed(
|
|
HttpClient.HttpClient,
|
|
HttpClient.make((request) =>
|
|
Effect.gen(function* () {
|
|
yield* Ref.update(attempts, (value) => value + 1)
|
|
const index = yield* Ref.getAndUpdate(cursor, (value) => value + 1)
|
|
return HttpClientResponse.fromWeb(request, responses[index] ?? responses[responses.length - 1])
|
|
}),
|
|
),
|
|
)
|
|
}),
|
|
),
|
|
),
|
|
)
|
|
|
|
const randomMidpoint = {
|
|
nextDoubleUnsafe: () => 0.5,
|
|
nextIntUnsafe: () => 0,
|
|
}
|
|
|
|
describe("RequestExecutor", () => {
|
|
it.effect("returns redacted diagnostics for retryable rate limits", () =>
|
|
Effect.gen(function* () {
|
|
const executor = yield* RequestExecutor.Service
|
|
const error = yield* executor.execute(request).pipe(Effect.flip)
|
|
|
|
expect(error).toBeInstanceOf(ProviderRequestError)
|
|
if (!(error instanceof ProviderRequestError)) throw new Error("expected ProviderRequestError")
|
|
expect(error).toMatchObject({
|
|
status: 429,
|
|
retryable: true,
|
|
retryAfterMs: 0,
|
|
requestId: "req_123",
|
|
request: {
|
|
method: "POST",
|
|
url: "https://provider.test/v1/chat?api_key=%3Credacted%3E&key=%3Credacted%3E&debug=1",
|
|
headers: { authorization: "<redacted>", "x-safe": "visible" },
|
|
},
|
|
response: {
|
|
status: 429,
|
|
headers: {
|
|
"retry-after-ms": "0",
|
|
"x-request-id": "req_123",
|
|
"x-api-key": "<redacted>",
|
|
},
|
|
},
|
|
})
|
|
expect(error.body).toBe("rate limited")
|
|
}).pipe(
|
|
Effect.provide(
|
|
responsesLayer([
|
|
...Array.from({ length: 3 }, () => new Response("rate limited", {
|
|
status: 429,
|
|
headers: { "retry-after-ms": "0", "x-request-id": "req_123", "x-api-key": "secret" },
|
|
})),
|
|
]),
|
|
),
|
|
),
|
|
)
|
|
|
|
it.effect("retries retryable status responses before returning the stream", () =>
|
|
Effect.gen(function* () {
|
|
const executor = yield* RequestExecutor.Service
|
|
const response = yield* executor.execute(request)
|
|
|
|
expect(response.status).toBe(200)
|
|
expect(yield* response.text).toBe("ok")
|
|
}).pipe(
|
|
Effect.provide(
|
|
responsesLayer([
|
|
new Response("busy", { status: 503, headers: { "retry-after-ms": "0" } }),
|
|
new Response("ok", { status: 200 }),
|
|
]),
|
|
),
|
|
),
|
|
)
|
|
|
|
it.effect("does not retry non-retryable status responses and truncates large bodies", () =>
|
|
Effect.gen(function* () {
|
|
const executor = yield* RequestExecutor.Service
|
|
const error = yield* executor.execute(request).pipe(Effect.flip)
|
|
|
|
expect(error).toBeInstanceOf(ProviderRequestError)
|
|
if (!(error instanceof ProviderRequestError)) throw new Error("expected ProviderRequestError")
|
|
expect(error.retryable).toBe(false)
|
|
expect(error.bodyTruncated).toBe(true)
|
|
expect(error.body).toHaveLength(16_384)
|
|
}).pipe(
|
|
Effect.provide(
|
|
responsesLayer([
|
|
new Response("x".repeat(20_000), { status: 401 }),
|
|
new Response("should not retry", { status: 200 }),
|
|
]),
|
|
),
|
|
),
|
|
)
|
|
|
|
it.effect("redacts common secret fields in response bodies", () =>
|
|
Effect.gen(function* () {
|
|
const executor = yield* RequestExecutor.Service
|
|
const error = yield* executor.execute(request).pipe(Effect.flip)
|
|
|
|
expect(error).toBeInstanceOf(ProviderRequestError)
|
|
if (!(error instanceof ProviderRequestError)) throw new Error("expected ProviderRequestError")
|
|
expect(error.body).toContain('"key":"<redacted>"')
|
|
expect(error.body).toContain('api_key=<redacted>')
|
|
expect(error.body).not.toContain("body-secret")
|
|
expect(error.body).not.toContain("query-secret")
|
|
}).pipe(
|
|
Effect.provide(
|
|
responsesLayer([
|
|
new Response('{"error":{"message":"bad","key":"body-secret","detail":"api_key=query-secret"}}', {
|
|
status: 400,
|
|
}),
|
|
]),
|
|
),
|
|
),
|
|
)
|
|
|
|
it.effect("uses exponential jittered delay when retry-after is absent", () =>
|
|
Effect.gen(function* () {
|
|
const attempts = yield* Ref.make(0)
|
|
return yield* Effect.gen(function* () {
|
|
const executor = yield* RequestExecutor.Service
|
|
const fiber = yield* executor.execute(request).pipe(Effect.flip, Effect.forkChild)
|
|
|
|
yield* Effect.yieldNow
|
|
expect(yield* Ref.get(attempts)).toBe(1)
|
|
|
|
yield* TestClock.adjust(499)
|
|
yield* Effect.yieldNow
|
|
expect(yield* Ref.get(attempts)).toBe(1)
|
|
|
|
yield* TestClock.adjust(1)
|
|
yield* Effect.yieldNow
|
|
expect(yield* Ref.get(attempts)).toBe(2)
|
|
|
|
yield* TestClock.adjust(999)
|
|
yield* Effect.yieldNow
|
|
expect(yield* Ref.get(attempts)).toBe(2)
|
|
|
|
yield* TestClock.adjust(1)
|
|
const error = yield* Fiber.join(fiber)
|
|
|
|
expect(error).toBeInstanceOf(ProviderRequestError)
|
|
expect(yield* Ref.get(attempts)).toBe(3)
|
|
}).pipe(
|
|
Effect.provide(
|
|
countedResponsesLayer(attempts, [
|
|
new Response("busy", { status: 503 }),
|
|
new Response("still busy", { status: 503 }),
|
|
new Response("done retrying", { status: 503 }),
|
|
]),
|
|
),
|
|
)
|
|
}).pipe(Effect.provideService(Random.Random, randomMidpoint)),
|
|
)
|
|
|
|
it.effect("does not retry after a successful response reaches stream parsing", () =>
|
|
Effect.gen(function* () {
|
|
const attempts = yield* Ref.make(0)
|
|
const model = OpenAIChat.model({ id: "gpt-4o-mini", baseURL: "https://api.openai.test/v1" })
|
|
const error = yield* LLMClient.generate(LLM.request({ model, prompt: "Say hello." })).pipe(
|
|
Effect.provide(
|
|
dynamicResponse((input) =>
|
|
Ref.update(attempts, (value) => value + 1).pipe(
|
|
Effect.as(
|
|
input.respond(
|
|
sseRaw(
|
|
`data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}`,
|
|
"data: not-json",
|
|
),
|
|
{ headers: { "content-type": "text/event-stream" } },
|
|
),
|
|
),
|
|
),
|
|
),
|
|
),
|
|
Effect.flip,
|
|
)
|
|
|
|
expect(error).toBeInstanceOf(ProviderChunkError)
|
|
expect(yield* Ref.get(attempts)).toBe(1)
|
|
}),
|
|
)
|
|
})
|