diff --git a/packages/opencode/test/provider/header-timeout.test.ts b/packages/opencode/test/provider/header-timeout.test.ts index 543bb5f30d..8860c39c32 100644 --- a/packages/opencode/test/provider/header-timeout.test.ts +++ b/packages/opencode/test/provider/header-timeout.test.ts @@ -3,9 +3,8 @@ import { createServer, type Server } from "node:http" import { streamText } from "ai" import { Effect, Layer } from "effect" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" -import { disposeAllInstances, provideTmpdirInstance, provideTmpdirServer } from "../fixture/fixture" +import { disposeAllInstances, provideTmpdirInstance } from "../fixture/fixture" import { testEffect } from "../lib/effect" -import { reply, TestLLMServer } from "../lib/llm-server" import { testProviderConfig } from "../lib/test-provider" import { Env } from "@/env" import { Plugin } from "@/plugin" @@ -22,70 +21,66 @@ const it = testEffect( Provider.defaultLayer, Env.defaultLayer, Plugin.defaultLayer, - TestLLMServer.layer, CrossSpawnSpawner.defaultLayer, ), ) it.live("headerTimeout does not abort delayed SSE body after headers arrive", () => - provideTmpdirServer( - ({ llm }) => - Effect.gen(function* () { - yield* llm.push(reply().wait(Bun.sleep(250)).text("late").stop()) + Effect.gen(function* () { + const server = yield* Effect.acquireRelease( + Effect.promise(() => delayedBodyServer(250)), + (server) => Effect.sync(() => server.server.close()), + ) - const provider = yield* Provider.Service - const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model")) - const result = streamText({ - model: yield* provider.getLanguage(model), - messages: [{ role: "user", content: "hello" }], - }) + yield* provideTmpdirInstance( + () => + Effect.gen(function* () { + const provider = yield* Provider.Service + const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model")) + const result = streamText({ + model: yield* provider.getLanguage(model), + messages: [{ role: "user", content: "hello" }], + }) - expect(yield* Effect.promise(() => result.text)).toBe("late") - }), - { - config: (url) => { - const config = testProviderConfig(url) - return { - ...config, - provider: { - test: { - ...config.provider.test, - options: { ...config.provider.test.options, headerTimeout: 50 }, - }, - }, - } - }, - }, - ), + expect(yield* Effect.promise(() => result.text)).toBe("late") + }), + { config: providerConfig(server.url, { headerTimeout: 50 }) }, + ) + }), ) it.live("chunkTimeout raises a response stream error when SSE body stalls", () => - provideTmpdirServer( - ({ llm }) => - Effect.gen(function* () { - yield* llm.push(reply().wait(Bun.sleep(250)).text("late").stop()) + Effect.gen(function* () { + const server = yield* Effect.acquireRelease( + Effect.promise(() => delayedBodyServer(250)), + (server) => Effect.sync(() => server.server.close()), + ) - const provider = yield* Provider.Service - const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model")) - const result = streamText({ - model: yield* provider.getLanguage(model), - onError() {}, - messages: [{ role: "user", content: "hello" }], - }) + yield* provideTmpdirInstance( + () => + Effect.gen(function* () { + const provider = yield* Provider.Service + const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model")) + const result = streamText({ + model: yield* provider.getLanguage(model), + onError() {}, + messages: [{ role: "user", content: "hello" }], + }) - const error = yield* Effect.promise(async () => { - try { - for await (const part of result.fullStream) { - if (part.type === "error") return part.error + const error = yield* Effect.promise(async () => { + try { + for await (const part of result.fullStream) { + if (part.type === "error") return part.error + } + } catch (error) { + return error } - } catch (error) { - return error - } - }) - expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) - }), - { config: (url) => providerConfig(url, { chunkTimeout: 50 }) }, - ), + }) + expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) + }), + { config: providerConfig(server.url, { chunkTimeout: 50 }) }, + ) + }), ) it.live("headerTimeout aborts when response headers do not arrive", () => @@ -205,6 +200,20 @@ async function delayedHeaderServer(delay: number): Promise<{ server: Server; url return { server, url: `http://127.0.0.1:${address.port}` } } +async function delayedBodyServer(delay: number): Promise<{ server: Server; url: string }> { + const server = createServer((_, res) => { + res.writeHead(200, { "content-type": "text/event-stream" }) + res.flushHeaders() + setTimeout(() => { + res.end('data: {"choices":[{"delta":{"content":"late"}}]}\n\ndata: [DONE]\n\n') + }, delay) + }) + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)) + const address = server.address() + if (!address || typeof address === "string") throw new Error("server did not bind to a TCP port") + return { server, url: `http://127.0.0.1:${address.port}` } +} + function withAuthContent(self: Effect.Effect, value: Record = defaultAuthContent()) { return Effect.acquireUseRelease( Effect.sync(() => {