Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 25 additions & 9 deletions internal/mcpproxy/mcpproxy.go
Original file line number Diff line number Diff line change
Expand Up @@ -352,26 +352,42 @@ func (m *mcpRequestContext) initializeSession(ctx context.Context, routeName fil
}
if rawMsg == nil {
parser := newSSEEventParser(sseReader, backend.Name)
for {
var readErr error
for rawMsg == nil {
event, parseErr := parser.next()
// 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.
// 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.
if event != nil {
// TODO: there's no session here what should we do?
if len(event.messages) < 1 {
return nil, errors.New("failed to get message from MCP sse event")
// Some backends emit non-response events (keep-alives with an
// empty data line, notifications) before the initialize result.
// Skip those and keep reading until we find the JSON-RPC response.
for _, msg := range event.messages {
if _, ok := msg.(*jsonrpc.Response); ok {
rawMsg = msg
}
}
// Last event is the actual response.
rawMsg = event.messages[len(event.messages)-1]
}
if rawMsg != nil {
// Found the response; a trailing EOF on this same event is not a failure.
break
}
if parseErr != nil {
if errors.Is(parseErr, io.EOF) || strings.Contains(parseErr.Error(), "context deadline exceeded") {
break
readErr = parseErr
if !errors.Is(parseErr, io.EOF) && !strings.Contains(parseErr.Error(), "context deadline exceeded") {
m.l.Error("failed to read MCP GET response body", slog.String("error", parseErr.Error()))
}
m.l.Error("failed to read MCP GET response body", slog.String("error", parseErr.Error()))
break
}
}
if rawMsg == nil {
// The SSE stream ended (EOF/deadline) or errored before any JSON-RPC
// response arrived. Surface a clear error instead of falling through to
// the misleading "MCP message is not a response: <nil>".
if readErr != nil && !errors.Is(readErr, io.EOF) && !strings.Contains(readErr.Error(), "context deadline exceeded") {
return nil, fmt.Errorf("failed to read MCP initialize response from backend %q: %w", backend.Name, readErr)
}
return nil, fmt.Errorf("MCP initialize stream from backend %q ended before a JSON-RPC response was received", backend.Name)
}
}

msg, ok := rawMsg.(*jsonrpc.Response)
Expand Down
65 changes: 65 additions & 0 deletions internal/mcpproxy/mcpproxy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -399,6 +399,33 @@ func TestInitializeSession_InitializeFailure(t *testing.T) {
require.Contains(t, err.Error(), "failed with status code")
}

func TestInitializeSession_SSEEndsBeforeResponse(t *testing.T) {
// Backend returns a 200 text/event-stream whose initialize response contains
// only non-response events (a keep-alive) and then closes before ever sending
// the JSON-RPC response. This must produce a clear error rather than the
// misleading "MCP message is not a response: <nil>".
backendServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.Header().Set(sessionIDHeader, "test-session-123")
w.Header().Set("Content-Type", "text/event-stream")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`id: keepalive_0000
data:

`))
}))
defer backendServer.Close()

proxy := newTestMCPProxy()
proxy.backendListenerAddr = backendServer.URL

sessionID, err := proxy.initializeSession(t.Context(), "route1", filterapi.MCPBackend{Name: "test-backend"}, &mcp.InitializeParams{}, time.Now())

require.Error(t, err)
require.Empty(t, sessionID)
require.Contains(t, err.Error(), "ended before a JSON-RPC response was received")
require.NotContains(t, err.Error(), "is not a response")
}

func TestInitializeSession_NotificationsInitializedFailure(t *testing.T) {
// Mock backend server.
var callCount perBackendCallCount
Expand Down Expand Up @@ -471,3 +498,41 @@ func TestInvokeJSONRPCRequest_NoSessionID(t *testing.T) {
require.Equal(t, http.StatusOK, resp.StatusCode)
require.NoError(t, resp.Body.Close())
}

// Issue #2219: when a backend's initialize SSE response starts with a non-response
// event (an empty/keep-alive data line) before the real JSON-RPC response, session
// creation must still succeed rather than failing with "MCP message is not a response".
func TestNewSession_SSEWithLeadingKeepAlive(t *testing.T) {
var callCount perBackendCallCount
backendServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
backend := r.Header.Get(internalapi.MCPBackendHeader)
if callCount.inc(backend)%2 == 1 {
// Odd calls: initialize requests. Emit a leading empty keep-alive
// event, then the real response event.
w.Header().Set(sessionIDHeader, "test-session-123")
w.Header().Set("Content-Type", "text/event-stream")
w.WriteHeader(http.StatusOK)
_, _ = w.Write([]byte(`id: keepalive_0000
data:

event: message
id: msg_0001
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"}}}

`))
} else {
// Even calls: notifications/initialized requests.
w.WriteHeader(http.StatusAccepted)
}
}))
defer backendServer.Close()

proxy := newTestMCPProxy()
proxy.backendListenerAddr = backendServer.URL

s, err := proxy.newSession(t.Context(), &mcp.InitializeParams{}, "test-route", "", nil, time.Now())

require.NoError(t, err)
require.NotNil(t, s)
require.NotEmpty(t, s.clientGatewaySessionID())
}
5 changes: 5 additions & 0 deletions internal/mcpproxy/sse.go
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,11 @@ func (s *sseEventParser) parseEvent(chunk []byte) (*sseEvent, error) {
ret.id = string(bytes.TrimSpace(line[len(sseIDPrefix):]))
case bytes.HasPrefix(line, sseDataPrefix):
data := bytes.TrimSpace(line[len(sseDataPrefix):])
if len(data) == 0 {
// Empty data line (e.g. a keep-alive/heartbeat event). There is
// nothing to decode, so skip it rather than failing the event.
continue
}
msg, err := jsonrpc.DecodeMessage(data)
if err != nil {
return nil, fmt.Errorf("failed to decode jsonrpc message from sse data: %w", err)
Expand Down
23 changes: 23 additions & 0 deletions internal/mcpproxy/sse_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -366,3 +366,26 @@ func TestSSEEventParser_EndOfStream(t *testing.T) {
})
}
}

// Issue #2219: a leading keep-alive event with an empty data line (emitted by
// some FastMCP backends, e.g. firecrawl with protocolVersion 2025-11-25, before
// the real response) must be skipped, not treated as a JSON decode error.
func TestSSEEventParser_EmptyDataLineSkipped(t *testing.T) {
raw := []byte("id: keepalive_0000\ndata:\n\n" +
"event: message\nid: msg_0001\ndata: {\"jsonrpc\":\"2.0\",\"id\":\"1\",\"result\":{}}\n\n")
p := newSSEEventParser(bytes.NewReader(raw), "mybackend")

// First event carries only an empty data line: no messages, no error.
ev1, err := p.next()
require.NoError(t, err)
require.NotNil(t, ev1)
require.Empty(t, ev1.messages)
require.Equal(t, "keepalive_0000", ev1.id)

// Second event is the real JSON-RPC response.
ev2, err := p.next()
require.NoError(t, err)
require.Len(t, ev2.messages, 1)
_, ok := ev2.messages[0].(*jsonrpc.Response)
require.True(t, ok)
}
Loading