Repository navigation
fix(mcp): deliver server-push SSE on both profiles (GET establish + list_changed) - #1386
Conversation
There was a problem hiding this comment.
Code Review
This pull request implements robust Server-Sent Events (SSE) framing for the MCP server-push channel, resolving issues where headers were not flushed immediately and streams stalled under Fastify and Harper HTTP adapters. It introduces a new sse.ts utility to serialize event queues into primed Node Readable streams and adds comprehensive integration and unit tests. The review feedback highlights critical improvement opportunities to prevent memory leaks by properly cleaning up event listeners on stream closure, and to prevent potential process crashes by handling stream errors in the Fastify adapter.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
|
Reviewed; no blockers found. |
040e590 to
aafcfc7
Compare
…ist_changed) The MCP GET server-push channel did not work end-to-end on either profile, despite `tools/list_changed`/`resources/list_changed` being advertised and the `listChanged` dispatcher firing correctly: - Application profile (Harper HTTP server): the GET never returned — Node defers header transmission after writeHead until the first body byte, and the SSE queue yields nothing until a push, so the client hung with headers unsent. The raw IterableEventQueue objects were never serialized to SSE text either. - Operations profile (Fastify): the GET established headers but never streamed pushed frames — `reply.send(iterableEventQueue)` sends headers then stalls after one chunk, so a delivered `list_changed` sat in the queue forever. Root cause: both adapters handed the raw queue to the server and assumed it would SSE-frame + stream + flush it; neither does. The `listChanged` logic was correct (the queue just had no consumer). Fix — `components/mcp/sse.ts` `toSseStream`: a primed `PassThrough` of SSE wire text, driven by the queue's `'data'`/`'close'` events (NOT its async iterator — a generator awaiting the next frame can't be torn down on disconnect, since Node defers the generator's `.return()` until the pending await settles, leaking the socket + registry entry). Leading `:` comment flushes headers; frames serialized to `event:`/`data:`/`id:` (multi-line `data` split per spec); queue `'close'` ends the stream; stream `'close'` (disconnect) unsubscribes + signals the queue so the session registry drops the entry. Adapters: harperHttp returns the primed stream as the body; fastify `hijack()`s + writes/flushes headers on `reply.raw` + pipes it (Fastify won't drain the raw queue), merging Fastify-accumulated headers (e.g. `@fastify/cors`) onto the raw response. Hijack skips Fastify onSend/onResponse hooks for the long-lived stream (auth + handler already ran). Two SSE-lifecycle fixes this surfaced: - `server/http.ts`: don't re-wrap an existing Node stream in `Readable.from()`. It's redundant and broke destroy propagation — destroying the wrapper doesn't close the wrapped SSE PassThrough, so the application-profile session leaked on disconnect. A real stream already has `.pipe`, so it flows through unchanged. - `sessionRegistry.pruneIdleSessions`: also treat `queue.hasDataListeners` as a live consumer. The event-driven SSE path doesn't set `resolveNext`, so without this an open-but-idle stream (no recent `tools/call`) could be pruned. Tests: integration (`sse-listchanged.test.ts`) boots both profiles (app GET establishes; ops delivers `tools/list_changed` after a role flip); unit (`sse.test.ts`) covers framing/priming/multi-line/session-close/disconnect; `sessionRegistry.test.js` covers the event-driven live-consumer prune guard. Verified against the external MCP conformance suite (164 pass). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
aafcfc7 to
b5ca041
Compare
Address review findings on the GET-SSE framing helper: - toSseStream now removes its own queue 'close' listener (not just the 'data' listener) when the stream tears down before the queue closes, so a long-lived IterableEventQueue can't accumulate dead listeners per reconnect. SseFrameSource declares off/removeListener for 'close'. - A default 'error' guard is attached to the PassThrough so an error can never go unhandled and crash the worker — covers both the operations pipe and Harper's own HTTP server pipe (application profile), neither of which forwards a source error to a handler. - The operations adapter also destroys its raw socket on stream 'error' to tear the response down cleanly. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…down tests Codex flagged the 30ms fixed sleeps in the SSE teardown/listener-cleanup tests as a CI-flake risk: on a loaded runner the close propagation can be delayed past the window. Wait on the queue's 'close' event instead. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
The `await res.text()` embedded in the strictEqual assertion message was evaluated eagerly (template-literal args run before the call), consuming the response body even on a 200. The following `res.json()` then threw "Body is unusable: Body has already been read", failing N1/N2 on every runtime. Read the body once into a local and JSON.parse it. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
On Bun the operations API is served by delegating each request to Fastify via `fastify.inject()`, which resolves only when the response *ends* and returns a buffered payload. A long-lived SSE response (the MCP server-push GET, #619) never ends, so `await inject()` hung forever and the client never received headers — N2 hung only on Bun while passing on every Node runtime and on the application profile (which uses Harper's own HTTP server, not the inject bridge). Pass `payloadAsStream: true` so inject() resolves as soon as the hijacked reply writes headers and exposes the body as a Readable. Event-stream responses are returned as a streaming Response (frames reach the client as produced); all other responses keep the prior buffered behavior (drained to a single payload so Content-Length stays set). Bun-only path — the Node operations server listens on a real socket and never hits this bridge. Known limitation: a Bun client disconnect does not propagate back through the inject bridge to tear down the hijacked reply, so a dropped SSE session is reclaimed by the session-registry idle-prune backstop rather than immediately as on the Node/application path. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
The previous commit streamed operations SSE through the Bun inject() bridge but left disconnects unpropagated, so a dropped SSE session leaked: the session's queue 'data' listener stayed attached, and the registry idle-prune guard (which skips sessions with a live data listener) then never reclaimed it. When Bun cancels the response body it destroys the inject response stream; hook its 'close' to destroy the inject response object — the same object the SSE adapter listens on — which runs the adapter teardown, unsubscribes the queue 'data' listener, and drops the registry entry. Dropped Bun SSE sessions are now cleaned up immediately, same as the Node/application path. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
On server-side teardown (DELETE / superseding GET / idle-prune) the session
queue emits 'close'; toSseStream's onClose ends the stream, whose own 'close'
handler then emitted 'close' on the same queue a second time — duplicate
teardown for any on('close') listener. Track whether the source already
closed and skip the re-emit in that case; the emit now happens only for
client-side stream destruction (disconnect), which is what the session
registry relies on to drop the entry. Adds a unit test asserting the queue's
'close' fires exactly once on server-side teardown.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
|
Warning You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again! |
✅ Verified against the external MCP conformance suiteRan the external Streamable-HTTP conformance/edge suite ( Both server-push findings this PR targets are now confirmed fixed — the suite's
Two corroborating signals it's a real fix, not a calibration artifact:
Everything else stayed green: both-profile parity, CRUD breadth, rate limiting (per-tool + session scope), TTL eviction, auth/RBAC + session-binding, schema/Ajv validation, adversarial resource URIs. The single skip is the app-profile-500 regression guard, which correctly self-skips now that the profile works. Out of scope for this PR and still open (low priority): no 🤖 Generated with Claude Code |
kriszyp
left a comment
There was a problem hiding this comment.
Approving. The failing CI here is unrelated to this change — the only red test is a blob.test.js case that lives on main but isn't in this PR (CI runs the merge commit, and this PR touches no blob code; its own MCP/SSE unit tests pass on every Node and Bun shard). The fix itself is clean: toSseStream is a primed, event-driven Readable (correctly not an async generator, so it tears down on disconnect), both profiles are exercised (GET-establish timeout guard + tools/list_changed arrival after the super_user flip), and the stream error guards keep an SSE error from crashing the process. Nice work. 🙏
Two optional follow-ups, no blockers: consider a raw.on('error') on the fastify hijack path for airtight crash-safety, and note hasDataListeners is sticky on the queue (latent sharp edge, not exploitable on current paths).
— Reviewed by Claude (claude-opus-4-8)
…sDataListeners
Addresses the two optional follow-ups from the approving review:
- fastify hijack path: add raw.on('error', () => stream.destroy()) so a socket
error (e.g. the client resetting the connection mid-write) can't surface as an
unhandled 'error' on the raw response — symmetric to the existing
stream.on('error', () => raw.destroy()).
- IterableEventQueue.hasDataListeners is no longer sticky: it was stuck true for
the life of the queue once any 'data' listener had attached. Recompute it from
the live data-listener count when a listener is removed (override removeListener
AND off, since EventEmitter.off aliases the base removeListener and would bypass
the override — and SSE teardown unsubscribes via off). Net effect is also more
correct: send() with no live listener now buffers instead of emitting into the
void. New unit tests cover attach/drain, removal-clears-flag (via off and
removeListener), multi-listener, and non-data-listener isolation.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
|
Addressed both optional follow-ups from the review in
MCP unit suite green; the remaining red CI is the unrelated 🤖 Generated with AI (Claude Opus 4.8); reviewed by Kyle. |
kriszyp
left a comment
There was a problem hiding this comment.
Re-approving at 001e9da. The only change since my approval at 62ddec5 is crash-safety hardening — the raw error guard on the fastify hijack path and the non-sticky hasDataListeners I'd suggested, plus tests. No new risk. The red CI is an unrelated northwind JOIN-permissions flake (the SSE tests pass in that shard); re-running it and merging on green.
— Reviewed by Claude (claude-opus-4-8)
Fixes the native MCP server's GET server-push (SSE) channel, which never delivered server→client push on either profile despite advertising
tools/list_changed/resources/list_changedand thelistChangeddispatcher firing. The first two bugs were found by an external conformance suite and confirmed onmainby instrumentation; a Bun-only operations-profile hang was then found and fixed via CI.Bugs
writeHead) until the first body byte; the SSE queue yields nothing until a push, so the GET hung with headers unsent until the client gave up. (The raw queue objects were never serialized to SSE text either.)list_changednever delivered. The GET established headers butreply.send(IterableEventQueue)sends headers then stalls after one chunk, so a dispatchedtools/list_changedsat in the queue undelivered. Instrumentation confirmed the dispatcher fired and calledqueue.send(), but at send time the queue had no consumer (resolveNext=false).fastify.inject(), which resolves only when the response ends and buffers the payload. A long-lived SSE response never ends, soawait inject()hung forever and headers never reached the client. Bun-only — the application profile uses Harper's own HTTP server and was unaffected (it establishes in ~97ms on Bun).Root cause (first two): the adapters handed the raw
IterableEventQueueto their server assuming it would SSE-frame + stream + flush it; neither does. ThelistChangedlogic itself was correct.Fix
New
components/mcp/sse.tstoSseStream(queue)— wraps the queue in a primed, event-driven NodeReadableof SSE wire text::comment makes the body non-empty so headers flush immediately (fixes the application hang);event:/data:/id:(multi-linedatasplit per the SSE spec);'data'/'close'(not an async generator — a generator suspended atawait next()can't be torn down on disconnect). On stream close (client disconnect) it unsubscribes and signals the queue so the session registry drops the entry; a server-side close (DELETE / superseding GET / idle-prune) ends the stream without re-emitting the queue's'close'; a default'error'guard keeps an unhandled stream error from crashing the worker.Adapters / servers:
server/http.tsno longer re-wraps an existing Node stream inReadable.from()— that wrapper swallowed.destroy(), leaking the socket/session on disconnect. Shared HTTP path; output behavior unchanged, only teardown improved.reply.hijack()+ write/flush headers onreply.raw+ pipe the framed Readable, with an'error'handler that destroys the raw socket.server/http.tsbunDelegateToNodeServer):inject()now usespayloadAsStreamso it resolves atwriteHeadand exposes the body as a stream.text/event-streamresponses stream to the client; everything else is drained to a single buffered payload (Content-Length preserved — equivalent to before). Client disconnect is propagated back to the hijacked reply so the session is cleaned up immediately. Bun-only.hasDataListeners), so a connected-but-idle stream isn't reaped.Where to look / things to weigh
fastify.tshijack path — the main correctness surface. It merges Fastify-accumulated headers (reply.getHeaders(), e.g.@fastify/cors) onto the raw response so a CORS-allowed browser isn't blocked. Tradeoff (deliberate):hijack()skips Fastify'sonSend/onResponsehooks for this long-lived request, so per-request completion logging/metrics don't fire until the stream closes. Auth + the MCP handler already ran; this is the standard Fastify SSE pattern. Flagging for a reviewer who knows the operations-server telemetry expectations.server/http.ts— two changes on shared, non-MCP paths. (1) The Node body handler no longer re-wraps an existingReadable(affects every streaming HTTP response; output identical, only.destroy()propagation fixed). (2) The Bun inject bridge now reads every operations response as a stream viapayloadAsStream(non-SSE responses drained to a buffer, equivalent to the prior behavior). Both are the areas most worth a careful look.Tests
integrationTests/mcp/sse-listchanged.test.ts: boots both profiles; asserts the application GET establishes (200,text/event-stream, no hang) and the operations profile deliverstools/list_changedafter a role flip. Green on all runtimes in CI, including all six Bun shards (the Bun fix is validated end-to-end there).unitTests/components/mcp/sse.test.ts: framing, priming, multi-linedata, event-driven teardown, listener cleanup on disconnect, the'error'guard, and no-duplicate-'close'on server-side teardown.it.failstests for these bugs were promoted to passing positiveits.Reviews
Codex + Gemini cross-model review across several passes. Findings, all fixed with tests: multi-line-
datarobustness, CORS-header drop, stream-close socket leak, flaky fixed-sleep tests, the Bun disconnect session-leak (added disconnect propagation), and the duplicate'close'emit. One Codex concern — application SSE hanging on Bun — was adjudicated a false positive: N1 establishes in ~97ms on Bun, which is incompatible with the buffering path.Not in scope (deliberately deferred)
Two minor conformance gaps the external suite documents are left for a follow-up — they're orthogonal to the SSE channel: no
WWW-Authenticateheader on a 401 handshake (security/auth.ts), and lenient JSON-RPCidshapes ({}/[]accepted;jsonrpc.ts).Related: corrects #1349 §1 (the "GET-SSE push Implemented" claim).
🤖 Generated by an LLM (Claude Opus 4.8).