Module: MCPClient::HttpTransportBase

Includes:
CalledToolDefinition, CacheSupport, EraDetection, ListenStream, ParamHeaders, RequestRecovery, SessionRecovery, StreamCapture, StreamRecovery, ToolListing, JsonRpcCommon
Included in:
ServerHTTP::JsonRpcTransport, ServerStreamableHTTP::JsonRpcTransport
Defined in:
lib/mcp_client/http_transport_base.rb,
lib/mcp_client/http_transport_base/tool_listing.rb,
lib/mcp_client/http_transport_base/cache_support.rb,
lib/mcp_client/http_transport_base/era_detection.rb,
lib/mcp_client/http_transport_base/listen_stream.rb,
lib/mcp_client/http_transport_base/param_headers.rb,
lib/mcp_client/http_transport_base/stream_capture.rb,
lib/mcp_client/http_transport_base/bounded_inflate.rb,
lib/mcp_client/http_transport_base/stream_recovery.rb,
lib/mcp_client/http_transport_base/request_recovery.rb,
lib/mcp_client/http_transport_base/session_recovery.rb,
lib/mcp_client/http_transport_base/sse_event_scanner.rb

Overview

Base module for HTTP-based JSON-RPC transports Contains common functionality shared between HTTP and Streamable HTTP transports

Defined Under Namespace

Modules: BoundedInflate, CacheSupport, EraDetection, ListenStream, ParamHeaders, RequestRecovery, SessionRecovery, StreamCapture, StreamRecovery, ToolListing Classes: NormalizedResponse, SseEventScanner

Constant Summary collapse

AUTH_PARAM =

One auth-param (name = token / quoted-string) as it appears in a WWW-Authenticate challenge (RFC 7235 §2.1, optional whitespace around '=').

/[A-Za-z0-9._~+-]+\s*=\s*(?:"(?:[^"\\]|\\.)*"|[^,\s]*)/
AUTH_PARAMS_RUN =

A run of comma/space separated auth-params anchored at the start of a string. The run ends before a token that is NOT followed by '=' — the auth-scheme introducing the next challenge — while commas inside quoted values are consumed by the quoted-string branch, not treated as boundaries.

