Module: MCPClient::HttpTransportBase
- Includes:
- CalledToolDefinition, CacheSupport, EraDetection, ListenStream, ParamHeaders, RequestRecovery, SessionRecovery, StreamCapture, StreamRecovery, ToolListing, JsonRpcCommon
- 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
-
#discover_timeout ⇒ Numeric
readonly
Seconds allowed for the server/discover probe.
-
#protocol_mode ⇒ Symbol
readonly
The configured protocol mode (:auto, :modern or :legacy).
Attributes included from JsonRpcCommon
#request_meta, #send_client_info
Instance Method Summary collapse
-
#rpc_notify(method, params = {}) ⇒ void
Send a JSON-RPC notification (no response expected).
-
#rpc_request(method, params = {}, timeout: nil) ⇒ Object
Generic JSON-RPC request: send method with params and return result.
-
#send_cancellation_notification(request_id) ⇒ void
Best-effort notifications/cancelled for a request the client stopped waiting on.
-
#send_request_and_parse(method, params, timeout, deadline = nil) ⇒ Object
One request/response exchange with its own JSON-RPC id.
-
#terminate_session ⇒ Boolean
Terminate the current session with the server Sends an HTTP DELETE request with the session ID to properly close the session.
-
#valid_server_url?(url) ⇒ Boolean
Validate the server's base URL for security.
-
#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.
Methods included from SessionRecovery
Methods included from ToolListing
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.
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).
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)
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
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).
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 = begin send_http_request(notif) ensure (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.
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
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&.(req) ((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
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.
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 |