Răsfoiți Sursa

test(opencode): stabilize provider timeout fixtures

Dax Raad 2 luni în urmă
părinte
comite
ac0d5b8d7b
1 a modificat fișierele cu 64 adăugiri și 55 ștergeri
  1. 64 55
      packages/opencode/test/provider/header-timeout.test.ts

+ 64 - 55
packages/opencode/test/provider/header-timeout.test.ts

@@ -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(
-    ({ 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", () =>
 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())
-
-        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
+  Effect.gen(function* () {
+    const server = yield* Effect.acquireRelease(
+      Effect.promise(() => delayedBodyServer(250)),
+      (server) => Effect.sync(() => server.server.close()),
+    )
+
+    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
+              }
+            } 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", () =>
 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}` }
   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(() => {