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

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.

Parameters:

  • error (Faraday::ConnectionFailed, Faraday::SSLError) —

    the socket failure

  • request (Hash) —

    the JSON-RPC message that was being sent

Returns:



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.

Parameters:

  • capture (Hash, nil) —

    the capture state of the failed exchange

Returns:

  • (Integer)


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.

Parameters:

  • error (Faraday::ConnectionFailed, Faraday::SSLError) —

    the socket failure

Returns:

  • (Boolean) —

    true when the exchange had started when it broke



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.

Parameters:

  • body (String) —

    a response body

Returns:

  • (String) —

    the body with LF line terminators



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.

Parameters:

  • partial_body (String, nil) —

    the bytes captured before the failure

  • request (Hash) —

    the JSON-RPC message that was being sent

  • error (Faraday::Error) —

    the socket failure or timeout

  • capture (Hash, nil) (defaults to: nil) —

    the capture state of the failed exchange

Returns:



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.

Parameters:

Returns:

Raises:



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.

Parameters:

  • body (String) —

    an LF-normalized SSE body

Returns:

  • (Array<String>) —

    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.

Parameters:

  • cause (Exception, nil) —

    the exception Faraday wrapped

Returns:

  • (Boolean) —

    true when TLS failed before the request was sent



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.message.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.

Parameters:

  • request (Hash) —

    the JSON-RPC message that was being sent

  • capture (Hash) —

    the capture state of the exchange

Returns:

Raises:



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.

Parameters:

  • body (String) —

    a response body

Returns:

  • (String)


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