Skip to content

Commit c39627c

Browse files
pjdurdennacx
authored andcommitted
fix(mcpproxy): skip non-response SSE events during initialize (envoyproxy#2219) (envoyproxy#2267)
**Description** Fixes envoyproxy#2219. The MCP proxy fails `initialize` with `500 failed to create MCP session to any backend` and logs `MCP message is not a response: <nil>` when a backend's `initialize` SSE response begins with a non-response event (an empty/keep-alive `data:` line) before the JSON-RPC result. Some FastMCP backends do this depending on the requested `protocolVersion` (e.g. firecrawl with `2025-11-25`). Two layers caused it: 1. `sseEventParser.parseEvent` (`internal/mcpproxy/sse.go`) called `jsonrpc.DecodeMessage` on an empty `data:` line, which errors, so the whole leading event was returned as a parse error. 2. The initialize reader loop (`internal/mcpproxy/mcpproxy.go`) treated that non-EOF error as fatal and broke with `rawMsg` still nil, then reported it as "not a response". This change: - Skips empty `data:` lines in the SSE parser (keep-alive/heartbeat events carry no message to decode). Non-empty but malformed data still errors as before. - Makes the initialize reader skip non-response events and keep reading until it finds the JSON-RPC `*Response`, instead of requiring the first event to be the response. **Testing** - `TestSSEEventParser_EmptyDataLineSkipped` — a leading event with an empty data line yields an event with no messages and no error; the following event still decodes. - `TestNewSession_SSEWithLeadingKeepAlive` — a backend whose initialize SSE response starts with an empty keep-alive event before the real response now initializes successfully (this test reproduced the `MCP message is not a response: <nil>` failure before the fix). ``` go test ./internal/mcpproxy/ go vet ./internal/mcpproxy/ go build ./... ``` Done with AI assistance; I have reviewed every line and understand the change. --------- Signed-off-by: pjdurden <prajjwalchittori1@gmail.com> Co-authored-by: Ignasi Barrera <ignasi@tetrate.io> Signed-off-by: Daniel Chernovsky <daniel.chernovsky@doubleverify.com>
1 parent 0017bdb commit c39627c

4 files changed

Lines changed: 118 additions & 9 deletions

File tree

internal/mcpproxy/mcpproxy.go

Lines changed: 25 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -352,26 +352,42 @@ func (m *mcpRequestContext) initializeSession(ctx context.Context, routeName fil
352352
}
353353
if rawMsg == nil {
354354
parser := newSSEEventParser(sseReader, backend.Name)
355-
for {
355+
var readErr error
356+
for rawMsg == nil {
356357
event, parseErr := parser.next()
357358
// TODO: handle reconnect. We need to re-arrange the event ID so that it will also contain the backend name and the original session ID.
358359
// Since event ID can be arbitrary string, we can shove each backend's last even ID into the event ID just like the session ID.
359360
if event != nil {
360-
// TODO: there's no session here what should we do?
361-
if len(event.messages) < 1 {
362-
return nil, errors.New("failed to get message from MCP sse event")
361+
// Some backends emit non-response events (keep-alives with an
362+
// empty data line, notifications) before the initialize result.
363+
// Skip those and keep reading until we find the JSON-RPC response.
364+
for _, msg := range event.messages {
365+
if _, ok := msg.(*jsonrpc.Response); ok {
366+
rawMsg = msg
367+
}
363368
}
364-
// Last event is the actual response.
365-
rawMsg = event.messages[len(event.messages)-1]
369+
}
370+
if rawMsg != nil {
371+
// Found the response; a trailing EOF on this same event is not a failure.
372+
break
366373
}
367374
if parseErr != nil {
368-
if errors.Is(parseErr, io.EOF) || strings.Contains(parseErr.Error(), "context deadline exceeded") {
369-
break
375+
readErr = parseErr
376+
if !errors.Is(parseErr, io.EOF) && !strings.Contains(parseErr.Error(), "context deadline exceeded") {
377+
m.l.Error("failed to read MCP GET response body", slog.String("error", parseErr.Error()))
370378
}
371-
m.l.Error("failed to read MCP GET response body", slog.String("error", parseErr.Error()))
372379
break
373380
}
374381
}
382+
if rawMsg == nil {
383+
// The SSE stream ended (EOF/deadline) or errored before any JSON-RPC
384+
// response arrived. Surface a clear error instead of falling through to
385+
// the misleading "MCP message is not a response: <nil>".
386+
if readErr != nil && !errors.Is(readErr, io.EOF) && !strings.Contains(readErr.Error(), "context deadline exceeded") {
387+
return nil, fmt.Errorf("failed to read MCP initialize response from backend %q: %w", backend.Name, readErr)
388+
}
389+
return nil, fmt.Errorf("MCP initialize stream from backend %q ended before a JSON-RPC response was received", backend.Name)
390+
}
375391
}
376392

377393
msg, ok := rawMsg.(*jsonrpc.Response)

internal/mcpproxy/mcpproxy_test.go

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -399,6 +399,33 @@ func TestInitializeSession_InitializeFailure(t *testing.T) {
399399
require.Contains(t, err.Error(), "failed with status code")
400400
}
401401

402+
func TestInitializeSession_SSEEndsBeforeResponse(t *testing.T) {
403+
// Backend returns a 200 text/event-stream whose initialize response contains
404+
// only non-response events (a keep-alive) and then closes before ever sending
405+
// the JSON-RPC response. This must produce a clear error rather than the
406+
// misleading "MCP message is not a response: <nil>".
407+
backendServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
408+
w.Header().Set(sessionIDHeader, "test-session-123")
409+
w.Header().Set("Content-Type", "text/event-stream")
410+
w.WriteHeader(http.StatusOK)
411+
_, _ = w.Write([]byte(`id: keepalive_0000
412+
data:
413+
414+
`))
415+
}))
416+
defer backendServer.Close()
417+
418+
proxy := newTestMCPProxy()
419+
proxy.backendListenerAddr = backendServer.URL
420+
421+
sessionID, err := proxy.initializeSession(t.Context(), "route1", filterapi.MCPBackend{Name: "test-backend"}, &mcp.InitializeParams{}, time.Now())
422+
423+
require.Error(t, err)
424+
require.Empty(t, sessionID)
425+
require.Contains(t, err.Error(), "ended before a JSON-RPC response was received")
426+
require.NotContains(t, err.Error(), "is not a response")
427+
}
428+
402429
func TestInitializeSession_NotificationsInitializedFailure(t *testing.T) {
403430
// Mock backend server.
404431
var callCount perBackendCallCount
@@ -471,3 +498,41 @@ func TestInvokeJSONRPCRequest_NoSessionID(t *testing.T) {
471498
require.Equal(t, http.StatusOK, resp.StatusCode)
472499
require.NoError(t, resp.Body.Close())
473500
}
501+
502+
// Issue #2219: when a backend's initialize SSE response starts with a non-response
503+
// event (an empty/keep-alive data line) before the real JSON-RPC response, session
504+
// creation must still succeed rather than failing with "MCP message is not a response".
505+
func TestNewSession_SSEWithLeadingKeepAlive(t *testing.T) {
506+
var callCount perBackendCallCount
507+
backendServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
508+
backend := r.Header.Get(internalapi.MCPBackendHeader)
509+
if callCount.inc(backend)%2 == 1 {
510+
// Odd calls: initialize requests. Emit a leading empty keep-alive
511+
// event, then the real response event.
512+
w.Header().Set(sessionIDHeader, "test-session-123")
513+
w.Header().Set("Content-Type", "text/event-stream")
514+
w.WriteHeader(http.StatusOK)
515+
_, _ = w.Write([]byte(`id: keepalive_0000
516+
data:
517+
518+
event: message
519+
id: msg_0001
520+
data: {"jsonrpc":"2.0","id":"ff3964c5-4c79-4567-96e2-29e905754e58","result":{"capabilities":{"logging":{},"tools":{"listChanged":true}},"protocolVersion":"2025-06-18","serverInfo":{"name":"dumb-echo-server","version":"0.1.0"}}}
521+
522+
`))
523+
} else {
524+
// Even calls: notifications/initialized requests.
525+
w.WriteHeader(http.StatusAccepted)
526+
}
527+
}))
528+
defer backendServer.Close()
529+
530+
proxy := newTestMCPProxy()
531+
proxy.backendListenerAddr = backendServer.URL
532+
533+
s, err := proxy.newSession(t.Context(), &mcp.InitializeParams{}, "test-route", "", nil, time.Now())
534+
535+
require.NoError(t, err)
536+
require.NotNil(t, s)
537+
require.NotEmpty(t, s.clientGatewaySessionID())
538+
}

