From 61bbc34eb261a27d2c56c8196a050929f8ef4e63 Mon Sep 17 00:00:00 2001 From: marius-kilocode Date: Wed, 24 Jun 2026 17:35:13 +0200 Subject: [PATCH] fix(cli): release disconnected event streams --- .changeset/release-disconnected-streams.md | 5 ++ packages/opencode/src/kilocode/server/sse.ts | 51 ++++++++++++++ .../instance/httpapi/handlers/global.ts | 13 +++- .../server/httpapi-global-sse.test.ts | 68 +++++++++++++++++++ 4 files changed, 135 insertions(+), 2 deletions(-) create mode 100644 .changeset/release-disconnected-streams.md create mode 100644 packages/opencode/src/kilocode/server/sse.ts create mode 100644 packages/opencode/test/kilocode/server/httpapi-global-sse.test.ts diff --git a/.changeset/release-disconnected-streams.md b/.changeset/release-disconnected-streams.md new file mode 100644 index 00000000000..651417c7aaf --- /dev/null +++ b/.changeset/release-disconnected-streams.md @@ -0,0 +1,5 @@ +--- +"@kilocode/cli": patch +--- + +Release disconnected event streams so long-running servers do not retain queued session diffs. diff --git a/packages/opencode/src/kilocode/server/sse.ts b/packages/opencode/src/kilocode/server/sse.ts new file mode 100644 index 00000000000..752d9db5bda --- /dev/null +++ b/packages/opencode/src/kilocode/server/sse.ts @@ -0,0 +1,51 @@ +import { Effect } from "effect" +import type { HttpServerRequest } from "effect/unstable/http" + +type Socket = { + readonly destroyed: boolean + once(event: "close", handler: () => void): unknown + off(event: "close", handler: () => void): unknown +} + +type Incoming = { + readonly aborted: boolean + readonly socket: Socket + once(event: "aborted", handler: () => void): unknown + off(event: "aborted", handler: () => void): unknown +} + +function isIncoming(source: object): source is Incoming { + return "aborted" in source && "socket" in source +} + +function aborted(signal: AbortSignal) { + return Effect.callback((resume) => { + if (signal.aborted) return resume(Effect.void) + const handler = () => resume(Effect.void) + signal.addEventListener("abort", handler, { once: true }) + return Effect.sync(() => signal.removeEventListener("abort", handler)) + }) +} + +// The stream scope owns the GlobalBus subscription, but some client disconnects do not +// interrupt that scope through the current Node HTTP adapter. Every orphaned stream then +// keeps an unbounded event queue alive. Since session.diff events can contain complete, +// multi-megabyte patches, extension reconnects can retain tens of gigabytes in those queues. +// Convert both Web and Node transport disconnects into an Effect that can interrupt the +// stream and run its acquireRelease finalizers immediately. +export function disconnect(request: HttpServerRequest.HttpServerRequest) { + const source = request.source + if (source instanceof Request) return aborted(source.signal) + if (!isIncoming(source)) return Effect.never + + return Effect.callback((resume) => { + if (source.aborted || source.socket.destroyed) return resume(Effect.void) + const handler = () => resume(Effect.void) + source.once("aborted", handler) + source.socket.once("close", handler) + return Effect.sync(() => { + source.off("aborted", handler) + source.socket.off("close", handler) + }) + }) +} 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 c9f202fd9bc..819766ac456 100644 --- a/packages/opencode/src/server/routes/instance/httpapi/handlers/global.ts +++ b/packages/opencode/src/server/routes/instance/httpapi/handlers/global.ts @@ -3,6 +3,7 @@ import { GlobalBus, type GlobalEvent as GlobalBusEvent } from "@/bus/global" import { EffectBridge } from "@/effect/bridge" import { Bus } from "@/bus" import { Installation } from "@/installation" +import { disconnect } from "@/kilocode/server/sse" // kilocode_change import { disposeAllInstancesAndEmitGlobalDisposed } from "@/server/global-lifecycle" import { InstallationVersion } from "@opencode-ai/core/installation/version" import * as Log from "@opencode-ai/core/util/log" @@ -33,7 +34,9 @@ function parseBody(body: string) { } } -function eventResponse() { +// kilocode_change start +function eventResponse(request: HttpServerRequest.HttpServerRequest) { + // kilocode_change end log.info("global event connected") const events = Stream.callback((queue) => { const handler = (event: GlobalBusEvent) => Queue.offerUnsafe(queue, event) @@ -53,6 +56,11 @@ function eventResponse() { Stream.map(eventData), Stream.pipeThroughChannel(Sse.encode()), Stream.encodeText, + // kilocode_change start - prevent disconnected SSE clients from retaining full diff payloads + // Explicit interruption closes the stream scope, unregisters its GlobalBus listener, and + // releases the unbounded callback queue even when transport cancellation is not propagated. + Stream.interruptWhen(disconnect(request)), + // kilocode_change end Stream.ensuring(Effect.sync(() => log.info("global event disconnected"))), ), { @@ -77,7 +85,8 @@ export const globalHandlers = HttpApiBuilder.group(RootHttpApi, "global", (handl }) const event = Effect.fn("GlobalHttpApi.event")(function* () { - return eventResponse() + const request = yield* HttpServerRequest.HttpServerRequest // kilocode_change + return eventResponse(request) // kilocode_change }) const configGet = Effect.fn("GlobalHttpApi.configGet")(function* () { diff --git a/packages/opencode/test/kilocode/server/httpapi-global-sse.test.ts b/packages/opencode/test/kilocode/server/httpapi-global-sse.test.ts new file mode 100644 index 00000000000..9f2d2f78fa8 --- /dev/null +++ b/packages/opencode/test/kilocode/server/httpapi-global-sse.test.ts @@ -0,0 +1,68 @@ +import { NodeHttpServer } from "@effect/platform-node" +import { describe, expect } from "bun:test" +import { Context, Effect, Fiber, Layer, Option, Stream } from "effect" +import { HttpClient, HttpRouter } from "effect/unstable/http" +import { HttpApiBuilder } from "effect/unstable/httpapi" +import { Auth } from "../../../src/auth" +import { GlobalBus } from "../../../src/bus/global" +import { Config } from "../../../src/config/config" +import { Installation } from "../../../src/installation" +import { ServerAuth } from "../../../src/server/auth" +import { RootHttpApi } from "../../../src/server/routes/instance/httpapi/api" +import { GlobalPaths } from "../../../src/server/routes/instance/httpapi/groups/global" +import { controlHandlers } from "../../../src/server/routes/instance/httpapi/handlers/control" +import { globalHandlers } from "../../../src/server/routes/instance/httpapi/handlers/global" +import { authorizationLayer } from "../../../src/server/routes/instance/httpapi/middleware/authorization" +import { schemaErrorLayer } from "../../../src/server/routes/instance/httpapi/middleware/schema-error" +import { pollWithTimeout, testEffect } from "../../lib/effect" + +const apiLayer = HttpRouter.serve( + HttpApiBuilder.layer(RootHttpApi).pipe( + Layer.provide([controlHandlers, globalHandlers]), + Layer.provide([authorizationLayer, schemaErrorLayer]), + ), + { disableListenLog: true, disableLogger: true }, +).pipe( + Layer.provideMerge(NodeHttpServer.layerTest), + Layer.provide(Layer.mock(Auth.Service)({})), + Layer.provide(Layer.mock(Config.Service)({})), + Layer.provide( + Layer.mock(Installation.Service)({ + method: () => Effect.succeed("npm"), + latest: () => Effect.succeed("9.9.9"), + upgrade: () => Effect.void, + }), + ), + Layer.provide(ServerAuth.Config.layer({ password: Option.none(), username: "opencode" })), + // Raw HttpApi routes expose an opaque handler context at the web boundary. + // oxlint-disable-next-line typescript-eslint/no-unsafe-type-assertion + Layer.provide(Layer.succeedContext(Context.empty() as Context.Context)), +) +const it = testEffect(apiLayer) + +describe("global SSE lifecycle", () => { + it.live( + "removes event listeners after clients disconnect", + () => + Effect.gen(function* () { + const count = GlobalBus.listenerCount("event") + + for (const attempt of [1, 2, 3]) { + const response = yield* HttpClient.get(GlobalPaths.event) + expect(response.status).toBe(200) + const fiber = yield* response.stream.pipe(Stream.runDrain, Effect.forkChild({ startImmediately: true })) + + yield* pollWithTimeout( + Effect.sync(() => (GlobalBus.listenerCount("event") === count + 1 ? true : undefined)), + `global event stream ${attempt} did not subscribe`, + ) + yield* Fiber.interrupt(fiber) + yield* pollWithTimeout( + Effect.sync(() => (GlobalBus.listenerCount("event") === count ? true : undefined)), + `global event stream ${attempt} did not remove its listener`, + ) + } + }), + 15_000, + ) +})