Class: MCPClient::ServerStreamableHTTP

Inherits:
ServerBase
  • Object
show all
Includes:
RequestMetaScope, JsonRpcTransport
Defined in:
lib/mcp_client/server_streamable_http.rb,
lib/mcp_client/server_streamable_http/json_rpc_transport.rb

Overview

Implementation of MCP server that communicates via Streamable HTTP transport (MCP 2025-06-18) This transport uses HTTP POST for RPC calls with optional SSE responses, and GET for event streams Compliant with MCP specification version 2025-06-18

Key features:

  • Supports server-sent events (SSE) for real-time notifications
  • Handles ping/pong keepalive mechanism
  • Thread-safe connection management
  • Automatic reconnection with exponential backoff

Defined Under Namespace

Modules: JsonRpcTransport

Constant Summary collapse

DEFAULT_READ_TIMEOUT =

Default values for connection settings

30
DEFAULT_MAX_RETRIES =
3
SSE_CONNECTION_TIMEOUT =

SSE connection settings

300
SSE_RECONNECT_DELAY =

5 minutes

1
SSE_MAX_RECONNECT_DELAY =

Initial reconnect delay in seconds

30
THREAD_JOIN_TIMEOUT =

Maximum reconnect delay in seconds

5
MIN_RESUMPTION_RECONNECT_DELAY =

Floor for the delay between resumption GETs. SEP-1699's polling pattern wants fast reconnects, so this is far smaller than the events-stream floor — it only prevents a peer-supplied "retry: 0" from turning the deadline window into a back-to-back request loop.

0.01
MAX_EVENT_ID_LENGTH =

Maximum length of a server-supplied SSE event id retained as the resumption cursor. The id is echoed in the Last-Event-ID header of subsequent requests, so an unbounded value means unbounded retained memory and oversized outbound headers.

1024
EVENT_ID_PATTERN =

Characters allowed in a retained event id: printable ASCII, since the value becomes an HTTP header value. Notably excludes CR/LF. The range starts at 0x20 because a space is legal inside a field value, and rejecting ids like "cursor 42" would silently strand resumption on a stale cursor.

/\A[\x20-\x7E]+\z/
MAX_CONCURRENT_RESPONSE_POSTS =

Ceiling on concurrent threads POSTing server-initiated responses (pongs, roots/sampling/elicitation replies, error responses). Each server request on the events stream costs one blocking HTTP POST in its own thread; without a bound, a peer flooding requests could accumulate threads and connections until the host is exhausted. Responses beyond the budget are dropped, with saturation logged at most once per SATURATION_LOG_INTERVAL seconds.

8
SATURATION_LOG_INTERVAL =

Minimum gap between "response budget saturated" warnings

5
MIN_EVENTS_RECONNECT_DELAY =

Floor for server-supplied retry directives on the long-lived events stream. The directive is peer-controlled: honoring "retry: 0" literally would let a hostile server that closes every stream drive a tight reconnect loop (sustained CPU/TLS/connection churn). Waiting longer than the directive stays SEP-1699 compliant — the retry field is a lower bound on the reconnect delay, not an exact schedule.

0.1
MAX_SSE_BUFFER_BYTES =

Maximum bytes an SSE parse buffer (events stream or resumption GET) may hold while waiting for an event terminator. The stream is peer-controlled: without a cap, a hostile server could withhold the blank-line delimiter forever and grow the buffer until the host runs out of memory. Generous enough for any legitimate JSON-RPC event.

32 * 1024 * 1024

Constants included from JsonRpcTransport

JsonRpcTransport::DECOMPRESS_CHUNK_BYTES, JsonRpcTransport::MAX_DECOMPRESSED_BODY_BYTES

Constants included from HttpTransportBase

HttpTransportBase::AUTH_PARAM, HttpTransportBase::AUTH_PARAMS_RUN, HttpTransportBase::INTERRUPTED_EXCHANGE_ERRORS, HttpTransportBase::INTERRUPTED_EXCHANGE_FARADAY_ERRORS, HttpTransportBase::PROTOCOL_MODES

Constants included from HttpTransportBase::CacheSupport

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

Constants included from RequestAuthorization

RequestAuthorization::ANONYMOUS_AUTHORIZATION, RequestAuthorization::UNRECORDED_AUTHORIZATION

