Module: MCPClient::HttpTransportBase::StreamRecovery
- Included in:
- MCPClient::HttpTransportBase
- Defined in:
- lib/mcp_client/http_transport_base/stream_recovery.rb
Overview
What a response stream that broke leaves behind, and what to make of it.
MCP 2026-07-28 has no resumption: "a broken response stream loses the in-flight request; clients MUST re-issue it as a new request with a new request ID". Faraday discards a partially read body and raises, so without capturing the bytes as they arrive a break after the final event is indistinguishable from one before it -- and re-issuing then runs a completed call a second time.
Defined Under Namespace
Classes: ResponseBodyCapture
Instance Method Summary collapse
-
#connection_failure_error(error, request) ⇒ MCPClient::Errors::MCPError
Translate a Faraday socket failure into the MCP error the caller must act on.
-
#delivered_status(capture) ⇒ Integer
The status a salvaged answer arrived under.
-
#interrupted_exchange?(error) ⇒ Boolean
Faraday wraps every socket failure in ConnectionFailed (or, for TLS, in SSLError), whether the connection was never established or it broke with a request in flight; only the wrapped exception distinguishes them.
-
#normalize_sse_newlines(body) ⇒ String
Per the SSE specification a line is terminated by CRLF, CR or LF alone; normalizing to LF lets one set of framing rules serve all three.
-
#salvaged_response(partial_body, request, error, capture = nil) ⇒ NormalizedResponse?
The response that did arrive before the socket died, when the stream carried this request's complete answer.
-
#settled_salvage(salvaged) ⇒ NormalizedResponse
A salvaged answer read the way the unbroken path reads one: an error status it arrived under still becomes the typed JSON-RPC error, so a recognized modern error keeps the status the era rule needs.
-
#sse_data_payloads(body) ⇒ Array<String>
The joined data payload of each event.
-
#tls_handshake_failure?(cause) ⇒ Boolean
OpenSSL names the failing operation in its message.
-
#truncated_body_outcome(request, capture) ⇒ NormalizedResponse
What a body that stopped short of its Content-Length settles: the answer if it is all there anyway (the bytes that arrived carry this request's response, and the rest was framing), otherwise the loss the re-issue rule is written for.
-
#without_bom(body) ⇒ String
The UTF-8 decode step of the SSE algorithm drops one leading byte-order mark; the field it precedes must still be recognized.
Instance Method Details
#connection_failure_error(error, request) ⇒ MCPClient::Errors::MCPError
Translate a Faraday socket failure into the MCP error the caller must act on.
A response stream that dies mid-body is what a broken stream actually looks like on the wire: Faraday raises rather than handing back a truncated body, so it never reaches the SSE parser that recognises a stream which closed between events. MCP 2026-07-28 has no resumption — "a broken response stream loses the in-flight request; clients MUST re-issue it as a new request with a new request ID" (changelog, major change 9) — and the rule does not care where the break landed. Raising ResponseStreamClosedError puts both breaks on the one re-issue path.
A failure that never got the request out, and a notification (which has no response to lose), stay a plain ConnectionError.
175 176 177 178 179 180 181 182 183 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 175 def connection_failure_error(error, request) if modern? && request.is_a?(Hash) && request.key?('id') && interrupted_exchange?(error) return MCPClient::Errors::ResponseStreamClosedError.new( "Response stream closed before delivering the response: #{error.message}" ) end MCPClient::Errors::ConnectionError.new("Server connection lost: #{error.message}") end |
#delivered_status(capture) ⇒ Integer
The status a salvaged answer arrived under. A well-formed -32022 in a 400 body identifies a modern server and is retried with an advertised version, while the same body under 200 is a permissive legacy echo: rebuilding every salvaged answer as 200 would turn the first into the second. The adapter fills the captured env in as it reads, so its status is the status line this response really carried.
255 256 257 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 255 def delivered_status(capture) (capture.is_a?(Hash) && capture[:mcp_env]&.status) || 200 end |
#interrupted_exchange?(error) ⇒ Boolean
Faraday wraps every socket failure in ConnectionFailed (or, for TLS, in SSLError), whether the connection was never established or it broke with a request in flight; only the wrapped exception distinguishes them.
190 191 192 193 194 195 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 190 def interrupted_exchange?(error) cause = (error.wrapped_exception if error.respond_to?(:wrapped_exception)) || error.cause return false if tls_handshake_failure?(cause) INTERRUPTED_EXCHANGE_ERRORS.any? { |klass| cause.is_a?(klass) } end |
#normalize_sse_newlines(body) ⇒ String
Per the SSE specification a line is terminated by CRLF, CR or LF alone; normalizing to LF lets one set of framing rules serve all three.
291 292 293 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 291 def normalize_sse_newlines(body) without_bom(body).gsub(/\r\n|\r/, "\n") end |
#salvaged_response(partial_body, request, error, capture = nil) ⇒ NormalizedResponse?
The response that did arrive before the socket died, when the stream carried this request's complete answer.
Faraday discards a partially read body and raises, so without the streamed capture a break after the final SSE event is indistinguishable from a break before it — and re-issuing there would run a tools/call the server already executed a second time. MCP 2026-07-28's re-issue rule is about an in-flight request that was lost; a delivered response settles its request, however the socket ends afterwards. A socket that stalls after the final event until the timeout is the same case from the other direction: the answer arrived, the framing after it did not.
225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 225 def salvaged_response(partial_body, request, error, capture = nil) return nil unless modern? && request.is_a?(Hash) && request.key?('id') return nil unless error.is_a?(Faraday::TimeoutError) || interrupted_exchange?(error) body = partial_body.to_s body = inflate_delivered_gzip(body) if body.b.start_with?(SseEventScanner::GZIP_MAGIC) return nil if body.nil? || body.empty? sse = sse_framed_body?(body) # A truncated stream's last event has no terminating blank line, so it # was never dispatched (HTML SSE parsing rules) and must be dropped # before asking whether the answer arrived. body = complete_sse_events(body) if sse return nil if body.empty? || !body_carries_response?(body, sse, request['id']) @logger.warn("Response stream ended after the response arrived (#{error.message}); " \ "keeping the delivered #{request['method']} response instead of re-issuing it") NormalizedResponse.new(delivered_status(capture), { 'content-type' => sse ? 'text/event-stream' : 'application/json' }, body, capture) end |
#settled_salvage(salvaged) ⇒ NormalizedResponse
A salvaged answer read the way the unbroken path reads one: an error status it arrived under still becomes the typed JSON-RPC error, so a recognized modern error keeps the status the era rule needs. Returning it unread would settle a 400 rejection as if it were a 200 result.
282 283 284 285 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 282 def settled_salvage(salvaged) handle_http_error_response(salvaged) unless (200..299).cover?(salvaged.status.to_i) salvaged end |
#sse_data_payloads(body) ⇒ Array<String>
Returns the joined data payload of each event.
308 309 310 311 312 313 314 315 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 308 def sse_data_payloads(body) body.split("\n\n").filter_map do |event| lines = event.lines.map(&:chomp).select { |line| line.start_with?('data:') } next if lines.empty? lines.map { |line| line.sub(/\Adata:\s*/, '') }.join("\n") end end |
#tls_handshake_failure?(cause) ⇒ Boolean
OpenSSL names the failing operation in its message. A handshake that never completed ("SSL_connect ... certificate verify failed") means the request never left this client, so there is nothing in flight to replace; a body that dies mid-read ("SSL_read: unexpected eof while reading") is a broken response stream like any other.
204 205 206 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 204 def tls_handshake_failure?(cause) cause.is_a?(OpenSSL::SSL::SSLError) && cause..to_s.include?('SSL_connect') end |
#truncated_body_outcome(request, capture) ⇒ NormalizedResponse
What a body that stopped short of its Content-Length settles: the answer if it is all there anyway (the bytes that arrived carry this request's response, and the rest was framing), otherwise the loss the re-issue rule is written for.
267 268 269 270 271 272 273 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 267 def truncated_body_outcome(request, capture) error = Faraday::ConnectionFailed.new(EOFError.new('response body stopped short of its Content-Length')) salvaged = salvaged_response(capture[:mcp_body_buffer], request, error, capture) return settled_salvage(salvaged) if salvaged raise connection_failure_error(error, request) end |
#without_bom(body) ⇒ String
The UTF-8 decode step of the SSE algorithm drops one leading byte-order mark; the field it precedes must still be recognized.
299 300 301 302 303 304 |
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 299 def without_bom(body) bom = body.encoding == Encoding::BINARY ? SseEventScanner::BOM : "\uFEFF".encode(body.encoding) body.start_with?(bom) ? body[bom.length..] : body rescue EncodingError body end |