Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/release-disconnected-streams.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@kilocode/cli": patch
---

Release disconnected event streams so long-running servers do not retain queued session diffs.
51 changes: 51 additions & 0 deletions packages/opencode/src/kilocode/server/sse.ts
Original file line number Diff line number Diff line change
@@ -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<void>((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<void>((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)
})
})
}
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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<GlobalBusEvent>((queue) => {
const handler = (event: GlobalBusEvent) => Queue.offerUnsafe(queue, event)
Expand All @@ -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"))),
),
{
Expand All @@ -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* () {
Expand Down
68 changes: 68 additions & 0 deletions packages/opencode/test/kilocode/server/httpapi-global-sse.test.ts
Original file line number Diff line number Diff line change
@@ -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<unknown>)),
)
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,
)
})
Loading