Constants included from HttpTransportBase::ListenStream

HttpTransportBase::ListenStream::COMMENT_LINE, HttpTransportBase::ListenStream::EVENT_TERMINATOR, HttpTransportBase::ListenStream::LINE_TERMINATOR, HttpTransportBase::ListenStream::LISTEN_CLOSE_POLL_INTERVAL, HttpTransportBase::ListenStream::LISTEN_MAX_BUFFER_BYTES, HttpTransportBase::ListenStream::LISTEN_MAX_RECONNECT_DELAY, HttpTransportBase::ListenStream::LISTEN_OPEN_TIMEOUT, HttpTransportBase::ListenStream::LISTEN_RECONNECT_DELAY, HttpTransportBase::ListenStream::LISTEN_STREAM_TIMEOUT, HttpTransportBase::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

MCPClient::SessionPin::SESSION_PINS, MCPClient::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

MCPClient::SubscriptionSupport::CONTROL_NOTIFICATIONS, MCPClient::SubscriptionSupport::DEFAULT_ACK_TIMEOUT

Constants included from JsonRpcCommon::ErrorBodies

JsonRpcCommon::ErrorBodies::MAX_ERROR_BODY_BYTES

Constants included from RequestMetaScope

RequestMetaScope::HOST_CALLBACKS, RequestMetaScope::SCOPED_OPERATIONS

Constants inherited from ServerBase

MCPClient::ServerBase::MAX_LIST_PAGES, MCPClient::ServerBase::RELATED_TASK_META_KEY

Instance Attribute Summary collapse

Attributes included from HttpTransportBase

#discover_timeout, #protocol_mode

Attributes included from JsonRpcCommon

#request_meta, #send_client_info

Attributes inherited from ServerBase

#instructions, #logger, #name, #read_timeout

Instance Method Summary collapse

Methods included from HttpTransportBase

#rpc_notify, #rpc_request, #send_cancellation_notification, #send_request_and_parse, #valid_server_url?, #valid_session_id?

Methods included from HttpTransportBase::SessionRecovery

#resend_after_session_restart

Methods included from HttpTransportBase::ToolListing

#tools_generation

Methods included from HttpTransportBase::ListenStream

#cancel_subscription, #close_listen_streams, #ensure_session_ready, #open_subscription

Methods included from HttpTransportBase::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?

Methods inherited from ServerBase

#cancel_subscription, #capability?, #client_info=, #discovery_refresh_needed?, #listen, #merge_related_task_meta, #on_cache_invalidation, #on_notification, #ping, #require_capability!, #resource_not_found_error, #resource_not_found_response?, #rpc_notify, #rpc_request, #session_epoch

Constructor Details

#initialize(base_url:, **options) ⇒ ServerStreamableHTTP

Returns a new instance of ServerStreamableHTTP.

Parameters:

  • base_url (String) —

    The base URL of the MCP server

  • options (Hash) —

    Server configuration options (same as ServerHTTP)

Raises:

  • (ArgumentError)


105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
# File 'lib/mcp_client/server_streamable_http.rb', line 105

