Class: MCPClient::HttpTransportBase::StreamRecovery::ResponseBodyCapture

Inherits:
Faraday::Middleware
  • Object
show all
Defined in:
lib/mcp_client/http_transport_base/stream_recovery.rb

Overview

Innermost Faraday middleware: streams the response body into a per-request buffer so that

  1. a socket failure mid-body still leaves the bytes that did arrive (Faraday discards a partially read body and raises), letting a response that was fully delivered settle its request instead of being re-issued and executed twice; and
  2. a deadline can be enforced while the body is arriving, which a socket timeout alone cannot do for a stream that keeps dripping keep-alives.

It restores the buffer as the response body, and being the innermost handler its on_complete runs before any user middleware (raise_error and friends) looks at that body.

Instance Method Summary collapse

Instance Method Details

#deliver_live_event(state, listener, event) ⇒ void

This method returns an undefined value.

Hand one event to the stream listener. A failure there is the exchange's failure, but raising it here would abort the read — and on MCP 2026-07-28 a client closing the response stream is the cancellation signal — so the first failure is held for the transport to raise once the body has been read (StreamCapture#stream_listener_error).

Parameters:

  • state (Hash) —

    the exchange's capture state

  • listener (Proc) —

    the stream listener

  • event (String) —

    one complete SSE event



86
87
88
89
90
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 86

def deliver_live_event(state, listener, event)
  listener.call(event)
rescue StandardError => e
  state[:mcp_stream_error] ||= e
end

#note_response_arrival(state, event) ⇒ void

This method returns an undefined value.

Flag the event carrying the answer to this request, so whoever dates the response (CacheSupport's recorder) can date it from the chunk that completed the result rather than from whatever opened the stream: a keep-alive or a progress notification is not the result, and a TTL that ran from it would expire results that took a while to compute (MCP 2026-07-28 caching, "Freshness Calculation").

Parameters:

  • state (Hash) —

    the exchange's capture state

  • event (String) —

    one complete SSE event



101
102
103
104
105
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 101

def note_response_arrival(state, event)
  return if state[:mcp_response_seen] || !state.key?(:mcp_response_id)

  state[:mcp_response_seen] = true if response_event?(event, state[:mcp_response_id])
end

#on_complete(env) ⇒ void

This method returns an undefined value.

Parameters:

  • env (Faraday::Env) —

    the completed request environment



122
123
124
125
126
127
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 122

def on_complete(env)
  state = env.request&.context
  buffer = state && state[:mcp_body_buffer]
  env.body = buffer.dup if buffer && env.body.to_s.empty?
  state[:mcp_short_body] = short_body?(env, state, buffer) if state
end

#on_request(env) ⇒ void

This method returns an undefined value.

Parameters:

  • env (Faraday::Env) —

    the outgoing request environment



30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 30

def on_request(env)
  state = env.request&.context
  buffer = state && state[:mcp_body_buffer]
  return unless buffer

  # The retry middleware sits above this one and replays the whole inner
  # stack, so each attempt must start from an empty buffer (and from an
  # empty event scanner: the count of events dispatched while the body
  # arrived belongs to the attempt whose body is finally parsed).
  buffer.clear
  listener = state[:mcp_stream_listener]
  scanner = listener && SseEventScanner.new(max_inflated_bytes: state[:mcp_inflate_limit])
  state[:mcp_live_events] = 0
  # The adapter fills this same env in as it reads: its status is set
  # from the status line, so a salvaged answer can be rebuilt under the
  # status it really arrived with. The era rule reads a recognized
  # modern error only under the status it came with.
  state[:mcp_env] = env
  env.request.on_data = lambda do |chunk, _size, _env|
    # Before the chunk is kept, never after: bytes that arrive past the
    # deadline are not part of an answer this request may settle on, and
    # buffering them first would let the salvage hand back an answer the
    # caller had already stopped waiting for.
    deadline = state[:mcp_deadline]
    raise Faraday::TimeoutError, 'Request exceeded its deadline' if deadline && monotonic_now > deadline

    buffer << chunk.to_s
    # Only a streamed body can be measured against its Content-Length
    # here; a response the adapter hands over whole never reaches this.
    state[:mcp_streamed] = true

    next unless scanner

    # Every complete event is handed over as it arrives, so a server
    # request or a progress notification on the stream is acted on
    # while the response is still open (a server that waits for its
    # ping to be answered before sending the result would otherwise
    # deadlock against a client that answers only at EOF).
    scanner.feed(chunk.to_s) do |event|
      note_response_arrival(state, event)
      deliver_live_event(state, listener, event)
    end
    state[:mcp_live_events] = scanner.count
  end
end

#response_event?(event, id) ⇒ Boolean

Returns whether the event is the JSON-RPC response to it.

Parameters:

  • event (String) —

    one complete SSE event

  • id (Integer, String) —

    the id of the request awaiting its answer

Returns:

  • (Boolean) —

    whether the event is the JSON-RPC response to it



110
111
112
113
114
115
116
117
118
# File 'lib/mcp_client/http_transport_base/stream_recovery.rb', line 110

def response_event?(event, id)
  data = event.lines.filter_map { |line| line[5..].to_s.sub(/\A /, '').chomp if line.start_with?('data:') }
  return false if data.empty?

  message = JSON.parse(data.join("\n"))
  message.is_a?(Hash) && !message.key?('method') && (message['id'] == id || message['id'].to_s == id.to_s)
rescue JSON::ParserError
  false
end