Module: MCPClient::ServerSSE::SseParser
- Includes:
- OriginPolicy
- Included in:
- MCPClient::ServerSSE
- Defined in:
- lib/mcp_client/server_sse/sse_parser.rb
Overview
Wire-level SSE parsing & dispatch ===
Instance Method Summary collapse
-
#fail_endpoint_handshake!(message) ⇒ Object
Record the handshake failure cause and raise.
-
#handle_endpoint_event(data) ⇒ Object
Handle the special "endpoint" control frame (for SSE handshake).
-
#handle_message_event(event, arrived = nil) ⇒ Object
Handle a "message" SSE event (payload is JSON-RPC over SSE).
-
#parse_and_handle_sse_event(event_data, arrived = nil) ⇒ Object
Parse and handle a raw SSE event payload.
-
#parse_sse_event(event_data) ⇒ Hash?
Parse a raw SSE chunk into its :event, :data, :id fields.
-
#process_error_in_message?(data) ⇒ Boolean
Process a connection-level JSON-RPC error payload in the SSE stream.
-
#process_notification?(data) ⇒ Boolean
Process a JSON-RPC notification (no id => notification).
-
#process_response?(data, arrived = nil) ⇒ Boolean
Process a JSON-RPC response (id => response).
-
#process_server_request?(data) ⇒ Boolean
Process a JSON-RPC request from server (has both id AND method).
-
#resolve_endpoint_uri(data) ⇒ String
Resolve an endpoint URI reference against the SSE connection URL.
Methods included from OriginPolicy
#origin_of, #reject_cross_origin_redirect!, #same_origin?
Instance Method Details
#fail_endpoint_handshake!(message) ⇒ Object
Record the handshake failure cause and raise. The SSE worker thread swallows this exception with a generic rescue, so also record the failure (mirroring @auth_error) for the connect caller blocked in wait_for_connection to surface promptly.
236 237 238 239 240 241 242 243 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 236 def fail_endpoint_handshake!() @mutex.synchronize do @connection_error = @connection_established = false @connection_cv.broadcast end raise MCPClient::Errors::TransportError, end |
#handle_endpoint_event(data) ⇒ Object
Handle the special "endpoint" control frame (for SSE handshake).
The event data is a URI reference (MCP 2024-11-05 HTTP with SSE: the
server sends "an endpoint event containing a URI for the client to
use for sending messages") which must be resolved against the SSE
connection URL per RFC 3986 section 5.1.3, so relative endpoint URIs
POST to the URL the server actually designated.
194 195 196 197 198 199 200 201 202 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 194 def handle_endpoint_event(data) endpoint = resolve_endpoint_uri(data) @mutex.synchronize do @rpc_endpoint = endpoint @sse_connected = true @connection_established = true @connection_cv.broadcast end end |
#handle_message_event(event, arrived = nil) ⇒ Object
Handle a "message" SSE event (payload is JSON-RPC over SSE)
33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 33 def (event, arrived = nil) return if event[:data].empty? begin # Dated from the arrival of the chunk it came in, before any event # of that chunk was decoded or dispatched. arrived ||= monotonic_now if respond_to?(:monotonic_now, true) data = JSON.parse(event[:data]) return if (data) return if process_server_request?(data) return if process_notification?(data) process_response?(data, arrived) rescue MCPClient::Errors::ConnectionError raise rescue JSON::ParserError => e @logger.warn("Failed to parse JSON from event data: #{describe_parse_error(e, event[:data])}") rescue StandardError => e @logger.error("Error processing SSE event: #{e.}") end end |
#parse_and_handle_sse_event(event_data, arrived = nil) ⇒ Object
Parse and handle a raw SSE event payload.
16 17 18 19 20 21 22 23 24 25 26 27 28 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 16 def parse_and_handle_sse_event(event_data, arrived = nil) event = parse_sse_event(event_data) return if event.nil? case event[:event] when 'endpoint' handle_endpoint_event(event[:data]) when 'ping' # no-op when 'message' (event, arrived) end end |
#parse_sse_event(event_data) ⇒ Hash?
Parse a raw SSE chunk into its :event, :data, :id fields
163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 163 def parse_sse_event(event_data) event = { event: 'message', data: '', id: nil } data_lines = [] has_content = false event_data.each_line do |line| line = line.chomp next if line.empty? # blank line next if line.start_with?(':') # SSE comment has_content = true if line.start_with?('event:') event[:event] = line[6..].strip elsif line.start_with?('data:') data_lines << line[5..].strip elsif line.start_with?('id:') event[:id] = line[3..].strip end end event[:data] = data_lines.join("\n") has_content ? event : nil end |
#process_error_in_message?(data) ⇒ Boolean
Process a connection-level JSON-RPC error payload in the SSE stream. Error RESPONSES (id-bearing) belong to a pending request and are delivered to the waiting caller via process_response? instead, per the MCP lifecycle "Error Handling" section (implementations SHOULD handle error cases such as protocol version mismatch), so they must not be swallowed here.
64 65 66 67 68 69 70 71 72 73 74 75 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 64 def (data) return false unless data['error'] return false if data['id'] = data['error']['message'] || 'Unknown server error' error_code = data['error']['code'] () if (, error_code) @logger.error("Server error: #{}") true end |
#process_notification?(data) ⇒ Boolean
Process a JSON-RPC notification (no id => notification)
90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 90 def process_notification?(data) return false unless data['method'] && !data.key?('id') # notifications/message is the Logging utility, Deprecated as a whole # in 2026-07-28 (SEP-2577). The notice belongs to the transport, not # to MCPClient::Client: a host that registered on_notification on a # ServerSSE receives log messages without a Client ever existing. It # is raised here rather than only in # {MCPClient::SubscriptionSupport#route_notification} because this # transport is the one that does not go through it — see below. warn_logging_deprecated if data['method'] == 'notifications/message' # The legacy SSE transport carries no subscriptions/listen stream, so # there is no delivery to run ahead of — but the transport's own caches # and a host that registered its invalidation on the dedicated hook # must still be told, in the order routing uses them # (see {MCPClient::ServerBase#on_cache_invalidation}). invalidate_cache_for_notification(data['method'], data['params']) if respond_to?( :invalidate_cache_for_notification, true ) notify_cache_invalidation(data['method'], data['params']) @notification_callback&.call(data['method'], data['params']) true end |
#process_response?(data, arrived = nil) ⇒ Boolean
Process a JSON-RPC response (id => response)
118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 118 def process_response?(data, arrived = nil) return false unless data['id'] # Deliver the response to the waiting caller via @sse_results only. # We intentionally do NOT write @tools_data here: request_tools_list is # the sole writer of that cache and sets it to the COMPLETE, fully # paginated list. Writing each page as it arrives would let a concurrent # list_tools observe a partial (page-1-only) cache mid-pagination. @mutex.synchronize do # The stream is peer-controlled: only ids some caller is actually # waiting on are stored. Without this check a server could stream # unsolicited responses with fresh ids and grow @sse_results without # bound for the lifetime of the client. unless @pending_request_ids.include?(data['id']) @logger.debug("Discarding unsolicited response id #{data['id'].inspect}") return true end # Dated from arrival: the waiter polls and may wake much later. (@sse_result_arrivals ||= {})[data['id']] = arrived || monotonic_now if respond_to?(:monotonic_now, true) @sse_results[data['id']] = if data['error'] # JSON-RPC error response: store the error under a Symbol key # (JSON.parse only produces String keys, so this cannot collide # with a success result) for the waiter to raise ServerError. { error: data['error'] } elsif data.key?('result') # An explicit null (or false) result is an answer; only the # member's presence decides that. data['result'] else # JSON-RPC 2.0 section 5: "Either the result member or error # member MUST be included". An envelope carrying neither answers # nothing, and reading its absent result as a delivered nil # would turn a malformed response into a successful call. { no_answer: true } end end true end |
#process_server_request?(data) ⇒ Boolean
Process a JSON-RPC request from server (has both id AND method)
80 81 82 83 84 85 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 80 def process_server_request?(data) return false unless data['method'] && data.key?('id') handle_server_request(data) true end |
#resolve_endpoint_uri(data) ⇒ String
Resolve an endpoint URI reference against the SSE connection URL. The resolved endpoint MUST stay on the SSE connection's origin: the event payload is server-controlled input, and honoring a cross-origin target would redirect every JSON-RPC POST — including the configured Authorization/API-key headers and callback response bodies — to a server the caller never chose.
212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 |
# File 'lib/mcp_client/server_sse/sse_parser.rb', line 212 def resolve_endpoint_uri(data) endpoint = URI.join(@base_url, data) base = URI.parse(@base_url) unless same_origin?(base, endpoint) fail_endpoint_handshake!( "Cross-origin endpoint in SSE endpoint event: #{data.inspect} " \ "does not match the connection origin #{origin_of(base)}" ) end endpoint.to_s rescue URI::Error => e # The endpoint event is the handshake's core payload; an unresolvable # URI must fail the handshake rather than deferring a broken POST # target to the first request. @logger.error("Failed to resolve endpoint URI #{data.inspect} against #{@base_url}: #{e.}") fail_endpoint_handshake!("Invalid endpoint URI in SSE endpoint event: #{data.inspect} (#{e.})") end |