|
|
@@ -4,6 +4,7 @@ import { Headers, HttpClientRequest } from "effect/unstable/http"
|
|
|
import {
|
|
|
LLM,
|
|
|
AIError,
|
|
|
+ HttpOptions,
|
|
|
LLMEvent,
|
|
|
LLMRequest,
|
|
|
Message,
|
|
|
@@ -217,19 +218,19 @@ describe("OpenAI Responses route", () => {
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
- it.effect("prepares OpenAI Responses WebSocket target", () =>
|
|
|
+ it.effect("prepares one OpenAI Responses route for either transport", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const prepared = yield* compileRequest(
|
|
|
LLMRequest.update(request, {
|
|
|
- model: OpenAIResponses.webSocketRoute
|
|
|
+ model: OpenAIResponses.route
|
|
|
.with({ endpoint: { baseURL: "https://api.openai.test/v1/" }, auth: Auth.bearer("test") })
|
|
|
.model({ id: "gpt-4.1-mini" }),
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
- expect(prepared.route).toBe("openai-responses-websocket")
|
|
|
+ expect(prepared.route).toBe("openai-responses")
|
|
|
expect(prepared.protocol).toBe("openai-responses")
|
|
|
- expect(prepared.metadata).toEqual({ transport: "websocket-json" })
|
|
|
+ expect(prepared.metadata).toEqual({ transport: "http-json" })
|
|
|
expect(prepared.body).toMatchObject({ model: "gpt-4.1-mini", store: false, stream: true })
|
|
|
}),
|
|
|
)
|
|
|
@@ -237,7 +238,11 @@ describe("OpenAI Responses route", () => {
|
|
|
it.effect("streams OpenAI Responses over WebSocket", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const sent: string[] = []
|
|
|
- const opened: Array<{ readonly url: string; readonly authorization: string | undefined }> = []
|
|
|
+ const opened: Array<{
|
|
|
+ readonly url: string
|
|
|
+ readonly authorization: string | undefined
|
|
|
+ readonly protocol: string | undefined
|
|
|
+ }> = []
|
|
|
let closed = false
|
|
|
const deps = Layer.succeed(
|
|
|
RequestExecutor.Service,
|
|
|
@@ -250,10 +255,15 @@ describe("OpenAI Responses route", () => {
|
|
|
Effect.succeed({
|
|
|
sendText: (message) =>
|
|
|
Effect.sync(() => {
|
|
|
- opened.push({ url: input.url, authorization: input.headers.authorization })
|
|
|
+ opened.push({
|
|
|
+ url: input.url,
|
|
|
+ authorization: input.headers.authorization,
|
|
|
+ protocol: input.headers["openai-beta"],
|
|
|
+ })
|
|
|
sent.push(message)
|
|
|
}),
|
|
|
messages: Stream.fromArray([
|
|
|
+ ProviderShared.encodeJson({ type: "response.created", response: { id: "resp_ws" } }),
|
|
|
ProviderShared.encodeJson({ type: "response.output_text.delta", item_id: "msg_1", delta: "Hi" }),
|
|
|
ProviderShared.encodeJson({ type: "response.completed", response: { id: "resp_ws" } }),
|
|
|
]),
|
|
|
@@ -264,16 +274,24 @@ describe("OpenAI Responses route", () => {
|
|
|
})
|
|
|
const response = yield* LLMClient.generate(
|
|
|
LLM.request({
|
|
|
- model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responsesWebSocket(
|
|
|
- "gpt-4.1-mini",
|
|
|
- ),
|
|
|
+ model: OpenAI.configure({
|
|
|
+ baseURL: "https://api.openai.test/v1/",
|
|
|
+ apiKey: "test",
|
|
|
+ headers: { "openai-beta": "custom-protocol" },
|
|
|
+ }).responses("gpt-4.1-mini"),
|
|
|
prompt: "Say hello.",
|
|
|
}),
|
|
|
{ webSocket },
|
|
|
).pipe(Effect.provide(LLMClient.layer.pipe(Layer.provide(deps))))
|
|
|
|
|
|
expect(response.text).toBe("Hi")
|
|
|
- expect(opened).toEqual([{ url: "wss://api.openai.test/v1/responses", authorization: "Bearer test" }])
|
|
|
+ expect(opened).toEqual([
|
|
|
+ {
|
|
|
+ url: "wss://api.openai.test/v1/responses",
|
|
|
+ authorization: "Bearer test",
|
|
|
+ protocol: "custom-protocol",
|
|
|
+ },
|
|
|
+ ])
|
|
|
expect(closed).toBe(true)
|
|
|
expect(sent).toHaveLength(1)
|
|
|
expect(JSON.parse(sent[0])).toEqual({
|
|
|
@@ -285,6 +303,136 @@ describe("OpenAI Responses route", () => {
|
|
|
}),
|
|
|
)
|
|
|
|
|
|
+ it.effect("rejects out-of-order and mismatched WebSocket response events", () =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const streams = [
|
|
|
+ Stream.fromArray([
|
|
|
+ ProviderShared.encodeJson({ type: "response.output_text.delta", item_id: "late", delta: "Late" }),
|
|
|
+ ProviderShared.encodeJson({ type: "response.completed", response: { id: "resp_old" } }),
|
|
|
+ ]),
|
|
|
+ Stream.fromArray([
|
|
|
+ ProviderShared.encodeJson({ type: "response.created", response: { id: "resp_new" } }),
|
|
|
+ ProviderShared.encodeJson({ type: "response.completed", response: { id: "resp_old" } }),
|
|
|
+ ]),
|
|
|
+ ]
|
|
|
+ const webSocket = WebSocketTransport.makeDirect({
|
|
|
+ open: () =>
|
|
|
+ Effect.succeed({
|
|
|
+ sendText: () => Effect.void,
|
|
|
+ messages: streams.shift() ?? Stream.die("unexpected WebSocket open"),
|
|
|
+ close: Effect.void,
|
|
|
+ }),
|
|
|
+ })
|
|
|
+ const deps = Layer.succeed(
|
|
|
+ RequestExecutor.Service,
|
|
|
+ RequestExecutor.Service.of({ execute: () => Effect.die("unexpected HTTP request") }),
|
|
|
+ )
|
|
|
+ const model = OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responses(
|
|
|
+ "gpt-4.1-mini",
|
|
|
+ )
|
|
|
+
|
|
|
+ const errors = yield* Effect.forEach(["late", "mismatch"], (prompt) =>
|
|
|
+ LLMClient.generate(LLM.request({ model, prompt }), { webSocket }).pipe(
|
|
|
+ Effect.provide(LLMClient.layer.pipe(Layer.provide(deps))),
|
|
|
+ Effect.flip,
|
|
|
+ ),
|
|
|
+ )
|
|
|
+
|
|
|
+ expect(errors.map((error) => error.reason._tag)).toEqual(["InvalidProviderOutput", "InvalidProviderOutput"])
|
|
|
+ expect(errors[0]?.message).toContain("before response.created")
|
|
|
+ expect(errors[1]?.message).toContain("response ID changed")
|
|
|
+ }),
|
|
|
+ )
|
|
|
+
|
|
|
+ it.effect("builds WebSocket and HTTP fallback from the same final request", () =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const attempts = yield* Ref.make(0)
|
|
|
+ const message = yield* Ref.make("")
|
|
|
+ const body = yield* Ref.make("")
|
|
|
+ const response = yield* LLMClient.generate(
|
|
|
+ LLM.request({
|
|
|
+ model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responses("gpt-4.1-mini"),
|
|
|
+ prompt: "Say hello.",
|
|
|
+ http: {
|
|
|
+ body: {
|
|
|
+ model: "overlaid-model",
|
|
|
+ metadata: { source: "overlay" },
|
|
|
+ stream_options: { include_usage: true },
|
|
|
+ background: true,
|
|
|
+ },
|
|
|
+ headers: { "x-request": "request" },
|
|
|
+ query: { mode: "test" },
|
|
|
+ },
|
|
|
+ }),
|
|
|
+ {
|
|
|
+ webSocket: {
|
|
|
+ execute: (exchange) =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ expect(exchange.connect.rotateAfterMs).toBe(55 * 60 * 1000)
|
|
|
+ expect(exchange.connect.headers["openai-beta"]).toBe("responses_websockets=2026-02-06")
|
|
|
+ expect(exchange.connect.headers["content-length"]).toBeUndefined()
|
|
|
+ yield* exchange.driver
|
|
|
+ .create(undefined)
|
|
|
+ .pipe(Effect.flatMap((create) => Ref.set(message, create.message)))
|
|
|
+ return { frames: exchange.fallback(), complete: Effect.void }
|
|
|
+ }),
|
|
|
+ },
|
|
|
+ },
|
|
|
+ ).pipe(
|
|
|
+ Effect.provide(
|
|
|
+ dynamicResponse((input) =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ yield* Ref.update(attempts, (value) => value + 1)
|
|
|
+ yield* Ref.set(body, input.text)
|
|
|
+ expect(input.request.url).toBe("https://api.openai.test/v1/responses?mode=test")
|
|
|
+ expect(input.request.headers.authorization).toBe("Bearer test")
|
|
|
+ expect(input.request.headers["x-request"]).toBe("request")
|
|
|
+ return input.respond(sseEvents({ type: "response.completed", response: {} }), {
|
|
|
+ headers: { "content-type": "text/event-stream" },
|
|
|
+ })
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ ),
|
|
|
+ )
|
|
|
+
|
|
|
+ const httpBody = JSON.parse(yield* Ref.get(body))
|
|
|
+ const { stream: _stream, stream_options: _streamOptions, background: _background, ...shared } = httpBody
|
|
|
+ expect(response.finishReason?.normalized).toBe("stop")
|
|
|
+ expect(yield* Ref.get(attempts)).toBe(1)
|
|
|
+ expect(JSON.parse(yield* Ref.get(message))).toEqual({ type: "response.create", ...shared })
|
|
|
+ expect(httpBody).toMatchObject({
|
|
|
+ model: "overlaid-model",
|
|
|
+ metadata: { source: "overlay" },
|
|
|
+ stream: true,
|
|
|
+ stream_options: { include_usage: true },
|
|
|
+ background: true,
|
|
|
+ })
|
|
|
+ }),
|
|
|
+ )
|
|
|
+
|
|
|
+ it.effect("uses exactly one HTTP request when no WebSocket executor is supplied", () =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const attempts = yield* Ref.make(0)
|
|
|
+ yield* LLMClient.generate(
|
|
|
+ LLMRequest.update(request, { http: new HttpOptions({ body: { input: "raw-http-input" } }) }),
|
|
|
+ ).pipe(
|
|
|
+ Effect.provide(
|
|
|
+ dynamicResponse((input) =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ yield* Ref.update(attempts, (value) => value + 1)
|
|
|
+ expect(JSON.parse(input.text).input).toBe("raw-http-input")
|
|
|
+ return input.respond(sseEvents({ type: "response.completed", response: {} }), {
|
|
|
+ headers: { "content-type": "text/event-stream" },
|
|
|
+ })
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ ),
|
|
|
+ )
|
|
|
+
|
|
|
+ expect(yield* Ref.get(attempts)).toBe(1)
|
|
|
+ }),
|
|
|
+ )
|
|
|
+
|
|
|
it.effect("closes a direct WebSocket execution after partial consumption", () =>
|
|
|
Effect.gen(function* () {
|
|
|
const closed = yield* Ref.make(false)
|
|
|
@@ -293,6 +441,7 @@ describe("OpenAI Responses route", () => {
|
|
|
Effect.succeed({
|
|
|
sendText: () => Effect.void,
|
|
|
messages: Stream.fromArray([
|
|
|
+ ProviderShared.encodeJson({ type: "response.created", response: { id: "resp_ws" } }),
|
|
|
ProviderShared.encodeJson({ type: "response.output_text.delta", item_id: "msg_1", delta: "Hi" }),
|
|
|
ProviderShared.encodeJson({ type: "response.completed", response: { id: "resp_ws" } }),
|
|
|
]),
|
|
|
@@ -302,9 +451,7 @@ describe("OpenAI Responses route", () => {
|
|
|
|
|
|
yield* LLMClient.stream(
|
|
|
LLM.request({
|
|
|
- model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responsesWebSocket(
|
|
|
- "gpt-4.1-mini",
|
|
|
- ),
|
|
|
+ model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responses("gpt-4.1-mini"),
|
|
|
prompt: "Say hello.",
|
|
|
}),
|
|
|
{ webSocket },
|
|
|
@@ -347,7 +494,7 @@ describe("OpenAI Responses route", () => {
|
|
|
const errors = yield* Effect.forEach(events, (event) =>
|
|
|
LLMClient.generate(
|
|
|
LLM.request({
|
|
|
- model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responsesWebSocket(
|
|
|
+ model: OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responses(
|
|
|
"gpt-4.1-mini",
|
|
|
),
|
|
|
prompt: "Say hello.",
|
|
|
@@ -391,11 +538,18 @@ describe("OpenAI Responses route", () => {
|
|
|
const failure = new AIError({
|
|
|
module: "test",
|
|
|
method: "receive",
|
|
|
- reason: new TransportReason({ message: "socket closed", phase: "close" }),
|
|
|
+ reason: new TransportReason({
|
|
|
+ message: "socket closed",
|
|
|
+ transport: "websocket",
|
|
|
+ operation: "read",
|
|
|
+ phase: "close",
|
|
|
+ }),
|
|
|
})
|
|
|
const streams = [
|
|
|
Stream.fail(failure),
|
|
|
- Stream.make(ProviderShared.encodeJson({ type: "response.created" })).pipe(Stream.concat(Stream.fail(failure))),
|
|
|
+ Stream.make(ProviderShared.encodeJson({ type: "response.created", response: { id: "resp_observed" } })).pipe(
|
|
|
+ Stream.concat(Stream.fail(failure)),
|
|
|
+ ),
|
|
|
]
|
|
|
const deps = Layer.succeed(
|
|
|
RequestExecutor.Service,
|
|
|
@@ -409,7 +563,7 @@ describe("OpenAI Responses route", () => {
|
|
|
close: Effect.void,
|
|
|
}),
|
|
|
})
|
|
|
- const model = OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responsesWebSocket(
|
|
|
+ const model = OpenAI.configure({ baseURL: "https://api.openai.test/v1/", apiKey: "test" }).responses(
|
|
|
"gpt-4.1-mini",
|
|
|
)
|
|
|
|
|
|
@@ -469,7 +623,7 @@ describe("OpenAI Responses route", () => {
|
|
|
yield* LLMClient.generate(
|
|
|
LLMRequest.update(request, {
|
|
|
model: Azure.configure({
|
|
|
- baseURL: "https://opencode-test.openai.azure.com/openai/v1/",
|
|
|
+ baseURL: "https://opencode-test.openai.azure.com/openai/",
|
|
|
apiKey: "azure-key",
|
|
|
headers: { authorization: "Bearer stale" },
|
|
|
}).responses("gpt-4.1-mini"),
|