def initialize(base_url:, **options)
  opts = default_options.merge(options)
  super(name: opts[:name])
  initialize_logger(opts[:logger])

  @max_retries = opts[:retries]
  @retry_backoff = opts[:retry_backoff]

  # Validate and normalize base_url
  raise ArgumentError, "Invalid or insecure server URL: #{base_url}" unless valid_server_url?(base_url)

  # Normalize base_url and handle cases where full endpoint is provided in base_url
  uri = URI.parse(base_url.chomp('/'))

  # Helper to build base URL without default ports
  build_base_url = lambda do |parsed_uri|
    port_part = if parsed_uri.port &&
                   !((parsed_uri.scheme == 'http' && parsed_uri.port == 80) ||
                     (parsed_uri.scheme == 'https' && parsed_uri.port == 443))
                  ":#{parsed_uri.port}"
                else
                  ''
                end
    "#{parsed_uri.scheme}://#{parsed_uri.host}#{port_part}"
  end

  @base_url = build_base_url.call(uri)
  @endpoint = if uri.path && !uri.path.empty? && uri.path != '/' && opts[:endpoint] == '/rpc'
                # If base_url contains a path and we're using default endpoint,
                # treat the path as the endpoint and use the base URL without path
                uri.path
              else
                # Standard case: base_url is just scheme://host:port, endpoint is separate
                opts[:endpoint]
              end

  # Set up headers for Streamable HTTP requests
  @headers = opts[:headers].merge({
                                    'Content-Type' => 'application/json',
                                    'Accept' => 'text/event-stream, application/json',
                                    'Accept-Encoding' => 'gzip',
                                    'User-Agent' => "ruby-mcp-client/#{MCPClient::VERSION}",
                                    'Cache-Control' => 'no-cache'
                                  })

  @read_timeout = opts[:read_timeout]
  configure_protocol_mode(opts[:protocol], opts[:discover_timeout])
  @faraday_config = opts[:faraday_config]
  @max_decompressed_body_bytes = validate_decompression_limit(opts[:max_decompressed_body_bytes])
  @tools = nil
  @tools_data = nil
  @prompts = nil
  @prompts_data = nil
  @resources = nil
  @resources_data = nil
  @request_id = 0
  @mutex = Monitor.new
  @connection_established = false
  @initialized = false
  @http_conn = nil
  @session_id = nil
  @last_event_id = nil
  @sse_retry_ms = nil
  @pending_stream_responses = {}
  @response_post_count = 0
  # Saturation bookkeeping for the response-POST budget
  @dropped_response_posts = 0
  @last_saturation_log_at = nil
  @oauth_provider = opts[:oauth_provider]

  # SSE events connection state
  @events_connection = nil
  @events_thread = nil
  @buffer = +'' # Buffer for partial SSE event data
  # How much of @buffer has already been searched for an event terminator
  @buffer_scanned = 0
  @elicitation_request_callback = nil # MCP 2025-06-18
  @roots_list_request_callback = nil # MCP 2025-06-18
  @sampling_request_callback = nil # MCP 2025-11-25
end

Instance Attribute Details

#base_url ⇒ String (readonly)

Returns The base URL of the MCP server.

Returns:

  • (String) —

    The base URL of the MCP server



93
94
95
# File 'lib/mcp_client/server_streamable_http.rb', line 93

def base_url
  @base_url
end

#capabilities ⇒ Hash? (readonly)

Server capabilities from initialize response

Returns:

  • (Hash, nil) —

    Server capabilities



101
102
103
# File 'lib/mcp_client/server_streamable_http.rb', line 101

def capabilities
  @capabilities
end

#endpoint ⇒ String (readonly)

Returns The JSON-RPC endpoint path.

Returns:

  • (String) —

    The JSON-RPC endpoint path



93
# File 'lib/mcp_client/server_streamable_http.rb', line 93

attr_reader :base_url, :endpoint, :tools

#server_info ⇒ Hash? (readonly)

Server information from initialize response

Returns:

  • (Hash, nil) —

    Server information



97
98
99
# File 'lib/mcp_client/server_streamable_http.rb', line 97

def server_info
  @server_info
end

#tools ⇒ Object (readonly)

Returns the value of attribute tools.



93
# File 'lib/mcp_client/server_streamable_http.rb', line 93

attr_reader :base_url, :endpoint, :tools

Instance Method Details

#apply_request_headers(req, request) ⇒ Object

Override apply_request_headers to add session and SSE headers for MCP protocol



579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
# File 'lib/mcp_client/server_streamable_http.rb', line 579

def apply_request_headers(req, request)
  super

  # Modern servers have no session; the base class added the 2026-07-28
  # request metadata headers.
  return if modern?

  # Add session and protocol version headers for non-initialize requests
  return unless request['method'] != 'initialize'

  if @session_id
    req.headers['Mcp-Session-Id'] = @session_id
    @logger.debug("Adding session header: Mcp-Session-Id: #{@session_id}")
  end

  return unless @protocol_version

  req.headers['Mcp-Protocol-Version'] = @protocol_version
  @logger.debug("Adding protocol version header: Mcp-Protocol-Version: #{@protocol_version}")

  # NOTE: Last-Event-ID is deliberately NOT sent on POSTs — per SEP-1699,
  # resumption is always via HTTP GET with Last-Event-ID.
end

#call_tool(tool_name, parameters) ⇒ Object

Call a tool with the given parameters

Parameters:

  • tool_name (String) —

    the name of the tool to call

  • parameters (Hash) —

    the parameters to pass to the tool