internal/mcpproxy/sse.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -106,6 +106,11 @@ func (s *sseEventParser) parseEvent(chunk []byte) (*sseEvent, error) {
106106
ret.id = string(bytes.TrimSpace(line[len(sseIDPrefix):]))
107107
case bytes.HasPrefix(line, sseDataPrefix):
108108
data := bytes.TrimSpace(line[len(sseDataPrefix):])
109+
if len(data) == 0 {
110+
// Empty data line (e.g. a keep-alive/heartbeat event). There is
111+
// nothing to decode, so skip it rather than failing the event.
112+
continue
113+
}
109114
msg, err := jsonrpc.DecodeMessage(data)
110115
if err != nil {
111116
return nil, fmt.Errorf("failed to decode jsonrpc message from sse data: %w", err)

internal/mcpproxy/sse_test.go

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -366,3 +366,26 @@ func TestSSEEventParser_EndOfStream(t *testing.T) {
366366
})
367367
}
368368
}
369+
370+
// Issue #2219: a leading keep-alive event with an empty data line (emitted by
371+
// some FastMCP backends, e.g. firecrawl with protocolVersion 2025-11-25, before
372+
// the real response) must be skipped, not treated as a JSON decode error.
373+
func TestSSEEventParser_EmptyDataLineSkipped(t *testing.T) {
374+
raw := []byte("id: keepalive_0000\ndata:\n\n" +
375+
"event: message\nid: msg_0001\ndata: {\"jsonrpc\":\"2.0\",\"id\":\"1\",\"result\":{}}\n\n")
376+
p := newSSEEventParser(bytes.NewReader(raw), "mybackend")
377+
378+
// First event carries only an empty data line: no messages, no error.
379+
ev1, err := p.next()
380+
require.NoError(t, err)
381+
require.NotNil(t, ev1)
382+
require.Empty(t, ev1.messages)
383+
require.Equal(t, "keepalive_0000", ev1.id)
384+
385+
// Second event is the real JSON-RPC response.
386+
ev2, err := p.next()
387+
require.NoError(t, err)
388+
require.Len(t, ev2.messages, 1)
389+
_, ok := ev2.messages[0].(*jsonrpc.Response)
390+
require.True(t, ok)
391+
}

0 commit comments

Comments
 (0)