Class: MCPClient::HttpTransportBase::StreamRecovery::ResponseBodyCapture
- Inherits:
-
Faraday::Middleware
- Object
- Faraday::Middleware
- MCPClient::HttpTransportBase::StreamRecovery::ResponseBodyCapture
- 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
- 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
- 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
-
#deliver_live_event(state, listener, event) ⇒ void
Hand one event to the stream listener.
-
#note_response_arrival(state, event) ⇒ void
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").
- #on_complete(env) ⇒ void
- #on_request(env) ⇒ void
-
#response_event?(event, id) ⇒ Boolean
Whether the event is the JSON-RPC response to it.
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).
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").
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.
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.
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.
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? = JSON.parse(data.join("\n")) .is_a?(Hash) && !.key?('method') && (['id'] == id || ['id'].to_s == id.to_s) rescue JSON::ParserError false end |