Returns:

  • (Object) —

    the result of the tool invocation (with string keys for backward compatibility)

Raises:



260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
# File 'lib/mcp_client/server_streamable_http.rb', line 260

def call_tool(tool_name, parameters)
  rpc_request('tools/call', build_named_request_params(tool_name, parameters))
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ValidationError
  # Re-raise connection/transport errors directly to match test expectations
  raise
rescue MCPClient::Errors::ServerError => e
  # 2026-07-28 protocol errors (typed -3202x, invalid result) carry
  # actionable data such as requiredCapabilities; keep them intact.
  raise if e.protocol_error?

  raise MCPClient::Errors::ToolCallError, "Error calling tool '#{tool_name}': #{e.message}"
rescue StandardError => e
  # For all other errors, wrap in ToolCallError
  raise MCPClient::Errors::ToolCallError, "Error calling tool '#{tool_name}': #{e.message}"
end

#call_tool_streaming(tool_name, parameters) ⇒ Enumerator

Stream tool call (default implementation returns single-value stream)

Parameters:

  • tool_name (String) —

    the name of the tool to call

  • parameters (Hash) —

    the parameters to pass to the tool

Returns:

  • (Enumerator) —

    stream of results



280
281
282
283
284
# File 'lib/mcp_client/server_streamable_http.rb', line 280

def call_tool_streaming(tool_name, parameters)
  Enumerator.new do |yielder|
    yielder << call_tool(tool_name, parameters)
  end
end

#cleanup ⇒ Object

Clean up the server connection Properly closes HTTP connections, stops threads, and clears cached state



635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
# File 'lib/mcp_client/server_streamable_http.rb', line 635

def cleanup
  # Only an established session ends a task's namespace; a sessionless
  # (MCP 2026-07-28) connection is merely closed, and a connection that
  # never came up has nothing to end — see #ending_session?.
  bump_session_epoch if ending_session?
  @mutex.synchronize do
    return unless @connection_established || @initialized

    @logger.info('Cleaning up Streamable HTTP connection')

    # Mark connection as closed to stop reconnection attempts
    @connection_established = false
    @initialized = false

    # Subscription streams (MCP 2026-07-28) end with the connection
    close_listen_streams

    # Attempt to terminate session before cleanup
    begin
      terminate_session if @session_id
    rescue StandardError => e
      @logger.warn("Failed to terminate session: #{e.message}")
    end

    # Stop events thread gracefully
    if @events_thread&.alive?
      @logger.debug('Stopping events thread...')
      @events_thread.kill
      @events_thread.join(THREAD_JOIN_TIMEOUT)
    end
    @events_thread = nil

    # Clear connections and state
    @http_conn = nil
    @events_connection = nil
    @session_id = nil
    @last_event_id = nil
    @sse_retry_ms = nil
    @pending_stream_responses.each_value(&:close)
    @pending_stream_responses.clear

    # Clear cached data
    @tools = nil
    @tools_data = nil
    @prompts = nil
    @prompts_data = nil
    @resources = nil
    @resources_data = nil
    @resources_result = nil
    @templates_result = nil
    @buffer = +''
    @buffer_scanned = 0

    @logger.info('Cleanup completed')
  end
  # Cached results and their hints belong to the connection (and its
  # authorization context) that was just torn down; outside @mutex, as
  # the cache has its own lock.
  clear_result_cache
ensure
  # Everything this transport left on this thread — the notes of the
  # entries it served and recorded, the credentials, parameters and
  # metadata of its requests — describes a slice that will never be
  # tagged and a request that will never be made. Dropped after the
  # session was terminated, never before: the DELETE that terminates it
  # is a request of this transport's own, and the recorder on its
  # connection would put its Authorization fingerprint straight back on
  # this thread.
  forget_transport_thread_state
end

#complete(ref:, argument:, context: nil) ⇒ Hash

Request completion suggestions from the server (MCP 2025-06-18)

Parameters:

  • ref (Hash) —

    reference object (e.g., { 'type' => 'ref/prompt', 'name' => 'prompt_name' })

  • argument (Hash) —

    the argument being completed (e.g., { 'name' => 'arg_name', 'value' => 'partial' })

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

    optional context for the completion (MCP 2025-11-25)

Returns:

  • (Hash) —

    completion result with 'values', optional 'total', and 'hasMore' fields

