diff --git a/packages/opencode/src/server/routes/instance/httpapi/handlers/event.ts b/packages/opencode/src/server/routes/instance/httpapi/handlers/event.ts index e770a7cfba1a..c6353122fee7 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/handlers/event.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/handlers/event.ts @@ -13,11 +13,15 @@ function eventData(data: unknown): Sse.Event { return { _tag: "Event", event: "message", - id: undefined, + id: eventID(data), data: JSON.stringify(data), } } +function eventID(data: unknown) { + return data && typeof data === "object" && "id" in data && typeof data.id === "string" ? data.id : undefined +} + function eventResponse(bus: Bus.Interface) { return Effect.gen(function* () { // Subscribe eagerly: the bus subscription is acquired in the request scope diff --git a/packages/opencode/src/server/routes/instance/httpapi/handlers/global.ts b/packages/opencode/src/server/routes/instance/httpapi/handlers/global.ts index f80869b64d3f..6362175fb9f5 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/handlers/global.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/handlers/global.ts @@ -20,11 +20,19 @@ function eventData(data: unknown): Sse.Event { return { _tag: "Event", event: "message", - id: undefined, + id: eventID(data), data: JSON.stringify(data), } } +function eventID(data: unknown) { + if (!data || typeof data !== "object" || !("payload" in data)) return undefined + const payload = data.payload + return payload && typeof payload === "object" && "id" in payload && typeof payload.id === "string" + ? payload.id + : undefined +} + function parseBody(body: string) { try { return JSON.parse(body || "{}") as unknown diff --git a/packages/opencode/test/server/httpapi-event.test.ts b/packages/opencode/test/server/httpapi-event.test.ts index 44d421ea0a12..3f2ddb55c4f8 100644 --- a/packages/opencode/test/server/httpapi-event.test.ts +++ b/packages/opencode/test/server/httpapi-event.test.ts @@ -5,6 +5,7 @@ import { Bus } from "../../src/bus" import { Event as ServerEvent } from "../../src/server/event" import { Server } from "../../src/server/server" import { EventPaths } from "../../src/server/routes/instance/httpapi/groups/event" +import { GlobalPaths } from "../../src/server/routes/instance/httpapi/groups/global" import { resetDatabase } from "../fixture/db" import { disposeAllInstances, TestInstance } from "../fixture/fixture" import { testEffectShared } from "../lib/effect" @@ -13,11 +14,25 @@ void Log.init({ print: false }) const EventData = Schema.Struct({ id: Schema.optional(Schema.String), + sseID: Schema.optional(Schema.String), type: Schema.String, properties: Schema.Record(Schema.String, Schema.Any), }) -const readEvent = (reader: ReadableStreamDefaultReader) => +function parseSseEvent(text: string) { + const id = text + .split("\n") + .find((line) => line.startsWith("id: ")) + ?.slice("id: ".length) + const data = text + .split("\n") + .filter((line) => line.startsWith("data: ")) + .map((line) => line.slice("data: ".length)) + .join("\n") + return { sseID: id, ...JSON.parse(data) } +} + +const readSseEvent = (reader: ReadableStreamDefaultReader) => Effect.gen(function* () { const result = yield* Effect.promise(() => reader.read()).pipe( Effect.timeoutOrElse({ @@ -26,11 +41,12 @@ const readEvent = (reader: ReadableStreamDefaultReader) => }), ) if (result.done || !result.value) return yield* Effect.fail(new Error("event stream closed")) - return Schema.decodeUnknownSync(EventData)( - JSON.parse(new TextDecoder().decode(result.value).replace(/^data: /, "")), - ) + return parseSseEvent(new TextDecoder().decode(result.value)) }) +const readEvent = (reader: ReadableStreamDefaultReader) => + readSseEvent(reader).pipe(Effect.map(Schema.decodeUnknownSync(EventData))) + const openEventStream = (directory: string) => Effect.gen(function* () { const response = yield* Effect.promise(async () => @@ -42,6 +58,15 @@ const openEventStream = (directory: string) => return { response, reader } }) +const openGlobalEventStream = () => + Effect.gen(function* () { + const response = yield* Effect.promise(async () => Server.Default().app.request(GlobalPaths.event)) + if (!response.body) return yield* Effect.die("missing SSE response body") + const reader = response.body.getReader() + yield* Effect.addFinalizer(() => Effect.promise(() => reader.cancel().catch(() => undefined))) + return { response, reader } + }) + afterEach(async () => { await disposeAllInstances() await resetDatabase() @@ -62,7 +87,9 @@ describe("event HttpApi", () => { expect(response.headers.get("cache-control")).toBe("no-cache, no-transform") expect(response.headers.get("x-accel-buffering")).toBe("no") expect(response.headers.get("x-content-type-options")).toBe("nosniff") - expect(yield* readEvent(reader)).toMatchObject({ type: "server.connected", properties: {} }) + const event = yield* readEvent(reader) + expect(event).toMatchObject({ type: "server.connected", properties: {} }) + expect(event.sseID).toBe(event.id) }), { git: true, config: { formatter: false, lsp: false } }, ) @@ -94,8 +121,21 @@ describe("event HttpApi", () => { expect(yield* readEvent(reader)).toMatchObject({ type: "server.connected", properties: {} }) yield* Bus.use.publish(ServerEvent.Connected, {}) - expect(yield* readEvent(reader)).toMatchObject({ type: "server.connected", properties: {} }) + const event = yield* readEvent(reader) + expect(event).toMatchObject({ type: "server.connected", properties: {} }) + expect(event.sseID).toBe(event.id) }), { git: true, config: { formatter: false, lsp: false } }, ) + + it.effect("sets SSE ids on global events", () => + Effect.gen(function* () { + const { response, reader } = yield* openGlobalEventStream() + expect(response.status).toBe(200) + + const event = yield* readSseEvent(reader) + expect(event).toMatchObject({ payload: { type: "server.connected" } }) + expect(event.sseID).toBe(event.payload.id) + }), + ) })