/\A(?:[\s,]*#{AUTH_PARAM})*/
INTERRUPTED_EXCHANGE_ERRORS =

Socket-level failures that can only occur once the exchange was under way: the peer reset or closed the connection, the response head was truncated, or the encoded body stopped short. Failures proving the request never reached the server (connection refused, DNS failure, unreachable network) are deliberately absent — there is nothing in flight to replace.

IOError covers EOFError and Net::HTTP's own "closed stream"; Zlib::Error covers a gzip body (Streamable HTTP always offers gzip) that stops before its footer; OpenSSL::SSL::SSLError covers an HTTPS body whose TLS session dies mid-read, which is what production Streamable HTTP actually raises — see tls_handshake_failure? for the one OpenSSL case that means the exchange never started.

[
  IOError, Errno::ECONNRESET, Errno::ECONNABORTED, Errno::EPIPE,
  Net::HTTPBadResponse, Net::ProtocolError, Zlib::Error, OpenSSL::SSL::SSLError
].freeze
INTERRUPTED_EXCHANGE_FARADAY_ERRORS =

Faraday exception classes that can carry a broken response stream. TLS failures are a sibling of ConnectionFailed, not a subclass, so both must be named for an HTTPS stream to reach the re-issue path at all.

[Faraday::ConnectionFailed, Faraday::SSLError].freeze
PROTOCOL_MODES =

How the server's protocol era is established (MCP 2026-07-28 Streamable HTTP "Backward Compatibility"): :auto attempts a modern request first and falls back to the initialize handshake on a legacy rejection, :modern never falls back, :legacy never probes.

i[auto modern legacy].freeze

Constants included from CacheSupport

CacheSupport::NOOP_APP, CacheSupport::PROBE_AUTHORIZATION_NEUTRAL_MIDDLEWARE, CacheSupport::PROBE_DEFAULT_INITIALIZE_OWNERS, CacheSupport::PROBE_PURE_MIDDLEWARE, CacheSupport::PROBE_STATIC_CLASSES, CacheSupport::PROBE_STATIC_DEPTH, CacheSupport::RESPONSE_RECEIVED_AT_KEY, CacheSupport::SENT_AUTHORIZATION_KEY, CacheSupport::UNKNOWN_CONTEXT

Constants included from RequestAuthorization

RequestAuthorization::ANONYMOUS_AUTHORIZATION, RequestAuthorization::UNRECORDED_AUTHORIZATION

Constants included from ListenStream

ListenStream::COMMENT_LINE, ListenStream::EVENT_TERMINATOR, ListenStream::LINE_TERMINATOR, ListenStream::LISTEN_CLOSE_POLL_INTERVAL, ListenStream::LISTEN_MAX_BUFFER_BYTES, ListenStream::LISTEN_MAX_RECONNECT_DELAY, ListenStream::LISTEN_OPEN_TIMEOUT, ListenStream::LISTEN_RECONNECT_DELAY, ListenStream::LISTEN_STREAM_TIMEOUT, ListenStream::THREAD_JOIN_TIMEOUT_FOR_LISTEN

Constants included from JsonRpcCommon

JsonRpcCommon::CORE_RESULT_TYPES, JsonRpcCommon::EXTENSION_ID_PATTERN, JsonRpcCommon::INPUT_RETRY_DELAY, JsonRpcCommon::INPUT_RETRY_MAX_DELAY, JsonRpcCommon::LEGACY_RESULT_TYPES, JsonRpcCommon::LOG_LEVELS, JsonRpcCommon::MAX_INPUT_ROUND_TRIPS, JsonRpcCommon::MAX_PEER_LOG_TEXT_LENGTH, JsonRpcCommon::MRTR_METHODS, JsonRpcCommon::NAME_HEADER_SOURCES, JsonRpcCommon::NON_IDEMPOTENT_METHODS, JsonRpcCommon::REMOVED_MODERN_NOTIFICATIONS, JsonRpcCommon::RESULT_TYPE_EXTENSIONS, JsonRpcCommon::TASKS_EXTENSION, JsonRpcCommon::TASK_METHODS

Constants included from SessionPin

SessionPin::SESSION_PINS, SessionPin::WRITE_GUARDS

Constants included from RequestMetadata

RequestMetadata::CACHE_NEUTRAL_META_KEYS, RequestMetadata::META_CLIENT_CAPABILITIES, RequestMetadata::META_CLIENT_INFO, RequestMetadata::META_LOG_LEVEL, RequestMetadata::META_PROTOCOL_VERSION, RequestMetadata::META_SERVER_INFO, RequestMetadata::META_SUBSCRIPTION_ID, RequestMetadata::OPAQUE_PARAMS, RequestMetadata::PROTECTED_META_KEYS, RequestMetadata::UNRECORDED_PARAMS

Constants included from ResultCaching

ResultCaching::CACHE_INIT_LOCK, ResultCaching::LEGACY_ENTRY, ResultCaching::LIST_CHANGE_NOTIFICATIONS, ResultCaching::LIST_METHOD_KINDS, ResultCaching::LIST_VALUE_KINDS, ResultCaching::MAX_CACHED_READS, ResultCaching::MAX_READ_GENERATIONS, ResultCaching::PLACEHOLDER_KINDS

Constants included from InputRoundTrips

InputRoundTrips::INPUT_REQUEST_HANDLERS

Constants included from SubscriptionSupport

SubscriptionSupport::CONTROL_NOTIFICATIONS, SubscriptionSupport::DEFAULT_ACK_TIMEOUT

Constants included from JsonRpcCommon::ErrorBodies

JsonRpcCommon::ErrorBodies::MAX_ERROR_BODY_BYTES

Instance Attribute Summary collapse

Attributes included from JsonRpcCommon

#request_meta, #send_client_info

Instance Method Summary collapse

Methods included from SessionRecovery

#resend_after_session_restart

Methods included from ToolListing

#tools_generation

Methods included from ListenStream

#cancel_subscription, #close_listen_streams, #ensure_session_ready, #open_subscription

Methods included from StreamRecovery

#connection_failure_error, #delivered_status, #interrupted_exchange?, #normalize_sse_newlines, #salvaged_response, #settled_salvage, #sse_data_payloads, #tls_handshake_failure?, #truncated_body_outcome, #without_bom

Methods included from JsonRpcCommon

#accepted_result_types, #apply_discover_result, #begin_era_probe, #build_jsonrpc_notification, #build_jsonrpc_request, #build_named_request_params, #cancellable_request?, #client_capabilities, #client_info_payload, #declare_extension, #declare_sampling_tools, #declared_extensions, #describe_body_size, #describe_jsonrpc_message, #describe_parse_error, #discovery_cache_scope, #discovery_clock, #discovery_fresh?, #encode_header_value, #era_probe_in_flight?, #host_request_meta, #implemented_extension_result_types, #initialization_params, #mcp_name_header_value, #merge_meta_spellings, #modern?, #modern_request_headers, #notify_cache_invalidation, #ping, #process_jsonrpc_response, #protocol_era, #protocol_version, #record_discovery_freshness, #record_server_info, #refused_undeclared_sampling_tools?, #registered_callback?, #reject_task_result_discover!, #reject_task_result_on_unsupported_method!, #request_meta_claim, #required_request_meta, #reserved_meta_supplied?, #resolve_input_round_trips, restore_wire_keys, result_type, #sampling_tools_supported?, #sanitize_log_text, #select_protocol_version, #send_client_info?, #settle_era_probe, #spend_held_request_meta, #split_request_meta, #supported_versions, #suppressed_modern_notification?, #tasks_extension_declared?, #validate_log_level!, #validate_protocol_version!, #validate_result_type!, #with_request_meta, #with_retry

Methods included from SessionPin

#check_session_pin!, #guarded_writes, #pinned_to_session, #unpinned_session

Methods included from RequestMetadata

#adoptable_request_meta_hold, #claimable_request_meta_hold, #close_request_meta_hold, #current_params_fingerprint, #deep_sort_keys, #held_request_meta, #held_request_meta_key, #holding_request_meta, #note_request_params, #note_request_params_pending, #offer_request_meta_hold, #offered_request_meta_key, #open_request_meta_hold, #outside_request_meta_hold, #params_fingerprint_of, #recorded_request_params, #release_held_request_meta, #request_params_fingerprint, #request_params_key, #restore_request_params, #take_offered_request_meta_hold, #withdraw_request_meta_hold

Methods included from ResultCaching

#assume_zero_ttl?, #attach_list_value, #authorization_fingerprint, authorization_header_value, #authorization_header_value, #bind_authorization_context, #bump_cache_epoch, #bump_cache_generation, #cache_entries, #cache_entries_mutex, #cache_entry_for, #cache_entry_fresh?, #cache_entry_hinted?, #cache_entry_token, #cache_epoch, #cache_fresh?, #cache_generation, #cache_info, #cached_list_value, #clear_response_received_at, #clear_result_cache, #discard_paginated_list, #empty_list_copy?, #entry_for_current_params?, #entry_in_current_context?, #entry_matches_authorization?, faraday_headers, #faraday_headers, #fetching_list_page, #forget_served_entries, #forget_transport_thread_state, #fresh_list_value, #hinted_list_value, #invalid_cursor_error?, #invalidate_cache, #invalidate_cache_for_notification, #invalidate_read_cache, #list_cache_epoch, #list_kind_for, #mixed_pages_placeholder, #monotonic_now, #note_legacy_served, #note_response_received_at, #note_served_entry, #on_cache_invalidation, #private_entry_for_current_context, #prune_read_entries, #read_cache_key, #read_resource_with_cache, #record_cache_hint, #record_list_cache_hint, #record_paginated_cache_hint, #recorded_entries_key, #release_serving_request_meta, #remember_recorded_entry, #response_received_at, #response_received_key, #sent_authorization_known?, #served_entries_key, #stale_fallback_for, #stale_list_entry, #stale_list_value, #store_read_entry, #take_served_entry, #transport_thread_local_keys

Methods included from InputRoundTrips

#fulfil_input_request, #fulfil_input_requests, #undeclared_sampling_tool_use?

Methods included from SubscriptionSupport

#await_acknowledgment_deadline, #close_subscription_gracefully, #confirm_resource_subscription, #deliver_subscription_notification, #discard_mapped_resource_subscription, #drop_unacknowledged_resource_subscriptions, #ensure_modern_listen!, #expire_unacknowledged_subscription, #handle_server_cancellation, #handle_subscription_acknowledgment, #handle_subscription_control, #handle_subscription_response, #listen, #live_resource_subscription, #notify_host, #open_listen, #open_resource_subscription, #rearm_acknowledgment_deadline, #recheck_mapped_resource_subscription, #register_subscription, #report_declined_subscription_types, #require_resource_watch, #resource_subscription_mutex, #resource_subscription_mutexes, #resource_subscriptions, #route_notification, #settled_resource_subscription, #subscribe_resource_via_listen, #subscription_ack_timeout, #subscription_by_id, #subscription_delivery_target, #subscription_for_notification, #subscriptions, #subscriptions_mutex, #unmap_resource_subscription, #unregister_subscription, #unregister_subscription_id, #unsubscribe_resource_via_listen

Methods included from DeprecationNotices

#warn_input_request_answer_deprecated, #warn_input_request_deprecated, #warn_logging_deprecated, #warn_request_log_level_deprecated, #warn_roots_deprecated, #warn_sampling_deprecated

Methods included from JsonRpcCommon::InputWaits

#on_input_required_wait, #reject_input_required_discover!, #resume_input_required

Methods included from JsonRpcCommon::ErrorBodies

#decoded_error_body, #gunzip_bounded, #jsonrpc_error_from_http_response, #jsonrpc_error_in_body, #oversized_error_body?

Instance Attribute Details

#discover_timeout ⇒ Numeric (readonly)

Returns seconds allowed for the server/discover probe.

Returns:

  • (Numeric) —

    seconds allowed for the server/discover probe



280
281
282
# File 'lib/mcp_client/http_transport_base.rb', line 280

def discover_timeout
  @discover_timeout
end

#protocol_mode ⇒ Symbol (readonly)

Returns the configured protocol mode (:auto, :modern or :legacy).

Returns:

  • (Symbol) —

    the configured protocol mode (:auto, :modern or :legacy)



277
278
279
# File 'lib/mcp_client/http_transport_base.rb', line 277

def protocol_mode
  @protocol_mode
end

Instance Method Details

#rpc_notify(method, params = {}) ⇒ void

This method returns an undefined value.

Send a JSON-RPC notification (no response expected)

Parameters:

  • method (String) —

    JSON-RPC method name

  • params (Hash) (defaults to: {}) —

    parameters for the notification



163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
# File 'lib/mcp_client/http_transport_base.rb', line 163

def rpc_notify(method, params = {})
  ensure_connected
  if suppressed_modern_notification?(method)
    @logger.debug("Not sending #{method}: removed in MCP #{protocol_version}")
    return
  end

  notif = build_jsonrpc_notification(method, params)

  begin
    send_http_request(notif)
  rescue MCPClient::Errors::ServerError, MCPClient::Errors::ConnectionError, Faraday::ConnectionFailed => e
    raise MCPClient::Errors::TransportError, "Failed to send notification: #{e.message}"
  end
end

#rpc_request(method, params = {}, timeout: nil) ⇒ Object

Generic JSON-RPC request: send method with params and return result

Parameters:

  • method (String) —

    JSON-RPC method name

  • params (Hash) (defaults to: {}) —

    parameters for the request

Returns:

  • (Object) —

    result from JSON-RPC response

Raises:



84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
# File 'lib/mcp_client/http_transport_base.rb', line 84

def rpc_request(method, params = {}, timeout: nil)
  freshly_probed = !@mutex.synchronize { @connection_established }
  ensure_connected
  if method == 'ping' && modern?
    # `ping` was removed in MCP 2026-07-28; the mandatory server/discover
    # request is the modern heartbeat, and the probe that just established
    # the connection already was one.
    return @last_discover_result if freshly_probed && @last_discover_result

    method = 'server/discover'
  end

  header_refresh_done = false
  # The multi round-trip resolver sits outside the per-attempt recovery,
  # so a retry carrying inputResponses/requestState keeps them through
  # version renegotiation, the HeaderMismatch refresh and a re-issued
  # stream. Each attempt is a request of its own, with its own id and its
  # own budget; the deadline lives in attempt_request.
  result = resolve_input_round_trips(method, params, timeout) do |attempt_params|
    attempt_request(method, attempt_params, timeout, header_refresh_done) { header_refresh_done = true }
  end
  # Every server/discover answer is validated and applied: a later
  # heartbeat may advertise new versions or capabilities.
  result = apply_discover_result(result) if method == 'server/discover'
  result
end

#send_cancellation_notification(request_id) ⇒ void

This method returns an undefined value.

Best-effort notifications/cancelled for a request the client stopped waiting on. Failures are swallowed.

It is sent for the abandoned request, on that request's own thread and after it, and it brings nothing back to cache: the credentials it carries are whatever the host holds by now -- a rotation, a refresh -- and they must not stand in for the ones the abandoned request went out with, which are what its failure is judged by (MCP 2026-07-28 caching, cacheScope "private": a stale copy may be served only to the context the failed request itself carried).

Parameters:

  • request_id (Integer) —

    id of the abandoned request



146
147
148
149
150
151
152
153
154
155
156
157
# File 'lib/mcp_client/http_transport_base.rb', line 146

def send_cancellation_notification(request_id)
  notif = build_jsonrpc_notification('notifications/cancelled',
                                     { 'requestId' => request_id, 'reason' => 'Request timed out' })
  abandoned = recorded_request_authorization
  begin
    send_http_request(notif)
  ensure
    restore_request_authorization(abandoned)
  end
rescue StandardError => e
  @logger.debug("Failed to send cancellation notification: #{e.message}")
end

#send_request_and_parse(method, params, timeout, deadline = nil) ⇒ Object

One request/response exchange with its own JSON-RPC id.

Parameters:

  • method (String) —

    JSON-RPC method name

  • params (Hash) —

    parameters for the request

  • timeout (Numeric, nil) —

    per-request timeout override

  • deadline (Float, nil) (defaults to: nil) —

    monotonic instant this exchange and the one replacement the re-issue rule allows must finish by

Returns:

  • (Object) —

    result from the JSON-RPC response



118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
# File 'lib/mcp_client/http_transport_base.rb', line 118

def send_request_and_parse(method, params, timeout, deadline = nil)
  request_id = @mutex.synchronize { @request_id += 1 }
  request = build_jsonrpc_request(method, params, request_id)
  # Computed before sending so a value that cannot be mirrored fails the
  # call locally (ValidationError) rather than mid-request.
  param_headers = modern? ? mcp_param_headers(request) : {}
  send_jsonrpc_request(request, timeout: timeout, deadline: deadline, extra_headers: param_headers)
rescue MCPClient::Errors::RequestTimeoutError
  # MCP lifecycle: on timeout the sender SHOULD cancel the abandoned
  # request. On modern Streamable HTTP closing the response stream IS the
  # cancellation signal and no notifications/cancelled is expected; legacy
  # servers still get the notification.
  send_cancellation_notification(request_id) if !modern? && cancellable_request?(method, params)
  raise
end

#terminate_session ⇒ Boolean

Terminate the current session with the server Sends an HTTP DELETE request with the session ID to properly close the session

Returns:

  • (Boolean) —

    true if termination was successful

Raises:



183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
# File 'lib/mcp_client/http_transport_base.rb', line 183

def terminate_session
  # MCP 2026-07-28 removed the session layer: a modern connection has no
  # session to terminate and MUST NOT send the DELETE, whatever a
  # non-conforming server (or a caller) put in @session_id.
  if modern?
    @session_id = nil
    return true
  end

  return true unless @session_id

  # The session is over from here whatever the DELETE answers (every
  # outcome below clears the id), and it ends without a #cleanup: the
  # epoch moves so nothing keyed by it — the tasks extension's task ids,
  # answered and pending input keys — outlives it into the session the
  # next request establishes, which may reuse those very ids.
  bump_session_epoch
  conn = http_connection

  begin
    @logger.debug("Terminating session: #{@session_id}")
    response = conn.delete(@endpoint) do |req|
      # Apply base headers but prioritize session termination headers
      @headers.each { |k, v| req.headers[k] = v }
      req.headers['Mcp-Session-Id'] = @session_id
      req.headers['Mcp-Protocol-Version'] = @protocol_version if @protocol_version
      # MCP: authorization MUST be included in every HTTP request
      @oauth_provider&.apply_authorization(req)
      note_request_authorization(authorization_header_value(req.headers))
    end

    if response.success?
      @logger.debug("Session terminated successfully: #{@session_id}")
      @session_id = nil
      true
    else
      @logger.warn("Session termination failed with HTTP #{response.status}")
      @session_id = nil # Clear session ID even on HTTP error
      false
    end
  rescue Faraday::Error => e
    @logger.warn("Session termination request failed: #{e.message}")
    # Clear session ID even if termination request failed
    @session_id = nil
    false
  end
end

#valid_server_url?(url) ⇒ Boolean

Validate the server's base URL for security

Parameters:

  • url (String) —

    the URL to validate

Returns:

  • (Boolean) —

    true if URL is considered safe



249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
# File 'lib/mcp_client/http_transport_base.rb', line 249

def valid_server_url?(url)
  return false unless url.is_a?(String)

  uri = URI.parse(url)

  # Only allow HTTP and HTTPS protocols
  return false unless %w[http https].include?(uri.scheme)

  # Must have a host
  return false if uri.host.nil? || uri.host.empty?

  # Don't allow localhost binding to all interfaces in production
  if uri.host == '0.0.0.0'
    @logger.warn('Server URL uses 0.0.0.0 which may be insecure. Consider using 127.0.0.1 for localhost.')
  end

  true
rescue URI::InvalidURIError
  false
end

#valid_session_id?(session_id) ⇒ Boolean

Validate session ID format Per MCP 2025-11-25, the server-assigned session ID "MUST only contain visible ASCII characters (ranging from 0x21 to 0x7E)" — e.g. a UUID, a JWT, or a cryptographic hash — and the client MUST echo whatever the server assigned. A generous length cap guards against abuse.

Parameters:

  • session_id (String) —

    the session ID to validate

Returns:

  • (Boolean) —

    true if session ID is valid



238
239
240
241
242
243
244
# File 'lib/mcp_client/http_transport_base.rb', line 238

def valid_session_id?(session_id)
  return false unless session_id.is_a?(String)

  # The 4096-char cap is header-size hygiene, not MCP grammar — the spec
  # imposes no length limit on session IDs.
  session_id.match?(/\A[\x21-\x7E]{1,4096}\z/)
end