Files
anomalyco_opencode/packages/opencode/test/lib/sse.ts
T
Kit Langton ddbd119dcb test: drop AppRuntime usage from tests
Removes all direct AppRuntime usage from packages/opencode/test —
`grep -r AppRuntime packages/opencode/test` now returns nothing.

Two patterns are applied:

1. Event tests rewritten in the httpapi-cors.test.ts style. The two
   /event SSE tests now serve HttpApiApp.routes on
   NodeHttpServer.layerTest and hit them via HttpClient. Pub/sub
   identity with the in-process routes is preserved via a new opt-in
   `testEffectShared` (in test/lib/effect.ts) that builds the test
   layer through the shared process-wide memoMap so Bus.defaultLayer
   resolves to the same Bus.Service the routes subscribed to.

   The SSE reader helpers move to test/lib/sse.ts and use HttpClient +
   Effect.Stream + Queue<SseEvent>.

   The D7 diagnostic case is removed: the AppRuntime-vs-test-runtime
   distinction it diagnosed no longer exists.

2. Surgical swap in the remaining four files
   (provider/{amazon-bedrock,provider}, session/llm,
   control-plane/workspace). Each `AppRuntime.runPromise(...)` becomes a
   module-level `ManagedRuntime.make(Service.defaultLayer, { memoMap })`.
   The shared memoMap preserves service identity, so behavior is
   unchanged.

Tests: event (9/9), amazon-bedrock (19/19), provider (84/84),
workspace (35/35); 147 pass across the 5 affected files. The 3
pre-existing failures in session/llm.test.ts are independent of this
change (verified by stashing the diff).
2026-05-20 20:50:32 -04:00

64 lines
2.1 KiB
TypeScript

import { Effect, Queue, Schema, Stream } from "effect"
import { HttpClient, HttpClientRequest } from "effect/unstable/http"
import { EventPaths } from "../../src/server/routes/instance/httpapi/groups/event"
export const SseEvent = Schema.Struct({
id: Schema.optional(Schema.String),
type: Schema.String,
properties: Schema.Record(Schema.String, Schema.Any),
})
export type SseEvent = Schema.Schema.Type<typeof SseEvent>
function decodeFrames(text: string): SseEvent[] {
return text
.split(/\n\n+/)
.map((part) => part.trim())
.filter((part) => part.length > 0)
.map((part) => Schema.decodeUnknownSync(SseEvent)(JSON.parse(part.replace(/^data: /, ""))))
}
/**
* Opens a scoped subscription to the instance `/event` SSE stream and returns
* a Queue of decoded events. The underlying request and decoder fiber are
* released when the test scope closes.
*/
export const openInstanceEventStream = (directory: string) =>
Effect.gen(function* () {
const response = yield* HttpClientRequest.get(EventPaths.event).pipe(
HttpClientRequest.setHeader("x-opencode-directory", directory),
HttpClient.execute,
)
const queue = yield* Queue.unbounded<SseEvent>()
yield* response.stream.pipe(
Stream.decodeText({ encoding: "utf-8" }),
Stream.flatMap((text) => Stream.fromIterable(decodeFrames(text))),
Stream.runForEach((event) => Queue.offer(queue, event)),
Effect.forkScoped,
)
return queue
})
export const readNextEvent = (queue: Queue.Queue<SseEvent>) =>
Queue.take(queue).pipe(
Effect.timeoutOrElse({
duration: "3 seconds",
orElse: () => Effect.fail(new Error("timed out reading SSE event")),
}),
)
export const collectUntilEvent = (queue: Queue.Queue<SseEvent>, predicate: (event: SseEvent) => boolean) =>
Effect.gen(function* () {
const events: SseEvent[] = []
while (true) {
const event = yield* readNextEvent(queue)
events.push(event)
if (predicate(event)) return events
}
}).pipe(
Effect.timeoutOrElse({
duration: "4 seconds",
orElse: () => Effect.fail(new Error("collectUntilEvent deadline exceeded")),
}),
)