Raises:



292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
# File 'lib/mcp_client/server_streamable_http.rb', line 292

def complete(ref:, argument:, context: nil)
  ensure_connected
  require_capability!('completions', method: 'completion/complete')
  params = { ref: ref, argument: argument }
  params[:context] = context if context
  result = require_complete_result!(rpc_request('completion/complete', params), 'completion/complete')
  result['completion'] || { 'values' => [] }
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError,
       MCPClient::Errors::CapabilityError
  raise
rescue MCPClient::Errors::ServerError => e
  # 2026-07-28 protocol errors (typed -3202x, invalid result) carry
  # actionable data such as requiredCapabilities; keep them intact.
  raise if e.protocol_error?

  raise MCPClient::Errors::ServerError, "Error requesting completion: #{e.message}"
rescue StandardError => e
  raise MCPClient::Errors::ServerError, "Error requesting completion: #{e.message}"
end

#connect ⇒ Boolean

Connect to the MCP server over Streamable HTTP

Returns:

  • (Boolean) —

    true if connection was successful

Raises:



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
# File 'lib/mcp_client/server_streamable_http.rb', line 189

def connect
  # Serialized: concurrent first requests must not each run the probe
  # and possibly settle on different eras (the monitor is reentrant, so
  # the request plumbing inside may take @mutex again).
  @mutex.synchronize do
    return true if @connection_established

    begin
      @connection_established = false
      @initialized = false

      # Test connectivity with a simple HTTP request
      test_connection

      # Establish the protocol era: server/discover for a modern server,
      # the initialize handshake for a legacy one.
      negotiate_protocol

      # Long-lived GET stream for server events: legacy only. MCP
      # 2026-07-28 removed the GET endpoint; change notifications arrive
      # on subscriptions/listen streams instead.
      start_events_connection unless modern?

      @connection_established = true
      @initialized = true

      true
    rescue MCPClient::Errors::ConnectionError => e
      cleanup
      raise e
    rescue StandardError => e
      cleanup
      raise MCPClient::Errors::ConnectionError, "Failed to connect to MCP server at #{@base_url}: #{e.message}"
    end
  end
end

#fetch_prompts_list ⇒ Array<MCPClient::Prompt>

Fetch and cache the prompt list.

Returns:



366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
# File 'lib/mcp_client/server_streamable_http.rb', line 366

def fetch_prompts_list
  ensure_connected

  prompts = request_prompts_list.map { |prompt_data| MCPClient::Prompt.from_json(prompt_data, server: self) }
  @mutex.synchronize do
    @prompts = attach_list_value(:prompts, prompts) ? prompts : nil
  end

  # This request's own list, never a re-read of @prompts (another
  # request may have stored its list in between).
  prompts
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError
  # Re-raise these errors directly
  raise
rescue StandardError => e
  raise MCPClient::Errors::PromptGetError, "Error listing prompts: #{e.message}"
end

#fetch_resources_list(cursor) ⇒ Hash

Fetch one page of resources/list, caching the first page.

Parameters:

  • cursor (String, nil)

Returns:

  • (Hash)


432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
# File 'lib/mcp_client/server_streamable_http.rb', line 432

def fetch_resources_list(cursor)
  params = {}
  params['cursor'] = cursor if cursor
  epoch = cache_epoch(:resources)
  answer = fetching_list_page(:resources, cursor) { rpc_request('resources/list', params) }
  result = require_complete_result!(answer, 'resources/list')
  record_cache_hint(:resources, result, epoch: epoch) unless cursor

  resources = (result['resources'] || []).map do |resource_data|
    MCPClient::Resource.from_json(resource_data, server: self)
  end

  resources_result = { 'resources' => resources, 'nextCursor' => result['nextCursor'] }

  # A list invalidated while in flight is returned but not cached.
  @mutex.synchronize do
    unless cursor
      @resources_result = attach_list_value(:resources, resources_result) ? resources_result : nil
    end
  end

  resources_result
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError
  # Re-raise these errors directly
  raise
rescue StandardError => e
  raise MCPClient::Errors::ResourceReadError, "Error listing resources: #{e.message}"
end

#fetch_templates_list(cursor) ⇒ Hash

Fetch one page of resources/templates/list, caching the first page.

