test(opencode): flush headers before stream timeout

This commit is contained in:
Dax Raad 2026-05-29 23:36:39 -04:00
commit 6e1b390532

View file

@ -3,9 +3,8 @@ import { createServer, type Server } from "node:http"
import { streamText } from "ai" import { streamText } from "ai"
import { Effect, Layer } from "effect" import { Effect, Layer } from "effect"
import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" 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 { testEffect } from "../lib/effect"
import { reply, TestLLMServer } from "../lib/llm-server"
import { testProviderConfig } from "../lib/test-provider" import { testProviderConfig } from "../lib/test-provider"
import { Env } from "@/env" import { Env } from "@/env"
import { Plugin } from "@/plugin" import { Plugin } from "@/plugin"
@ -22,70 +21,66 @@ const it = testEffect(
Provider.defaultLayer, Provider.defaultLayer,
Env.defaultLayer, Env.defaultLayer,
Plugin.defaultLayer, Plugin.defaultLayer,
TestLLMServer.layer,
CrossSpawnSpawner.defaultLayer, CrossSpawnSpawner.defaultLayer,
), ),
) )
it.live("headerTimeout does not abort delayed SSE body after headers arrive", () => it.live("headerTimeout does not abort delayed SSE body after headers arrive", () =>
provideTmpdirServer( Effect.gen(function* () {
({ llm }) => const server = yield* Effect.acquireRelease(
Effect.gen(function* () { Effect.promise(() => delayedBodyServer(250)),
yield* llm.push(reply().wait(deferredSleep(250)).text("late").stop()) (server) => Effect.sync(() => server.server.close()),
)
const provider = yield* Provider.Service yield* provideTmpdirInstance(
const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model")) () =>
const result = streamText({ Effect.gen(function* () {
model: yield* provider.getLanguage(model), const provider = yield* Provider.Service
messages: [{ role: "user", content: "hello" }], 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") expect(yield* Effect.promise(() => result.text)).toBe("late")
}), }),
{ { config: providerConfig(server.url, { headerTimeout: 50 }) },
config: (url) => { )
const config = testProviderConfig(url) }),
return {
...config,
provider: {
test: {
...config.provider.test,
options: { ...config.provider.test.options, headerTimeout: 50 },
},
},
}
},
},
),
) )
it.live("chunkTimeout raises a response stream error when SSE body stalls", () => it.live("chunkTimeout raises a response stream error when SSE body stalls", () =>
provideTmpdirServer( Effect.gen(function* () {
({ llm }) => const server = yield* Effect.acquireRelease(
Effect.gen(function* () { Effect.promise(() => delayedBodyServer(250)),
yield* llm.push(reply().wait(deferredSleep(250)).text("late").stop()) (server) => Effect.sync(() => server.server.close()),
)
const provider = yield* Provider.Service yield* provideTmpdirInstance(
const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model")) () =>
const result = streamText({ Effect.gen(function* () {
model: yield* provider.getLanguage(model), const provider = yield* Provider.Service
onError() {}, const model = yield* provider.getModel(ProviderID.make("test"), ModelID.make("test-model"))
messages: [{ role: "user", content: "hello" }], const result = streamText({
}) model: yield* provider.getLanguage(model),
onError() {},
messages: [{ role: "user", content: "hello" }],
})
const error = yield* Effect.promise(async () => { const error = yield* Effect.promise(async () => {
try { try {
for await (const part of result.fullStream) { for await (const part of result.fullStream) {
if (part.type === "error") return part.error if (part.type === "error") return part.error
}
} catch (error) {
return error
} }
} catch (error) { })
return error expect(error).toBeInstanceOf(ProviderError.ResponseStreamError)
} }),
}) { config: providerConfig(server.url, { chunkTimeout: 50 }) },
expect(error).toBeInstanceOf(ProviderError.ResponseStreamError) )
}), }),
{ config: (url) => providerConfig(url, { chunkTimeout: 50 }) },
),
) )
it.live("headerTimeout aborts when response headers do not arrive", () => it.live("headerTimeout aborts when response headers do not arrive", () =>
@ -192,12 +187,6 @@ function providerConfig(url: string, options: Record<string, unknown> = {}) {
} }
} }
function deferredSleep(ms: number): PromiseLike<void> {
return {
then: (resolve, reject) => Bun.sleep(ms).then(resolve, reject),
}
}
async function delayedHeaderServer(delay: number): Promise<{ server: Server; url: string }> { async function delayedHeaderServer(delay: number): Promise<{ server: Server; url: string }> {
const server = createServer((_, res) => { const server = createServer((_, res) => {
setTimeout(() => { setTimeout(() => {
@ -211,6 +200,20 @@ async function delayedHeaderServer(delay: number): Promise<{ server: Server; url
return { server, url: `http://127.0.0.1:${address.port}` } 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<void>((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<A, E, R>(self: Effect.Effect<A, E, R>, value: Record<string, unknown> = defaultAuthContent()) { function withAuthContent<A, E, R>(self: Effect.Effect<A, E, R>, value: Record<string, unknown> = defaultAuthContent()) {
return Effect.acquireUseRelease( return Effect.acquireUseRelease(
Effect.sync(() => { Effect.sync(() => {