diff --git a/packages/opencode/src/bus/global.ts b/packages/opencode/src/bus/global.ts index 43386dd6b20..dc23020d27c 100644 --- a/packages/opencode/src/bus/global.ts +++ b/packages/opencode/src/bus/global.ts @@ -8,3 +8,4 @@ export const GlobalBus = new EventEmitter<{ }, ] }>() +GlobalBus.setMaxListeners(50) // kilocode_change — surface warning if SSE listeners accumulate diff --git a/packages/opencode/src/server/routes/event.ts b/packages/opencode/src/server/routes/event.ts index 989b857710d..565486d3c3a 100644 --- a/packages/opencode/src/server/routes/event.ts +++ b/packages/opencode/src/server/routes/event.ts @@ -70,14 +70,27 @@ export const EventRoutes = () => stream.onAbort(stop) + // kilocode_change start + // On Windows, stream.onAbort() may never fire after a client disconnects + // (delayed TCP RST detection via IOCP). Without this try/catch, the + // GlobalBus listener, heartbeat interval, and AsyncQueue stay alive + // indefinitely for each dead connection — leaking memory on every + // SSE reconnect. Catching write errors lets us clean up eagerly. try { for await (const data of q) { if (data === null) return - await stream.writeSSE({ data }) + try { + await stream.writeSSE({ data }) + } catch { + log.info("event write failed, cleaning up dead stream") + stop() + return + } } } finally { stop() } + // kilocode_change end }) }, ) diff --git a/packages/opencode/src/server/routes/global.ts b/packages/opencode/src/server/routes/global.ts index 733a0ef6b27..ddfa4c9acad 100644 --- a/packages/opencode/src/server/routes/global.ts +++ b/packages/opencode/src/server/routes/global.ts @@ -56,14 +56,27 @@ async function streamEvents(c: Context, subscribe: (q: AsyncQueue stream.onAbort(stop) + // kilocode_change start + // On Windows, stream.onAbort() may never fire after a client disconnects + // (delayed TCP RST detection via IOCP). Without this try/catch, the + // GlobalBus listener, heartbeat interval, and AsyncQueue stay alive + // indefinitely for each dead connection — leaking memory on every + // SSE reconnect. Catching write errors lets us clean up eagerly. try { for await (const data of q) { if (data === null) return - await stream.writeSSE({ data }) + try { + await stream.writeSSE({ data }) + } catch { + log.info("global event write failed, cleaning up dead stream") + stop() + return + } } } finally { stop() } + // kilocode_change end }) }