Repository navigation
feat(server): add streamable HTTP stream resumability via pluggable EventStore - #930
Conversation
…ventStore Implements the Resumability and Redelivery section of the MCP Streamable HTTP transport spec on the server side. Opt-in via WithEventStore: every JSON-RPC message delivered on a session's SSE streams (POST-upgraded response streams and the standalone GET listening stream) is recorded in the store before it is written and carries the store-issued SSE id. Clients resume with GET + Last-Event-ID: recorded events after that id are replayed on the same stream in order, then delivery continues live. - EventStore interface (StoreEvent / ReplayEventsAfter) with an in-process reference implementation, NewInMemoryEventStore - Messages produced while no client is attached are still recorded, so a tool call whose POST connection broke keeps recording notifications and its final response for redelivery on resume - A resumed connection supersedes any connection still attached to the stream; unknown or cross-session ids get 400 Bad Request - Store-before-write ordering plus per-stream delivery locking makes replay/live handoff exactly-once - No behavior change when no store is configured Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
|
Connected to Huly®: MCP_G-494 |
WalkthroughAdds opt-in SSE resumability for Streamable HTTP using an ChangesSSE stream resumability
Estimated code review effort: 4 (Complex) | ~60 minutes Possibly related PRs
Suggested labels: Suggested reviewers: 🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Actionable comments posted: 5
🧹 Nitpick comments (4)
server/streamable_http_resume.go (2)
374-400: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick winHeartbeats are recorded into the EventStore forever, even though replaying stale pings on resume is pointless.
Every periodic ping goes through
st.deliver, which persists it viaStoreEventjust like real application messages. For long-lived idle connections this permanently grows storage with data nobody needs to replay. Consider a "live-only" delivery path (skipStoreEvent) for heartbeat pings.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@server/streamable_http_resume.go` around lines 374 - 400, Update startConnHeartbeat so periodic ping messages use a live-only delivery path that sends them to the active stream without calling StoreEvent. Preserve the existing heartbeat interval, request construction, and cancellation behavior while ensuring heartbeat pings are never persisted for resume replay.
421-433: 🔒 Security & Privacy | 🔵 Trivial | ⚡ Quick win
eventIDis written into the SSEid:line without sanitization.The
EventStoreinterface documents that "any non-empty value is acceptable" for event IDs, but a custom implementation returning an ID containing\n/\rwould corrupt the SSE framing (or inject extra fields) here, sinceeventIDis interpolated raw into"id: %s\n".🛡️ Proposed guard
func writeSSEEventRaw(w io.Writer, eventID string, data json.RawMessage) error { if eventID != "" { + if strings.ContainsAny(eventID, "\r\n") { + return fmt.Errorf("invalid SSE event id: contains control characters") + } if _, err := fmt.Fprintf(w, "id: %s\n", eventID); err != nil {🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@server/streamable_http_resume.go` around lines 421 - 433, Update writeSSEEventRaw to validate eventID before writing the id field, rejecting or safely handling values containing carriage-return or newline characters so they cannot alter SSE framing. Preserve the existing behavior for empty and otherwise valid event IDs, including the current error propagation.server/event_store.go (2)
19-30: 🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick winDistinguish "unknown ID" from other store failures with a sentinel error.
ReplayEventsAfter's single error return conflates "lastEventID not found" with any underlying store failure.handleResumeGet(server/streamable_http_resume.go) treats every error here as a 400 client error, which will misreport genuine backend failures (e.g. a Redis-backed store timing out) as "Unknown Last-Event-ID" instead of a 5xx.As per coding guidelines, "Return sentinel errors (e.g.,
ErrMethodNotFound), wrap withfmt.Errorf("context: %w", err), and useerrors.Is/Asfor checking."♻️ Proposed sentinel error
+// ErrUnknownEventID indicates lastEventID is not known for the given session. +var ErrUnknownEventID = errors.New("unknown event ID") type EventStore interface { ... - // ... It returns an error if lastEventID is not known - // for the given session, or if send returns an error. + // ... It returns an error wrapping ErrUnknownEventID if lastEventID is + // not known for the given session, or the error from send if it returns one. ReplayEventsAfter(...) (string, error) }🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@server/event_store.go` around lines 19 - 30, Introduce a package-level sentinel error for an unknown last event ID and make ReplayEventsAfter return it, wrapped with context where appropriate, only when the requested ID is absent. Update handleResumeGet to use errors.Is against that sentinel so only unknown IDs produce a 400 response; propagate other ReplayEventsAfter failures as server errors.Source: Coding guidelines
65-116: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winAdd godoc comments to the exported
StoreEvent/ReplayEventsAftermethods.Neither method has its own doc comment (only the interface documents them).
As per coding guidelines, "All exported types and functions MUST have godoc comments starting with the name."
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@server/event_store.go` around lines 65 - 116, Add GoDoc comments immediately before the exported InMemoryEventStore methods StoreEvent and ReplayEventsAfter, with each comment starting with the corresponding method name and briefly describing its behavior.Source: Coding guidelines
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@e2e/resumability_http_test.go`:
- Around line 407-411: Ensure the resumability tests wait for server-observed
detachment after cancel1() and cancelPost() before sending disconnected events
n3 and p2. Add and use a deterministic lifecycle hook or synchronization barrier
tied to connection teardown, avoiding time-based sleeps, while preserving the
existing event assertions.
- Around line 306-326: Update startToolCall and waitResponse to propagate
request errors alongside the HTTP response through the channel. Have
waitResponse fail immediately with the original error when
http.DefaultClient.Do(req) fails, while preserving response-body cleanup and
timeout handling for successful requests.
In `@server/event_store.go`:
- Around line 19-30: The EventStore lifecycle lacks session-history eviction,
allowing terminated sessions to retain messages and resume by stale
Last-Event-ID. Add a session purge method to EventStore, implement it in
InMemoryEventStore, and invoke it from the teardown paths
cleanupSessionState/cleanupResumableState in server/streamable_http_resume.go
and server/streamable_http.go, ensuring handleResumeGet rejects terminated
sessions before replay. Update server/event_store.go (lines 19-30),
server/streamable_http_resume.go (lines 402-419 and 321-372), and
server/streamable_http.go (lines 104-122 and 1227-1229) as applicable, with no
direct change required where the root teardown fix already covers the site.
In `@server/streamable_http_resume.go`:
- Around line 83-117: Update resumableStream.deliver and replayAndAttach so
blocking sink sends and replay writes can exit on ctx cancellation or a bounded
timeout, rather than holding st.mu indefinitely. Apply HTTP write deadlines
through http.ResponseController on the underlying writer where available, and
ensure cancellation releases the stream lock and lets startListeningPump
continue or terminate cleanly. Preserve event delivery and cleanup behavior for
successful writes while preventing stalled clients from blocking
cleanupResumableState.
In `@server/streamable_http.go`:
- Around line 650-666: Extract the duplicated resumable-stream delivery sequence
from handlePost into a deliverToResumableStream helper. Have it lazily
initialize rst, handle ctx cancellation and SSE header upgrades, then call
deliver with the supplied message and last flag; update all four call sites,
including the mid-flight, drain, and final-response paths, to use it while
preserving the existing lock and behavior.
---
Nitpick comments:
In `@server/event_store.go`:
- Around line 19-30: Introduce a package-level sentinel error for an unknown
last event ID and make ReplayEventsAfter return it, wrapped with context where
appropriate, only when the requested ID is absent. Update handleResumeGet to use
errors.Is against that sentinel so only unknown IDs produce a 400 response;
propagate other ReplayEventsAfter failures as server errors.
- Around line 65-116: Add GoDoc comments immediately before the exported
InMemoryEventStore methods StoreEvent and ReplayEventsAfter, with each comment
starting with the corresponding method name and briefly describing its behavior.
In `@server/streamable_http_resume.go`:
- Around line 374-400: Update startConnHeartbeat so periodic ping messages use a
live-only delivery path that sends them to the active stream without calling
StoreEvent. Preserve the existing heartbeat interval, request construction, and
cancellation behavior while ensuring heartbeat pings are never persisted for
resume replay.
- Around line 421-433: Update writeSSEEventRaw to validate eventID before
writing the id field, rejecting or safely handling values containing
carriage-return or newline characters so they cannot alter SSE framing. Preserve
the existing behavior for empty and otherwise valid event IDs, including the
current error propagation.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Repository UI
Review profile: CHILL
Plan: Pro
Run ID: f42f2547-2108-4295-ba14-f57c487e1ca2
📒 Files selected for processing (5)
e2e/resumability_http_test.goserver/event_store.goserver/streamable_http.goserver/streamable_http_resume.gowww/docs/pages/transports/http.mdx
- EventStore gains PurgeSession; session teardown (DELETE or idle sweep) now purges the session's events so stale Last-Event-IDs cannot resume a terminated session - Add ErrUnknownEventID sentinel; unknown IDs keep returning 400 while other store failures now surface as 500 - Write to attached connections outside the stream lock and bound resume replays with a write deadline, so a stalled client cannot block stream takeover, the session pump, or session cleanup - Deliver heartbeat pings live-only with no event ID instead of recording them for replay - Reject event IDs containing line breaks before they reach SSE framing - Extract the duplicated resumable delivery sequence in handlePost into one helper - Propagate request errors in the e2e POST helper; add godoc on InMemoryEventStore methods; cover session termination purge in e2e Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Description
Implements server-side stream resumability for the streamable HTTP transport, which the
server's docs currently list as unsupported. Recording is opt-in via a new
WithEventStoreoption: every JSON-RPC message delivered on a session's SSE streams (POST responses that
upgraded to SSE, and the standalone GET listening stream) is recorded before it is written
and carries the store-issued SSE
id. A client that loses its connection resumes with aGET carrying
Last-Event-ID: the server replays that stream's recorded events after thegiven id, in order with their original ids, then continues live delivery on the same
connection.
Notable behavior:
POST connection broke keeps recording its notifications and final response, all
redelivered on resume.
Last-Event-IDvalues get400 Bad Request.exactly-once (no gaps, no duplicates).
idfields, no recording).EventStoreis a two-method interface (StoreEvent/ReplayEventsAfter) so events canlive in shared storage for multi-node deployments;
NewInMemoryEventStore()ships as thesingle-process reference implementation. Event ids are opaque to the transport; retention
policy is the store's concern (matching the TypeScript SDK's EventStore design).
Scope: server side only. Client-side auto-resume (prior attempt in #380) would layer on
top of this and is left for a follow-up.
Tests: 13 e2e tests covering replay ordering and boundaries, capture-while-disconnected
(including the interrupted-POST response), per-stream isolation under concurrent tool
calls, listening-stream resume, supersede semantics, custom-store id opacity, invalid-id
handling, and no-store baseline.
go test ./... -raceandgolangci-lint runare clean.Type of Change
Checklist
MCP Spec Compliance
Summary by CodeRabbit
Last-Event-ID(with the session ID) for gap-free continuation.400 Bad Request.