Parameters:

  • cursor (String, nil)

Returns:

  • (Hash)


512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
# File 'lib/mcp_client/server_streamable_http.rb', line 512

def fetch_templates_list(cursor)
  params = {}
  params['cursor'] = cursor if cursor
  epoch = cache_epoch(:templates)
  answer = fetching_list_page(:templates, cursor) { rpc_request('resources/templates/list', params) }
  result = require_complete_result!(answer, 'resources/templates/list')
  record_cache_hint(:templates, result, epoch: epoch) unless cursor

  templates = (result['resourceTemplates'] || []).map do |template_data|
    MCPClient::ResourceTemplate.from_json(template_data, server: self)
  end
  templates_result = { 'resourceTemplates' => templates, 'nextCursor' => result['nextCursor'] }

  @mutex.synchronize do
    unless cursor
      @templates_result = attach_list_value(:templates, templates_result) ? templates_result : nil
    end
  end

  templates_result
end

#get_prompt(prompt_name, parameters) ⇒ Object

Get a prompt with the given parameters

Parameters:

  • prompt_name (String) —

    the name of the prompt to get

  • parameters (Hash) —

    the parameters to pass to the prompt

Returns:

  • (Object) —

    the result of the prompt (with string keys for backward compatibility)

Raises:



389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
# File 'lib/mcp_client/server_streamable_http.rb', line 389

def get_prompt(prompt_name, parameters)
  rpc_request('prompts/get', build_named_request_params(prompt_name, parameters))
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError
  # Re-raise connection/transport errors directly
  raise
rescue MCPClient::Errors::ServerError => e
  # 2026-07-28 protocol errors (typed -3202x, invalid result) carry
  # actionable data such as requiredCapabilities; keep them intact.
  raise if e.protocol_error?

  raise MCPClient::Errors::PromptGetError, "Error getting prompt '#{prompt_name}': #{e.message}"
rescue StandardError => e
  # For all other errors, wrap in PromptGetError
  raise MCPClient::Errors::PromptGetError, "Error getting prompt '#{prompt_name}': #{e.message}"
end

#handle_successful_response(response, request) ⇒ Object

Override handle_successful_response to capture session ID



604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
# File 'lib/mcp_client/server_streamable_http.rb', line 604

def handle_successful_response(response, request)
  super

  # Capture session ID from initialize response with validation
  return unless request['method'] == 'initialize' && response.success?

  session_id = response.headers['mcp-session-id'] || response.headers['Mcp-Session-Id']
  if session_id
    if valid_session_id?(session_id)
      capture_session_id(session_id)
      @logger.debug("Captured session ID: #{@session_id}")
    else
      @logger.warn("Invalid session ID format received: #{session_id.inspect}")
    end
  else
    @logger.warn('No session ID found in initialize response headers')
  end
end

#list_prompts ⇒ Array<MCPClient::Prompt>

List all prompts available from the MCP server

Returns:

Raises:



349
350
351
352
353
354
355
356
357
358
359
360
361
362
# File 'lib/mcp_client/server_streamable_http.rb', line 349

def list_prompts
  cached = fresh_list_value(:prompts) { @mutex.synchronize { @prompts } }
  return cached if cached

  @mutex.synchronize { @prompts_data = nil }
  begin
    ensure_connected
    refetch_or_serve_stale(:prompts, stale_list_entry(:prompts)) { fetch_prompts_list }
  rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError
    raise
  rescue StandardError => e
    raise MCPClient::Errors::PromptGetError, "Error listing prompts: #{e.message}"
  end
end

#list_resource_templates(cursor: nil) ⇒ Hash

List all resource templates available from the MCP server

Parameters:

  • cursor (String, nil) (defaults to: nil) —

    optional cursor for pagination

Returns:

  • (Hash) —

    result containing resourceTemplates array and optional nextCursor

Raises:



485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
# File 'lib/mcp_client/server_streamable_http.rb', line 485

def list_resource_templates(cursor: nil)
  # Only a list the server itself bounded is served from here: a
  # positive ttlMs means no second request, while a list with no hint
  # (a 2025-11-25 server) is asked for again, as it was before this
  # transport cached anything (MCP 2026-07-28 caching).
  cached = cursor ? nil : hinted_list_value(:templates)
  return cached if cached

  begin
    ensure_connected
    unless cursor
      return refetch_or_serve_stale(:templates, stale_list_entry(:templates)) do
        fetch_templates_list(nil)
      end
    end

    fetch_templates_list(cursor)
  rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError
    raise
  rescue StandardError => e
    raise MCPClient::Errors::ResourceReadError, "Error listing resource templates: #{e.message}"
  end
