Skip to content

feat(server): add streamable HTTP stream resumability via pluggable EventStore - #930

Merged
ezynda3 merged 2 commits into
mark3labs:mainfrom
Abhinav2408:feat/streamable-http-resumability
Jul 23, 2026
Merged

ezynda3 merged 2 commits into
mark3labs:mainfrom
Abhinav2408:feat/streamable-http-resumability

Conversation

@Abhinav2408

@Abhinav2408 Abhinav2408 commented Jul 15, 2026 •

Copy link
Copy Markdown
Contributor

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 WithEventStore
option: 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 a
GET carrying Last-Event-ID: the server replays that stream's recorded events after the
given id, in order with their original ids, then continues live delivery on the same
connection.

Notable behavior:

  • Messages produced while no client is attached are still recorded — a tool call whose
    POST connection broke keeps recording its notifications and final response, all
    redelivered on resume.
  • A resumed connection supersedes any connection still attached to the stream.
  • Unknown or cross-session Last-Event-ID values get 400 Bad Request.
  • Store-before-write plus per-stream delivery locking makes the replay→live handoff
    exactly-once (no gaps, no duplicates).
  • With no store configured, behavior is unchanged (no id fields, no recording).

EventStore is a two-method interface (StoreEvent / ReplayEventsAfter) so events can
live in shared storage for multi-node deployments; NewInMemoryEventStore() ships as the
single-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 ./... -race and golangci-lint run are clean.

Type of Change

  • New feature (non-breaking change that adds functionality)
  • MCP spec compatibility implementation

Checklist

  • My code follows the code style of this project
  • I have performed a self-review of my own code
  • I have added tests that prove my fix is effective or that my feature works
  • I have updated the documentation accordingly

MCP Spec Compliance

  • This PR implements a feature defined in the MCP specification
  • Link to relevant spec section: Resumability and Redelivery
  • Implementation follows the specification exactly

Summary by CodeRabbit

  • New Features
    • Added optional HTTP/SSE stream resumability via configurable event storage.
    • SSE messages and final tool responses can be replayed after reconnect using Last-Event-ID (with the session ID) for gap-free continuation.
    • Server rejects unknown, cross-session, or previously terminated event IDs with 400 Bad Request.
    • Resuming one stream doesn’t replay messages from other streams; interrupted tool-call requests continue correctly.
  • Documentation
    • Added a “Stream Resumability” section to HTTP transport Session Management docs.
  • Tests
    • Added extensive end-to-end coverage for resumability scenarios.
  • Chores
    • Improved write behavior by adding response write deadlines when streaming.

…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>
@mark-iii-labs-huly

Copy link
Copy Markdown

Connected to Huly®: MCP_G-494

@coderabbitai

coderabbitai Bot commented Jul 15, 2026 •

Copy link
Copy Markdown
Contributor

Review Change Stack

Walkthrough

Adds opt-in SSE resumability for Streamable HTTP using an EventStore, replay via Last-Event-ID, session and stream lifecycle management, custom-store support, documentation, and comprehensive end-to-end tests.

Changes

SSE stream resumability

Layer / File(s) Summary
Event store contract and in-memory implementation
server/event_store.go
Defines event storage, ordered replay, unknown-ID handling, copied payloads, and session purging.
Resumable stream lifecycle and replay
server/streamable_http_resume.go, server/streamable_http_handle.go
Implements event delivery, replay, connection attachment, listening pumps, heartbeats, cleanup, write deadlines, and SSE event IDs.
HTTP transport wiring and lifecycle
server/streamable_http.go, www/docs/pages/transports/http.mdx
Adds WithEventStore, integrates resumable streams into POST and GET handling, preserves session state for replay, and documents configuration and protocol behavior.
End-to-end resumability validation
e2e/resumability_http_test.go
Tests event IDs, replay, interrupted requests, stream isolation, invalid IDs, custom stores, plain JSON responses, and reconnect supersession.

Estimated code review effort: 4 (Complex) | ~60 minutes

Possibly related PRs

Suggested labels: area: mcp spec

Suggested reviewers: ezynda3

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Docstring Coverage ✅ Passed No functions found in the changed files to evaluate docstring coverage. Skipping docstring coverage check.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
Title check ✅ Passed The title clearly states the main server-side resumability change and the pluggable EventStore.
Description check ✅ Passed The description follows the template, covering change summary, type, checklist, and MCP compliance.
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 5

🧹 Nitpick comments (4)
server/streamable_http_resume.go (2)

374-400: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick win

Heartbeats 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 via StoreEvent just 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 (skip StoreEvent) 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

eventID is written into the SSE id: line without sanitization.

The EventStore interface documents that "any non-empty value is acceptable" for event IDs, but a custom implementation returning an ID containing \n/\r would corrupt the SSE framing (or inject extra fields) here, since eventID is 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 win

Distinguish "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 with fmt.Errorf("context: %w", err), and use errors.Is/As for 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 win

Add godoc comments to the exported StoreEvent/ReplayEventsAfter methods.

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

📥 Commits

Reviewing files that changed from the base of the PR and between f6f3485 and 363a696.

📒 Files selected for processing (5)
  • e2e/resumability_http_test.go
  • server/event_store.go
  • server/streamable_http.go
  • server/streamable_http_resume.go
  • www/docs/pages/transports/http.mdx

Comment thread e2e/resumability_http_test.go Outdated
Comment thread e2e/resumability_http_test.go
Comment thread server/event_store.go
Comment thread server/streamable_http_resume.go
Comment thread server/streamable_http.go
- 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants