http.ts 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  1. import { Effect, Layer, Ref } from "effect"
  2. import { HttpClient, HttpClientError, HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
  3. import { LLMClient, RequestExecutor } from "../../src/route.js"
  4. import type { Service as LLMClientService } from "../../src/route/client.js"
  5. import type { Service as RequestExecutorService } from "../../src/route/executor.js"
  6. export type HandlerInput = {
  7. readonly request: HttpClientRequest.HttpClientRequest
  8. readonly text: string
  9. readonly respond: (
  10. body: ConstructorParameters<typeof Response>[0],
  11. init?: ResponseInit,
  12. ) => HttpClientResponse.HttpClientResponse
  13. }
  14. export type Handler = (
  15. input: HandlerInput,
  16. ) => Effect.Effect<HttpClientResponse.HttpClientResponse, HttpClientError.HttpClientError>
  17. const handlerLayer = (handler: Handler): Layer.Layer<HttpClient.HttpClient> =>
  18. Layer.succeed(
  19. HttpClient.HttpClient,
  20. HttpClient.make((request) =>
  21. Effect.gen(function* () {
  22. const web = yield* HttpClientRequest.toWeb(request).pipe(Effect.orDie)
  23. const text = yield* Effect.promise(() => web.text())
  24. return yield* handler({
  25. request,
  26. text,
  27. respond: (body, init) => HttpClientResponse.fromWeb(request, new Response(body, init)),
  28. })
  29. }),
  30. ),
  31. )
  32. export type RuntimeEnv = RequestExecutorService | LLMClientService
  33. export interface SystemError extends Error {
  34. readonly code: string
  35. }
  36. export const systemError = (code: string, message: string): SystemError => Object.assign(new Error(message), { code })
  37. export const runtimeLayer = (layer: Layer.Layer<HttpClient.HttpClient>): Layer.Layer<RuntimeEnv> => {
  38. const requestExecutorLayer = RequestExecutor.layer.pipe(Layer.provide(layer))
  39. const llmClientLayer = LLMClient.layer.pipe(Layer.provide(requestExecutorLayer))
  40. return Layer.mergeAll(requestExecutorLayer, llmClientLayer)
  41. }
  42. const SSE_HEADERS = { "content-type": "text/event-stream" } as const
  43. /**
  44. * Layer that returns a single fixed response body. Use for stream-parser
  45. * fixture tests where the request shape is irrelevant. The body type widens
  46. * to whatever `Response` accepts so binary fixtures (`Uint8Array`,
  47. * `ReadableStream`, etc.) flow through without casts.
  48. */
  49. export const fixedResponse = (
  50. body: ConstructorParameters<typeof Response>[0],
  51. init: ResponseInit = { headers: SSE_HEADERS },
  52. ) => runtimeLayer(handlerLayer((input) => Effect.succeed(input.respond(body, init))))
  53. /**
  54. * Layer that builds a response per request. Useful for echo servers.
  55. */
  56. export const dynamicResponse = (handler: Handler) => runtimeLayer(handlerLayer(handler))
  57. /**
  58. * Layer that emits the supplied SSE chunks and then aborts mid-stream. Used to
  59. * exercise transport errors that surface during parsing.
  60. */
  61. export const truncatedStream = (chunks: ReadonlyArray<string>, error: Error = new Error("connection reset")) =>
  62. dynamicResponse((input) =>
  63. Effect.sync(() => {
  64. const encoder = new TextEncoder()
  65. let index = 0
  66. const stream = new ReadableStream({
  67. pull(controller) {
  68. const chunk = chunks[index]
  69. if (chunk !== undefined) {
  70. index++
  71. controller.enqueue(encoder.encode(chunk))
  72. return
  73. }
  74. controller.error(error)
  75. },
  76. })
  77. return input.respond(stream, { headers: SSE_HEADERS })
  78. }),
  79. )
  80. /**
  81. * Layer that returns successive bodies on each request. Useful for scripting
  82. * multi-step model exchanges (e.g. tool-call loops). The last body in the
  83. * array is reused if the test makes more requests than scripted.
  84. */
  85. export const scriptedResponses = (bodies: ReadonlyArray<string>, init: ResponseInit = { headers: SSE_HEADERS }) => {
  86. if (bodies.length === 0) throw new Error("scriptedResponses requires at least one body")
  87. return Layer.unwrap(
  88. Effect.gen(function* () {
  89. const cursor = yield* Ref.make(0)
  90. return dynamicResponse((input) =>
  91. Effect.gen(function* () {
  92. const index = yield* Ref.getAndUpdate(cursor, (n) => n + 1)
  93. return input.respond(bodies[index] ?? bodies[bodies.length - 1], init)
  94. }),
  95. )
  96. }),
  97. )
  98. }