end

#list_resources(cursor: nil) ⇒ Hash

List all resources available from the MCP server

Parameters:

  • cursor (String, nil) (defaults to: nil) —

    optional cursor for pagination

Returns:

  • (Hash) —

    result containing resources array and optional nextCursor

Raises:



409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
# File 'lib/mcp_client/server_streamable_http.rb', line 409

def list_resources(cursor: nil)
  cached = cursor ? nil : fresh_list_value(:resources) { @mutex.synchronize { @resources_result } }
  return cached if cached

  begin
    ensure_connected
    unless cursor
      return refetch_or_serve_stale(:resources, stale_list_entry(:resources)) do
        fetch_resources_list(nil)
      end
    end

    fetch_resources_list(cursor)
  rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError
    raise
  rescue StandardError => e
    raise MCPClient::Errors::ResourceReadError, "Error listing resources: #{e.message}"
  end
end

#list_tools ⇒ Array<MCPClient::Tool>

List all tools available from the MCP server

Returns:

Raises:



231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
# File 'lib/mcp_client/server_streamable_http.rb', line 231

def list_tools
  # MCP 2026-07-28 caching: a cached list is served only while fresh, and
  # only from the entry that carries its hint.
  cached = fresh_list_value(:tools) { @mutex.synchronize { @tools } }
  return cached if cached

  # Stale: the raw page cache must go too, or the re-fetch would be
  # answered from memory.
  @mutex.synchronize { @tools_data = nil }
  begin
    ensure_connected

    refetch_or_serve_stale(:tools, stale_list_entry(:tools)) { fetch_tools_list }
  rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError
    # Re-raise these errors directly
    raise
  rescue StandardError => e
    raise MCPClient::Errors::ToolCallError, "Error listing tools: #{e.message}"
  end
end

#log_level=(level) ⇒ Hash

Deprecated.

Logging is deprecated since MCP 2026-07-28 (SEP-2577); earliest removal is the first revision released on or after 2027-07-28. Have the server log to stderr (stdio) or use OpenTelemetry instead.

Set the logging level on the server (MCP 2025-06-18)

Parameters:

  • level (String) —

    the log level ('debug', 'info', 'notice', 'warning', 'error', 'critical', 'alert', 'emergency')

Returns:

  • (Hash) —

    empty result on success

Raises:



321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
# File 'lib/mcp_client/server_streamable_http.rb', line 321

def log_level=(level)
  MCPClient::Deprecations.warn(:logging, @logger)
  ensure_connected
  # MCP 2026-07-28 removed logging/setLevel: the level travels per request
  # in _meta["io.modelcontextprotocol/logLevel"].
  if modern?
    @log_level = validate_log_level!(level)
    return
  end

  require_capability!('logging', method: 'logging/setLevel')
  rpc_request('logging/setLevel', { level: level })
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError,
       MCPClient::Errors::CapabilityError, ArgumentError
  raise
rescue MCPClient::Errors::ServerError => e
  # 2026-07-28 protocol errors (typed -3202x, invalid result) carry
  # actionable data such as requiredCapabilities; keep them intact.
  raise if e.protocol_error?

  raise MCPClient::Errors::ServerError, "Error setting log level: #{e.message}"
rescue StandardError => e
  raise MCPClient::Errors::ServerError, "Error setting log level: #{e.message}"
end

#on_elicitation_request(&block) ⇒ void

This method returns an undefined value.

Register a callback for elicitation requests (MCP 2025-06-18)

Parameters:

  • block (Proc) —

    callback that receives (request_id, params) and returns response hash



709
710
711
# File 'lib/mcp_client/server_streamable_http.rb', line 709

def on_elicitation_request(&block)
  @elicitation_request_callback = block
end

#on_roots_list_request(&block) ⇒ void

Deprecated.

