-
Notifications
You must be signed in to change notification settings - Fork 10
sse streams
Verdict: two stream kinds. POST-SSE = per-request, lazy: JSON response unless handler emits server→client message first, then upgrades to text/event-stream and final response is last event. GET-SSE = session listening stream (stateful) or empty stream (stateless). Every event id is globally unique counter, suffixed #<streamKey> for POST streams → replay never crosses streams.
| Type | Role | Proof |
|---|---|---|
OutboundSseStream |
transport-neutral: start, started, writeEvent, comment, close, onClose, streamKey, channelId
|
OutboundSseStream |
PostSseStream |
Netty impl, state machine NEW/OPEN/CLOSED_UNOPENED/CLOSED_OPENED, all mutations on event loop |
PostSseStream |
NettySseConnection |
SseConnection for GET stream, close listener |
NettySseConnection |
SseManager |
open GET streams, priming, replay | SseManager |
SseHeartbeat |
:\r\n comment every interval on EL; skip when unwritable; failed tick stashes the cause then closes |
SseHeartbeat#enable, SseHeartbeat#send
|
SseSerializer |
id:/event:/data: framing into pooled buf, split on newlines |
SseSerializer |
OutboundSseStreamMessageRouter |
ThreadLocal (sessionId, stream) during dispatch → divert session notifications onto POST stream |
OutboundSseStreamMessageRouter#logger |
-
PostSseStreamcreated per POST;streamKey= one counter draw (not JSON-RPC id — clients reuse ids)PostSseStream#PostSseStream. That same draw is the priming event's id, so it precedes every id drawn later during dispatch — senders draw ids while the POST is still buffered (DefaultTachyonServer#sendSerializedNotification), so a freshly drawn priming id would outrank them and a resume from it would skip the notification. - Handler emits notification/progress/log/request →
start()→ writes 200 + SSE headers (Connection: close,X-Accel-Buffering: no) + enables heartbeat + priming eventid=<n>#<key>with empty data (SEP-1699) unless events queuedPostSseStream#doStart,HttpHelpers#setSseStreamHeaders.start()returns aCompletionStage<Void>that completes once that initial write (queued events included) is flushed — used bysubscriptions/listento time its ack, not just scheduling itOutboundSseStream#start,SubscriptionsListenHandler#handleAsync→ observability. A secondstart()mirrors the outcome of the call that opened the stream (comment()self-start included) instead of reporting success early;start()on a closed stream failsClosedChannelExceptionPostSseStream#doStart. -
ctx.notifications().comment(msg)self-starts stream → token-free keep-alivePostSseStream#comment,NotificationsImpl#comment. - Handler done, stream started ⇒ final response finalized on VT: append
ResponseEventto log (stateful), write,terminateAsync()(last chunk + close)McpOperationHandler#finalizePostSseResponse. - Stream never started ⇒
terminate()(neutralize so late message can't open a second response on pooled socket) then plain JSONMcpOperationHandler#completePostRequest. - Final write dropped (client gone) ⇒
redeliverOnReconnect: if session's current GET connection resumed this stream key, send live; else wait for replayMcpOperationHandler#redeliverOnReconnect.
Initial headers and every queued event are aggregated with Netty PromiseCombiner; a successful final write cannot hide an earlier failure PostSseStream#doStart. The close attribute preserves the first failure with setIfAbsent ChannelHandlerUtils#markCloseFailure.
- Stateless server:
openStatelessStream— headers,retry: 3000, heartbeat, priming, nothing elseSseManager#openStatelessStream. - Stateful: requires
MCP-Session-Id(400) and local or hydrated session (404)McpOperationHandler#handleGet. -
openStream: newNettySseConnectionreplaces session connection (previous one closed); close listener detaches only if still currentSseManager#openStream,Session#connection. -
Last-Event-IDpresent ⇒ rememberresumingStreamKey, replay on executor. - 2026-07-28 has no GET (protocol
matchesPOST only); usessubscriptions/listenPOST stream instead → feature-registries.
SseManager.replayEvents SseManager#replayEvents: parse <n>[#<key>], take all session events, keep sseId > n and streamKey == key (null = GET stream), convert via ServerEngine.toSseEvent (request/cancel events skipped) ServerEngine#wireEventId, stop when session.send false (throttled/closed).
NetworkConfig: readerIdleTimeout 60s closes silent non-SSE sockets; heartbeat 15s keeps SSE; keep heartbeat < reader idle and < session TTL (30s) NetworkConfig, NetworkConfig#DEFAULT_READER_IDLE_TIMEOUT. Idle tick on SSE channel = no-op McpOperationHandler#userEventTriggered.
-
close()writesretry: 3000then last chunk + close (client should reconnect) vsterminate()no retryPostSseStream#doClose,NettySseConnection#doClose. -
terminateAsyncalso completes on channel close so shutdown drain never hangs — but consults the recorded close cause first, because a failed terminating write closes the channel from its own listener and that fallback would otherwise report success for itPostSseStream#terminateAsync. -
doClosecancels heartbeats before the terminating chunk: the scheduled tick is otherwise cancelled only on channel close and could emit a comment the HTTP encoder no longer acceptsSseHeartbeat#cancel,NettySseConnection#doClose. - Fire-and-forget calls (
writeEvent,comment,close,terminate) swallow a shutting-down loop's rejection;startfails its stage andwriteEvent(long, byte[], Runnable)runsonDroppedinsteadPostSseStream#runOnEventLoopQuietly. - Every write (headers, priming/queued, events, comments, retry, last chunk) shares one failure listener: stash cause via
ChannelHandlerUtils#markCloseFailure, then close — soonClosereports a transport failure, not an ordinary disconnectPostSseStream#closeOnWriteFailure. A heartbeat is written by the scheduler, not that listener, and does the same for itself — it is how an idle stream finds a dead peerSseHeartbeat#send. - SSE responses always
Connection: close— server hard-closes on end; advertising keep-alive raced FIN vs next requestHttpHelpers#HttpHelpers.
Related: sessions, request-lifecycle, concurrency.
📄 source .llm-wiki/concepts/sse-streams.md · updated 2026-09-17 · verified at d831e9b1 · tags [concept, transport, sse]
🧭 Start
⚙️ Concepts (cross-cutting)
- request-lifecycle
- netty-pipeline
- protocol-versions
- sessions
- sse-streams
- feature-registries
- tasks
- extensions
- json-layer
- errors
- concurrency
- declarative-configuration
- configuration
- security-guards
- observability
- api-stability
📦 Modules