1
0

executor.test.ts 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692
  1. import { describe, expect } from "bun:test"
  2. import { Deferred, Effect, Fiber, Layer, Ref, Stream } from "effect"
  3. import { Headers, HttpClient, HttpClientError, HttpClientRequest, HttpClientResponse } from "effect/unstable/http"
  4. import { LLM, AIError } from "../src/index.js"
  5. import { LLMClient, RequestExecutor, WebSocketTransport, type WebSocketChannelExecutor } from "../src/route.js"
  6. import * as OpenAIChat from "../src/protocols/openai-chat.js"
  7. import * as OpenAI from "../src/providers/openai.js"
  8. import { dynamicResponse, fixedResponse, systemError } from "./lib/http.js"
  9. import { deltaChunk } from "./lib/openai-chunks.js"
  10. import { sseEvents, sseRaw } from "./lib/sse.js"
  11. import { it } from "./lib/effect.js"
  12. const request = HttpClientRequest.post("https://provider.test/v1/chat?api_key=secret&key=secret&debug=1").pipe(
  13. HttpClientRequest.setHeaders(Headers.fromInput({ authorization: "Bearer secret", "x-safe": "visible" })),
  14. )
  15. const secretRequest = HttpClientRequest.post("https://provider.test/v1/chat?api_key=query-secret-123&debug=1").pipe(
  16. HttpClientRequest.setHeaders(Headers.fromInput({ authorization: "Bearer header-secret-456" })),
  17. )
  18. const responsesLayer = (responses: ReadonlyArray<Response>) =>
  19. RequestExecutor.layer.pipe(
  20. Layer.provide(
  21. Layer.unwrap(
  22. Effect.gen(function* () {
  23. const cursor = yield* Ref.make(0)
  24. return Layer.succeed(
  25. HttpClient.HttpClient,
  26. HttpClient.make((request) =>
  27. Effect.gen(function* () {
  28. const index = yield* Ref.getAndUpdate(cursor, (value) => value + 1)
  29. return HttpClientResponse.fromWeb(request, responses[index] ?? responses[responses.length - 1])
  30. }),
  31. ),
  32. )
  33. }),
  34. ),
  35. ),
  36. )
  37. const countedResponsesLayer = (attempts: Ref.Ref<number>, responses: ReadonlyArray<Response>) =>
  38. RequestExecutor.layer.pipe(
  39. Layer.provide(
  40. Layer.unwrap(
  41. Effect.gen(function* () {
  42. const cursor = yield* Ref.make(0)
  43. return Layer.succeed(
  44. HttpClient.HttpClient,
  45. HttpClient.make((request) =>
  46. Effect.gen(function* () {
  47. yield* Ref.update(attempts, (value) => value + 1)
  48. const index = yield* Ref.getAndUpdate(cursor, (value) => value + 1)
  49. return HttpClientResponse.fromWeb(request, responses[index] ?? responses[responses.length - 1])
  50. }),
  51. ),
  52. )
  53. }),
  54. ),
  55. ),
  56. )
  57. const expectAIError = (error: unknown) => {
  58. expect(error).toBeInstanceOf(AIError)
  59. if (!(error instanceof AIError)) throw new Error("expected AIError")
  60. return error
  61. }
  62. const errorHttp = (error: AIError) => ("http" in error.reason ? error.reason.http : undefined)
  63. const largeProviderMessage = `Upstream request failed: ${"validation failed; ".repeat(1_000)}`
  64. describe("RequestExecutor", () => {
  65. it.effect("parses response body failures at the executor seam", () =>
  66. Effect.gen(function* () {
  67. const executor = yield* RequestExecutor.Service
  68. const error = yield* RequestExecutor.stream(executor, secretRequest).pipe(Stream.runDrain, Effect.flip)
  69. expectAIError(error)
  70. expect(error.reason).toMatchObject({
  71. _tag: "Transport",
  72. message: "ECONNRESET: disconnected query-secret-123 header-secret-456",
  73. transport: "http",
  74. operation: "read",
  75. code: "ECONNRESET",
  76. url: "https://provider.test/v1/chat?api_key=query-secret-123&debug=1",
  77. })
  78. }).pipe(
  79. Effect.provide(
  80. responsesLayer([
  81. new Response(
  82. new ReadableStream({
  83. start(controller) {
  84. controller.error(systemError("ECONNRESET", "disconnected query-secret-123 header-secret-456"))
  85. },
  86. }),
  87. ),
  88. ]),
  89. ),
  90. ),
  91. )
  92. it.effect("unwraps native transport failure causes", () =>
  93. Effect.gen(function* () {
  94. const executor = yield* RequestExecutor.Service
  95. const error = yield* RequestExecutor.stream(executor, secretRequest).pipe(Stream.runDrain, Effect.flip)
  96. expectAIError(error)
  97. expect(error.reason).toMatchObject({
  98. _tag: "Transport",
  99. message: "ECONNRESET: socket closed",
  100. operation: "read",
  101. code: "ECONNRESET",
  102. })
  103. }).pipe(
  104. Effect.provide(
  105. responsesLayer([
  106. new Response(
  107. new ReadableStream({
  108. pull(controller) {
  109. controller.error(new TypeError("fetch failed", { cause: systemError("ECONNRESET", "socket closed") }))
  110. },
  111. }),
  112. ),
  113. ]),
  114. ),
  115. ),
  116. )
  117. it.effect("preserves middleware error messages", () =>
  118. Effect.gen(function* () {
  119. const executor = yield* RequestExecutor.Service
  120. const error = yield* executor
  121. .execute(request, () => Effect.fail(new Error("plugin rejected request")))
  122. .pipe(Effect.flip)
  123. expectAIError(error)
  124. expect(error.reason.message).toBe("plugin rejected request")
  125. }).pipe(Effect.provide(responsesLayer([]))),
  126. )
  127. it.effect("reports the request sent by middleware", () =>
  128. Effect.gen(function* () {
  129. const executor = yield* RequestExecutor.Service
  130. const error = yield* executor
  131. .execute(request, (original, handler) =>
  132. handler(
  133. original.pipe(
  134. HttpClientRequest.setUrl("https://proxy.test/v1/chat?api_key=proxy-secret"),
  135. HttpClientRequest.setHeader("authorization", "Bearer proxy-secret"),
  136. ),
  137. ),
  138. )
  139. .pipe(Effect.flip)
  140. expectAIError(error)
  141. expect(error.reason).toMatchObject({
  142. _tag: "Transport",
  143. message: "ECONNRESET: proxy disconnected proxy-secret",
  144. url: "https://proxy.test/v1/chat?api_key=proxy-secret",
  145. http: {
  146. request: {
  147. url: "https://proxy.test/v1/chat?api_key=proxy-secret",
  148. headers: { authorization: "Bearer proxy-secret" },
  149. },
  150. },
  151. })
  152. }).pipe(
  153. Effect.provide(
  154. dynamicResponse((input) =>
  155. Effect.fail(
  156. new HttpClientError.HttpClientError({
  157. reason: new HttpClientError.TransportError({
  158. request: input.request,
  159. cause: systemError("ECONNRESET", "proxy disconnected proxy-secret"),
  160. }),
  161. }),
  162. ),
  163. ),
  164. ),
  165. ),
  166. )
  167. it.effect("classifies context overflow responses", () =>
  168. Effect.gen(function* () {
  169. const executor = yield* RequestExecutor.Service
  170. const error = yield* executor.execute(request).pipe(Effect.flip)
  171. expectAIError(error)
  172. expect(error.reason).toMatchObject({ _tag: "InvalidRequest", classification: "context-overflow" })
  173. }).pipe(
  174. Effect.provide(
  175. responsesLayer([
  176. new Response('{"error":{"code":"context_length_exceeded","message":"prompt too long"}}', {
  177. status: 400,
  178. }),
  179. ]),
  180. ),
  181. ),
  182. )
  183. it.effect("classifies generic HTTP 413 payload errors", () =>
  184. Effect.gen(function* () {
  185. const executor = yield* RequestExecutor.Service
  186. const error = yield* executor.execute(request).pipe(Effect.flip)
  187. expectAIError(error)
  188. expect(error.reason).toMatchObject({
  189. _tag: "InvalidRequest",
  190. classification: "payload-too-large",
  191. http: { response: { status: 413 } },
  192. })
  193. }).pipe(Effect.provide(responsesLayer([new Response("request too large", { status: 413 })]))),
  194. )
  195. it.effect("does not classify ordinary invalid requests as context overflow", () =>
  196. Effect.gen(function* () {
  197. const executor = yield* RequestExecutor.Service
  198. const error = yield* executor.execute(request).pipe(Effect.flip)
  199. expectAIError(error)
  200. expect(error.reason).toMatchObject({ _tag: "InvalidRequest" })
  201. expect("classification" in error.reason ? error.reason.classification : undefined).toBeUndefined()
  202. expect(error.reason.message).toBe("Provider request failed with HTTP 400")
  203. }).pipe(Effect.provide(responsesLayer([new Response("invalid parameter", { status: 400 })]))),
  204. )
  205. it.effect("preserves structured provider messages from large error bodies", () =>
  206. Effect.gen(function* () {
  207. const executor = yield* RequestExecutor.Service
  208. const error = yield* executor.execute(request).pipe(Effect.flip)
  209. expectAIError(error)
  210. expect(error.reason).toMatchObject({ _tag: "InvalidRequest", message: largeProviderMessage })
  211. expect(errorHttp(error)?.body).toContain(largeProviderMessage)
  212. expect(errorHttp(error)?.bodyTruncated).toBeUndefined()
  213. }).pipe(
  214. Effect.provide(
  215. responsesLayer([
  216. new Response(
  217. JSON.stringify({
  218. model: "gpt-5.6-sol",
  219. error: { type: "invalid_request", message: largeProviderMessage },
  220. }),
  221. { status: 400 },
  222. ),
  223. ]),
  224. ),
  225. ),
  226. )
  227. it.effect("falls back when structured provider messages are empty", () =>
  228. Effect.gen(function* () {
  229. const executor = yield* RequestExecutor.Service
  230. const error = yield* executor.execute(request).pipe(Effect.flip)
  231. expectAIError(error)
  232. expect(error.reason).toMatchObject({
  233. _tag: "InvalidRequest",
  234. message: "Provider request failed with HTTP 400",
  235. })
  236. }).pipe(Effect.provide(responsesLayer([new Response('{"error":{"message":" "}}', { status: 400 })]))),
  237. )
  238. it.effect("classifies provider rate limits hidden behind HTTP 400", () =>
  239. Effect.gen(function* () {
  240. const classify = (body: string) =>
  241. Effect.gen(function* () {
  242. const executor = yield* RequestExecutor.Service
  243. const error = yield* executor.execute(request).pipe(Effect.flip)
  244. expectAIError(error)
  245. expect(error.reason).toMatchObject({ _tag: "RateLimit" })
  246. }).pipe(Effect.provide(responsesLayer([new Response(body, { status: 400 })])))
  247. yield* classify("Request rate increased too quickly")
  248. yield* classify('{"type":"error","error":{"type":"too_many_requests"}}')
  249. yield* classify('{"type":"error","error":{"code":"rate_limit_exceeded"}}')
  250. }),
  251. )
  252. it.effect("classifies provider overloads hidden behind HTTP 400", () =>
  253. Effect.gen(function* () {
  254. const classify = (body: string) =>
  255. Effect.gen(function* () {
  256. const executor = yield* RequestExecutor.Service
  257. const error = yield* executor.execute(request).pipe(Effect.flip)
  258. expectAIError(error)
  259. expect(error.reason).toMatchObject({ _tag: "ProviderInternal" })
  260. }).pipe(Effect.provide(responsesLayer([new Response(body, { status: 400 })])))
  261. yield* classify('{"code":"resource_exhausted"}')
  262. yield* classify('{"code":"service_unavailable"}')
  263. }),
  264. )
  265. it.effect("returns complete diagnostics for rate limits", () =>
  266. Effect.gen(function* () {
  267. const executor = yield* RequestExecutor.Service
  268. const error = yield* executor.execute(request).pipe(Effect.flip)
  269. expectAIError(error)
  270. expect(error).toMatchObject({
  271. reason: {
  272. _tag: "RateLimit",
  273. retryAfterMs: 0,
  274. rateLimit: { retryAfterMs: 0 },
  275. http: {
  276. requestId: "req_123",
  277. request: {
  278. method: "POST",
  279. url: "https://provider.test/v1/chat?api_key=secret&key=secret&debug=1",
  280. headers: { authorization: "Bearer secret", "x-safe": "visible" },
  281. },
  282. response: {
  283. status: 429,
  284. headers: {
  285. "retry-after-ms": "0",
  286. "x-request-id": "req_123",
  287. "x-api-key": "secret",
  288. },
  289. },
  290. },
  291. },
  292. })
  293. expect(errorHttp(error)?.body).toBe("rate limited")
  294. }).pipe(
  295. Effect.provide(
  296. responsesLayer([
  297. new Response("rate limited", {
  298. status: 429,
  299. headers: { "retry-after-ms": "0", "x-request-id": "req_123", "x-api-key": "secret" },
  300. }),
  301. ]),
  302. ),
  303. ),
  304. )
  305. it.effect("preserves configured header names in diagnostics", () =>
  306. Effect.gen(function* () {
  307. const executor = yield* RequestExecutor.Service
  308. const error = yield* executor.execute(request).pipe(Effect.flip)
  309. expectAIError(error)
  310. expect(errorHttp(error)?.request.headers["x-safe"]).toBe("visible")
  311. expect(errorHttp(error)?.response?.headers["x-safe"]).toBe("response-secret")
  312. }).pipe(
  313. Effect.provide(responsesLayer([new Response("bad", { status: 400, headers: { "x-safe": "response-secret" } })])),
  314. Effect.provideService(Headers.CurrentRedactedNames, ["x-safe"]),
  315. ),
  316. )
  317. it.effect("extracts OpenAI-style rate-limit diagnostics", () =>
  318. Effect.gen(function* () {
  319. const executor = yield* RequestExecutor.Service
  320. const error = yield* executor.execute(request).pipe(Effect.flip)
  321. expectAIError(error)
  322. expect(error.reason).toMatchObject({ _tag: "RateLimit" })
  323. expect(error.reason._tag === "RateLimit" ? error.reason.rateLimit : undefined).toEqual({
  324. retryAfterMs: 0,
  325. limit: { requests: "500", tokens: "30000" },
  326. remaining: { requests: "499", tokens: "29900" },
  327. reset: { requests: "1s", tokens: "10s" },
  328. })
  329. }).pipe(
  330. Effect.provide(
  331. responsesLayer([
  332. new Response("rate limited", {
  333. status: 429,
  334. headers: {
  335. "retry-after-ms": "0",
  336. "x-ratelimit-limit-requests": "500",
  337. "x-ratelimit-limit-tokens": "30000",
  338. "x-ratelimit-remaining-requests": "499",
  339. "x-ratelimit-remaining-tokens": "29900",
  340. "x-ratelimit-reset-requests": "1s",
  341. "x-ratelimit-reset-tokens": "10s",
  342. },
  343. }),
  344. ]),
  345. ),
  346. ),
  347. )
  348. it.effect("extracts Anthropic-style rate-limit diagnostics", () =>
  349. Effect.gen(function* () {
  350. const executor = yield* RequestExecutor.Service
  351. const error = yield* executor.execute(request).pipe(Effect.flip)
  352. expectAIError(error)
  353. expect(error.reason).toMatchObject({ _tag: "ProviderInternal" })
  354. expect(errorHttp(error)?.rateLimit).toEqual({
  355. retryAfterMs: 0,
  356. limit: { requests: "100", "input-tokens": "10000" },
  357. remaining: { requests: "12", "input-tokens": "9000" },
  358. reset: { requests: "2026-05-06T12:00:00Z", "input-tokens": "2026-05-06T12:00:10Z" },
  359. })
  360. }).pipe(
  361. Effect.provide(
  362. responsesLayer([
  363. new Response("overloaded", {
  364. status: 529,
  365. headers: {
  366. "retry-after-ms": "0",
  367. "anthropic-ratelimit-requests-limit": "100",
  368. "anthropic-ratelimit-requests-remaining": "12",
  369. "anthropic-ratelimit-requests-reset": "2026-05-06T12:00:00Z",
  370. "anthropic-ratelimit-input-tokens-limit": "10000",
  371. "anthropic-ratelimit-input-tokens-remaining": "9000",
  372. "anthropic-ratelimit-input-tokens-reset": "2026-05-06T12:00:10Z",
  373. },
  374. }),
  375. ]),
  376. ),
  377. ),
  378. )
  379. it.effect("returns provider status failures without retrying", () =>
  380. Effect.gen(function* () {
  381. const attempts = yield* Ref.make(0)
  382. const error = yield* Effect.gen(function* () {
  383. const executor = yield* RequestExecutor.Service
  384. return yield* executor.execute(request).pipe(Effect.flip)
  385. }).pipe(
  386. Effect.provide(
  387. countedResponsesLayer(attempts, [
  388. new Response("busy", { status: 503, headers: { "retry-after-ms": "0" } }),
  389. new Response("ok", { status: 200 }),
  390. ]),
  391. ),
  392. )
  393. expectAIError(error)
  394. expect(error.reason).toMatchObject({ _tag: "ProviderInternal", status: 503 })
  395. expect(yield* Ref.get(attempts)).toBe(1)
  396. }),
  397. )
  398. it.effect("marks 504 and 529 status responses as provider-internal", () =>
  399. Effect.gen(function* () {
  400. const failWith = (status: number) =>
  401. Effect.gen(function* () {
  402. const executor = yield* RequestExecutor.Service
  403. const error = yield* executor.execute(request).pipe(Effect.flip)
  404. expectAIError(error)
  405. expect(error.reason).toMatchObject({ _tag: "ProviderInternal", status })
  406. }).pipe(
  407. Effect.provide(
  408. responsesLayer([
  409. new Response("provider failure", {
  410. status,
  411. headers: { "retry-after-ms": "0" },
  412. }),
  413. ]),
  414. ),
  415. )
  416. yield* failWith(504)
  417. yield* failWith(529)
  418. }),
  419. )
  420. it.effect("preserves large authentication error bodies", () =>
  421. Effect.gen(function* () {
  422. const executor = yield* RequestExecutor.Service
  423. const error = yield* executor.execute(request).pipe(Effect.flip)
  424. expectAIError(error)
  425. expect(error.reason).toMatchObject({ _tag: "Authentication" })
  426. expect(errorHttp(error)?.bodyTruncated).toBeUndefined()
  427. expect(errorHttp(error)?.body).toHaveLength(20_000)
  428. }).pipe(
  429. Effect.provide(
  430. responsesLayer([
  431. new Response("x".repeat(20_000), { status: 401 }),
  432. new Response("should not retry", { status: 200 }),
  433. ]),
  434. ),
  435. ),
  436. )
  437. it.effect("preserves response body fields", () =>
  438. Effect.gen(function* () {
  439. const executor = yield* RequestExecutor.Service
  440. const error = yield* executor.execute(request).pipe(Effect.flip)
  441. expectAIError(error)
  442. expect(errorHttp(error)?.body).toBe(
  443. '{"error":{"message":"bad","key":"body-secret","detail":"api_key=query-secret"}}',
  444. )
  445. }).pipe(
  446. Effect.provide(
  447. responsesLayer([
  448. new Response('{"error":{"message":"bad","key":"body-secret","detail":"api_key=query-secret"}}', {
  449. status: 400,
  450. }),
  451. ]),
  452. ),
  453. ),
  454. )
  455. it.effect("preserves echoed request values in response bodies", () =>
  456. Effect.gen(function* () {
  457. const executor = yield* RequestExecutor.Service
  458. const error = yield* executor.execute(secretRequest).pipe(Effect.flip)
  459. expectAIError(error)
  460. expect(errorHttp(error)?.body).toBe("provider echoed query-secret-123 and authorization header-secret-456")
  461. }).pipe(
  462. Effect.provide(
  463. responsesLayer([
  464. new Response("provider echoed query-secret-123 and authorization header-secret-456", { status: 400 }),
  465. ]),
  466. ),
  467. ),
  468. )
  469. it.effect("does not re-execute after a successful response reaches stream parsing", () =>
  470. Effect.gen(function* () {
  471. const attempts = yield* Ref.make(0)
  472. const model = OpenAIChat.route
  473. .with({ endpoint: { baseURL: "https://api.openai.test/v1" } })
  474. .model({ id: "gpt-4o-mini" })
  475. const error = yield* LLMClient.generate(LLM.request({ model, prompt: "Say hello." })).pipe(
  476. Effect.provide(
  477. dynamicResponse((input) =>
  478. Ref.update(attempts, (value) => value + 1).pipe(
  479. Effect.as(
  480. input.respond(
  481. sseRaw(
  482. `data: ${JSON.stringify(deltaChunk({ role: "assistant", content: "Hello" }))}`,
  483. "data: not-json",
  484. ),
  485. { headers: { "content-type": "text/event-stream" } },
  486. ),
  487. ),
  488. ),
  489. ),
  490. ),
  491. Effect.flip,
  492. )
  493. expectAIError(error)
  494. expect(error.reason).toMatchObject({ _tag: "InvalidProviderOutput" })
  495. expect(yield* Ref.get(attempts)).toBe(1)
  496. }),
  497. )
  498. })
  499. describe("WebSocket channel execution", () => {
  500. const model = OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responses("gpt-4.1-mini")
  501. const request = LLM.request({ model, prompt: "Say hello." })
  502. const frames = [
  503. JSON.stringify({ type: "response.output_text.delta", item_id: "msg_1", delta: "Hi" }),
  504. JSON.stringify({ type: "response.completed", response: { id: "resp_1" } }),
  505. ]
  506. it.effect("runs a channel driver through the direct executor", () =>
  507. Effect.gen(function* () {
  508. const sent = yield* Ref.make("")
  509. const closed = yield* Ref.make(false)
  510. const observed = yield* Ref.make(0)
  511. const webSocket = WebSocketTransport.makeDirect({
  512. open: () =>
  513. Effect.succeed({
  514. sendText: (message) => Ref.set(sent, message),
  515. messages: Stream.make("one", "done", "late"),
  516. close: Ref.set(closed, true),
  517. }),
  518. })
  519. const received = yield* Effect.scoped(
  520. Effect.gen(function* () {
  521. const execution = yield* webSocket.execute({
  522. id: "exchange_1",
  523. connect: { url: "wss://api.openai.test/v1/responses", headers: Headers.empty },
  524. fallback: () => Stream.empty,
  525. driver: {
  526. create: () => Effect.succeed({ message: "create", mode: "full" }),
  527. observe: (_create, frame) =>
  528. Ref.update(observed, (value) => value + 1).pipe(
  529. Effect.as(
  530. frame === "done" ? { type: "completed" as const, frame } : { type: "frame" as const, frame },
  531. ),
  532. ),
  533. },
  534. })
  535. return yield* Stream.runCollect(execution.frames)
  536. }),
  537. )
  538. expect(Array.from(received)).toEqual(["one", "done"])
  539. expect(yield* Ref.get(sent)).toBe("create")
  540. expect(yield* Ref.get(observed)).toBe(2)
  541. expect(yield* Ref.get(closed)).toBe(true)
  542. }),
  543. )
  544. it.effect("rejects a closed socket before attempting to send", () =>
  545. Effect.gen(function* () {
  546. class ClosedBeforeSend extends EventTarget {
  547. readyState = globalThis.WebSocket.OPEN
  548. sends = 0
  549. send() {
  550. this.sends++
  551. }
  552. close() {}
  553. }
  554. const socket = new ClosedBeforeSend()
  555. const connection = yield* WebSocketTransport.fromWebSocket(
  556. // oxlint-disable-next-line typescript-eslint/no-unsafe-type-assertion
  557. socket as unknown as globalThis.WebSocket,
  558. { url: "wss://api.openai.test/v1/responses", headers: Headers.empty },
  559. )
  560. socket.readyState = globalThis.WebSocket.CLOSED
  561. const error = yield* connection.sendText("create").pipe(Effect.flip)
  562. expect(error.reason).toMatchObject({ _tag: "Transport", phase: "send", delivery: "not-sent" })
  563. expect(socket.sends).toBe(0)
  564. yield* connection.close
  565. }),
  566. )
  567. it.effect("uses HTTP when no per-call WebSocket executor is provided", () =>
  568. Effect.gen(function* () {
  569. const response = yield* LLMClient.generate(request).pipe(Effect.provide(fixedResponse(sseEvents(...frames))))
  570. expect(response.text).toBe("Hi")
  571. }),
  572. )
  573. it.effect("commits channel execution only after complete consumption", () =>
  574. Effect.gen(function* () {
  575. const commits = yield* Ref.make(0)
  576. const executor = (input: Stream.Stream<string, AIError>): WebSocketChannelExecutor => ({
  577. execute: () =>
  578. Effect.succeed({
  579. frames: input,
  580. complete: Ref.update(commits, (value) => value + 1),
  581. }),
  582. })
  583. const response = yield* LLMClient.generate(request, {
  584. webSocket: executor(Stream.fromArray(frames)),
  585. }).pipe(Effect.provide(fixedResponse("")))
  586. expect(response.text).toBe("Hi")
  587. expect(yield* Ref.get(commits)).toBe(1)
  588. yield* LLMClient.generate(request, { webSocket: executor(Stream.make("not-json")) }).pipe(
  589. Effect.provide(fixedResponse("")),
  590. Effect.flip,
  591. )
  592. expect(yield* Ref.get(commits)).toBe(1)
  593. yield* LLMClient.stream(request, { webSocket: executor(Stream.fromArray(frames)) }).pipe(
  594. Stream.take(1),
  595. Stream.runDrain,
  596. Effect.provide(fixedResponse("")),
  597. )
  598. expect(yield* Ref.get(commits)).toBe(1)
  599. }),
  600. )
  601. it.effect("does not commit interrupted channel execution", () =>
  602. Effect.gen(function* () {
  603. const commits = yield* Ref.make(0)
  604. const started = yield* Deferred.make<void>()
  605. const executor: WebSocketChannelExecutor = {
  606. execute: () =>
  607. Effect.succeed({
  608. frames: Stream.fromEffect(
  609. Deferred.succeed(started, undefined).pipe(
  610. Effect.as(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })),
  611. ),
  612. ).pipe(Stream.concat(Stream.never)),
  613. complete: Ref.update(commits, (value) => value + 1),
  614. }),
  615. }
  616. const fiber = yield* LLMClient.stream(request, { webSocket: executor }).pipe(
  617. Stream.runDrain,
  618. Effect.provide(fixedResponse("")),
  619. Effect.forkChild({ startImmediately: true }),
  620. )
  621. yield* Deferred.await(started)
  622. yield* Fiber.interrupt(fiber)
  623. expect(yield* Ref.get(commits)).toBe(0)
  624. }),
  625. )
  626. })