Roots is deprecated since MCP 2026-07-28 (SEP-2577); earliest removal is the first revision released on or after 2027-07-28. Registering a handler is not itself use of Roots — a handler that answers with no root exposes nothing deprecated — but a handler that answers with a root adopts the deprecated feature and raises the notice. Pass directories or files through tool parameters, resource URIs or server configuration instead.

This method returns an undefined value.

Register a callback for roots/list requests (MCP 2025-06-18)

Parameters:

  • block (Proc) —

    callback that receives (request_id, params) and returns response hash



723
724
725
# File 'lib/mcp_client/server_streamable_http.rb', line 723

def on_roots_list_request(&block)
  @roots_list_request_callback = block
end

#on_sampling_request(&block) ⇒ void

Deprecated.

Sampling is deprecated since MCP 2026-07-28 (SEP-2577); earliest removal is the first revision released on or after 2027-07-28. Integrate directly with the LLM provider API instead of serving sampling/createMessage.

This method returns an undefined value.

Register a callback for sampling requests (MCP 2025-11-25)

Parameters:

  • block (Proc) —

    callback that receives (request_id, params) and returns response hash



735
736
737
# File 'lib/mcp_client/server_streamable_http.rb', line 735

def on_sampling_request(&block)
  @sampling_request_callback = block
end

#read_resource(uri) ⇒ Array<MCPClient::ResourceContent>

Read a resource by its URI

Parameters:

  • uri (String) —

    the URI of the resource to read

Returns:

Raises:



465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
# File 'lib/mcp_client/server_streamable_http.rb', line 465

def read_resource(uri)
  ensure_connected
  read_resource_with_cache(uri) { |sent| rpc_request('resources/read', { uri: sent }) }
rescue MCPClient::Errors::ServerError => e
  raise if e.protocol_error?
  raise resource_not_found_error(uri, e) if resource_not_found_response?(e)

  raise MCPClient::Errors::ResourceReadError, "Error reading resource '#{uri}': #{e.message}"
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError
  # Re-raise connection/transport errors directly
  raise
rescue StandardError => e
  # For all other errors, wrap in ResourceReadError
  raise MCPClient::Errors::ResourceReadError, "Error reading resource '#{uri}': #{e.message}"
end

#subscribe_resource(uri) ⇒ Boolean

Subscribe to resource updates

Parameters:

  • uri (String) —

    the URI of the resource to subscribe to

Returns:

  • (Boolean) —

    true if subscription successful

Raises:



538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
# File 'lib/mcp_client/server_streamable_http.rb', line 538

def subscribe_resource(uri)
  ensure_connected
  require_capability!('resources', 'subscribe', method: 'resources/subscribe')
  # MCP 2026-07-28 replaced resources/subscribe with a subscriptions/listen
  # stream carrying resourceSubscriptions.
  if modern?
    subscribe_resource_via_listen(uri)
    return true
  end

  rpc_request('resources/subscribe', { uri: uri })
  true
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError,
       MCPClient::Errors::CapabilityError
  raise
rescue StandardError => e
  raise MCPClient::Errors::ResourceReadError, "Error subscribing to resource '#{uri}': #{e.message}"
end

#terminate_session ⇒ Boolean

Terminate the current session (if any)

Returns:

  • (Boolean) —

    true if termination was successful or no session exists



625
626
627
628
629
630
631
# File 'lib/mcp_client/server_streamable_http.rb', line 625

def terminate_session
  @mutex.synchronize do
    return true unless @session_id

    super
  end
end

#unsubscribe_resource(uri) ⇒ Boolean

Unsubscribe from resource updates

Parameters:

  • uri (String) —

    the URI of the resource to unsubscribe from

Returns:

  • (Boolean) —

    true if unsubscription successful

Raises:



561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
# File 'lib/mcp_client/server_streamable_http.rb', line 561

def unsubscribe_resource(uri)
  ensure_connected
  require_capability!('resources', 'subscribe', method: 'resources/unsubscribe')
  if modern?
    unsubscribe_resource_via_listen(uri)
    return true
  end

  rpc_request('resources/unsubscribe', { uri: uri })
  true
rescue MCPClient::Errors::ConnectionError, MCPClient::Errors::TransportError, MCPClient::Errors::ServerError,
       MCPClient::Errors::CapabilityError
  raise
rescue StandardError => e
  raise MCPClient::Errors::ResourceReadError, "Error unsubscribing from resource '#{uri}': #{e.message}"
end