mirror of
https://github.com/Kilo-Org/kilocode.git
synced 2026-09-24 16:02:55 +08:00
fix(cli): re-apply dead-stream detection to SSE endpoints to prevent memory leak on Windows (#8952)
* fix(cli): re-apply dead-stream detection to SSE endpoints to prevent memory leak on Windows The upstream OpenCode v1.3.0 merge rewrote SSE routes with AsyncQueue but dropped our dead-stream detection. On Windows, stream.onAbort() may never fire after client disconnect (IOCP delays TCP RST detection), leaking a GlobalBus listener, heartbeat interval, and AsyncQueue per dead connection. Wrap writeSSE in try/catch to clean up eagerly on write failure. * fix(cli): log dead-stream cleanup in SSE endpoints
This commit is contained in:
@@ -8,3 +8,4 @@ export const GlobalBus = new EventEmitter<{
|
||||
},
|
||||
]
|
||||
}>()
|
||||
GlobalBus.setMaxListeners(50) // kilocode_change — surface warning if SSE listeners accumulate
|
||||
|
||||
@@ -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
|
||||
})
|
||||
},
|
||||
)
|
||||
|
||||
@@ -56,14 +56,27 @@ async function streamEvents(c: Context, subscribe: (q: AsyncQueue<string | null>
|
||||
|
||||
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
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user