Module: MCPClient::ServerStdio::JsonRpcTransport
- Includes:
- JsonRpcCommon
- Included in:
- MCPClient::ServerStdio
- Defined in:
- lib/mcp_client/server_stdio/json_rpc_transport.rb
Overview
JSON-RPC request/notification plumbing for stdio transport
Constant Summary
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 MCPClient::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 MCPClient::SubscriptionSupport
MCPClient::SubscriptionSupport::CONTROL_NOTIFICATIONS, MCPClient::SubscriptionSupport::DEFAULT_ACK_TIMEOUT
Constants included from JsonRpcCommon::ErrorBodies
JsonRpcCommon::ErrorBodies::MAX_ERROR_BODY_BYTES
Instance Attribute Summary
Attributes included from JsonRpcCommon
#request_meta, #send_client_info
Instance Method Summary collapse
-
#build_registered_request(method, params, req_id) ⇒ Hash
Build a JSON-RPC request under an id #next_id has already registered as outstanding.
-
#call_tool_streaming(tool_name, parameters) ⇒ Enumerator
Stream tool call fallback for stdio transport (yields single result).
-
#crash_looping?(carrier) ⇒ Boolean
Whether the process these subscriptions were last given to died too soon after receiving them for another process to be worth spawning.
-
#declared_protocol_version(req) ⇒ String?
The protocol version a built request declares in its
_meta. -
#defer_reestablished_attempt(subscription, id, error) ⇒ void
Put a subscription whose hand-over could not be written back on the queue the next session drains, and say so.
-
#discover_result?(result) ⇒ Boolean
Whether it has the DiscoverResult shape.
-
#dropped_requests ⇒ Set<Integer>
Ids of requests a teardown found outstanding: recorded under @mutex by cleanup, consumed by their waiters (see #wait_response).
-
#enqueue_reconnecting_locked(subscriptions) ⇒ Array<MCPClient::Subscription>
The whole queue afterwards.
-
#enqueue_reconnecting_subscriptions(subscriptions) ⇒ Array<MCPClient::Subscription>
Put subscriptions on the queue the next session drains, at most once each.
-
#ensure_initialized ⇒ void
Ensure the server process is started and initialized (handshake).
- #ensure_session_ready ⇒ void
-
#fail_open_attempt(subscription, id, error, io: nil) ⇒ void
Undo the listen attempt that just failed — but only when the subscription is this attempt's to undo.
-
#fail_reconnecting_subscriptions(error) ⇒ void
End the subscriptions waiting for a process that is not coming back.
- #fail_subscriptions(pending, error) ⇒ void
-
#fail_superseded_attempt(subscription, id, error) ⇒ void
A failure the subscription has already moved on from: harmless while the stream that replaced it stands, and the caller's answer when it does not.
-
#fall_forward_to_modern(error) ⇒ Hash
Resume the modern path after a fallback handshake was refused by a modern server: the rejection settles the era, so server/discover is re-issued with a version the server named and initialize is never sent again.
-
#fall_forward_to_modern?(error) ⇒ Boolean
Whether a rejected initialize handshake should send this connection back to the modern path.
-
#hand_over_to_established_process ⇒ void
Hand the queue to a process that is already established.
-
#identifies_modern_server?(msg) ⇒ Boolean
Whether a response, as it comes off the wire, could only have been written by a modern server: a result carrying a 2026-07-28 marker, or one of the spec-defined modern errors in its mandated shape.
-
#interpret_discover_answer(res) ⇒ Hash
Turn the probe's answer into a DiscoverResult, or into the failure that says what kind of server sent it.
-
#invalid_discover_answer(modern_answer, message) ⇒ StandardError
The failure the probe should propagate.
-
#legacy_after_probe(error) ⇒ void
Not a recognized modern error: a legacy server (or one that never answered).
-
#live_process? ⇒ Boolean
Whether a subprocess is connected and still running.
-
#modern_discover_answer?(res) ⇒ Boolean
Whether its result could only have come from a 2026-07-28 server.
-
#negotiate_protocol ⇒ void
Establish the server's protocol era (MCP 2026-07-28 basic/transports/stdio "Backward Compatibility"): probe with server/discover unless configured legacy-only, and fall back to the initialize handshake when the probe shows a legacy server.
-
#next_id ⇒ Integer
Generate a new unique request ID and mark it as awaiting a response.
-
#open_subscription(subscription) ⇒ void
Send the subscriptions/listen request for a subscription (a fresh id each time it is opened or re-opened).
-
#perform_discover ⇒ Hash
Send server/discover and apply the DiscoverResult.
-
#perform_initialize ⇒ void
Handshake: send initialize request and initialized notification.
-
#probe_modern_server ⇒ Boolean
Send the server/discover probe with this client's preferred modern version.
-
#queue_subscriptions_of_ended_process(subscriptions) ⇒ void
Hand the subscriptions of a process that is being torn down to the next one, forgetting the listen ids that process was holding: nothing written to it is outstanding any more, and none of those ids may be cancelled on the process that replaces it (see MCPClient::Subscription#record_outstanding_listen).
-
#reconnecting_mutex ⇒ Mutex
Guards the queue of subscriptions waiting for a process (created by #initialize, so no two threads ever race to make it).
-
#reconnecting_subscriptions ⇒ Array<MCPClient::Subscription>
The subscriptions on the queue that a process could still be re-sent to.
-
#release_retired_transport ⇒ void
Discard a transport whose subprocess exited after a successful handshake, so the next negotiation starts from a clean slate rather than on top of the dead process's handles.
-
#release_transport ⇒ void
Tear down a transport that is not going to be used again — a handshake that never completed, or a subprocess that exited under a completed one.
-
#reopen_refusal(carrier) ⇒ StandardError?
Why this session must not be given the open subscriptions, if it must not (see #reopen_subscriptions).
-
#reopen_subscriptions(session = @session) ⇒ void
After the process was re-established, re-send subscriptions/listen for every subscription the host still holds open ("the server holds no subscription state across reconnections").
-
#restart_for_open_subscriptions ⇒ void
Restart the process the reader just watched exit, for the subscriptions the host still holds.
-
#retry_discover_with_advertised_version(error) ⇒ void
After UnsupportedProtocolVersionError, pick a mutually supported version from the error's advertised list and re-issue the probe.
-
#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 wait for result.
-
#send_cancellation_notification(request_id) ⇒ void
Best-effort notifications/cancelled for a request the client stopped waiting on.
-
#send_if_current(req, generation) ⇒ Hash?
Write a registered request, unless the transport it was registered on has been replaced since.
-
#send_on_current_transport(method, params) {|req| ... } ⇒ Array(Integer, Hash)
Register, build and write a request on the transport that is current when it is written.
-
#send_request(req, generation = nil, io: @stdin) ⇒ Symbol
Send a JSON-RPC request.
-
#send_request_and_wait(method, params, timeout) {|version| ... } ⇒ Object
One request/response exchange with its own JSON-RPC id.
-
#subscription_failure(error) ⇒ MCPClient::Errors::MCPError
The error a subscription ends with.
-
#take_reconnecting_subscriptions ⇒ Array<MCPClient::Subscription>
Take the subscriptions waiting for a process, leaving the queue empty.
-
#transport_retired? ⇒ Boolean
Whether the subprocess behind the handshake exited.
-
#wait_response(id, timeout: nil) ⇒ Hash
Wait for a response with the given request ID.
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 MCPClient::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 MCPClient::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 Method Details
#build_registered_request(method, params, req_id) ⇒ Hash
Build a JSON-RPC request under an id #next_id has already registered as outstanding. Building can fail — the host's request_meta provider is evaluated here and may raise — and a request that was never built is never sent and never answered, so its marker has to go with it; otherwise every such failure leaks an entry into @awaiting.
835 836 837 838 839 840 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 835 def build_registered_request(method, params, req_id) build_jsonrpc_request(method, params, req_id) rescue StandardError @mutex.synchronize { @awaiting.delete(req_id) } raise end |
#call_tool_streaming(tool_name, parameters) ⇒ Enumerator
Stream tool call fallback for stdio transport (yields single result)
961 962 963 964 965 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 961 def call_tool_streaming(tool_name, parameters) Enumerator.new do |yielder| yielder << call_tool(tool_name, parameters) end end |
#crash_looping?(carrier) ⇒ Boolean
Whether the process these subscriptions were last given to died too soon after receiving them for another process to be worth spawning. The two moments are stamped on that process's own record, by its own lifecycle, so this answer cannot be spoiled by whatever another thread is doing to another process.
396 397 398 399 400 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 396 def crash_looping?(carrier) carrier&.( MCPClient::ServerStdio::SUBSCRIPTION_RESTART_MIN_INTERVAL ) || false end |
#declared_protocol_version(req) ⇒ String?
The protocol version a built request declares in its _meta. Read
back from the request rather than from the transport: building it
evaluates the host's metadata provider, during which a concurrent
request may settle the transport on a different version.
1093 1094 1095 1096 1097 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 1093 def declared_protocol_version(req) params = req['params'] = params.is_a?(Hash) ? params['_meta'] : nil .is_a?(Hash) ? [JsonRpcCommon::META_PROTOCOL_VERSION] : nil end |
#defer_reestablished_attempt(subscription, id, error) ⇒ void
This method returns an undefined value.
Put a subscription whose hand-over could not be written back on the queue the next session drains, and say so.
It has just been taken off that queue by #reopen_subscriptions and
out of the registry above, so leaving it alone would strand it: no
session would re-send it and no cleanup would find it again. The
queue is where a subscription waiting for a process belongs, and the
process that could not be written to is on its way out — its reader
reaches EOF and restarts, and the crash-loop bound then decides
whether another one is worth spawning.
The subscription is :reconnecting again already — that transition
was decided with the verdict, under the one lock
(MCPClient::Subscription#fail_attempt).
233 234 235 236 237 238 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 233 def defer_reestablished_attempt(subscription, id, error) enqueue_reconnecting_subscriptions([subscription]) @logger.debug("#{id ? "subscriptions/listen #{id}" : 'a subscriptions/listen request'} failed while the " \ "subscription was being handed to a new process (#{error.}); it will be re-sent to " \ 'the next one') end |
#discover_result?(result) ⇒ Boolean
Returns whether it has the DiscoverResult shape.
1029 1030 1031 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 1029 def discover_result?(result) result.is_a?(Hash) && result['supportedVersions'].is_a?(Array) end |
#dropped_requests ⇒ Set<Integer>
Ids of requests a teardown found outstanding: recorded under @mutex by cleanup, consumed by their waiters (see #wait_response).
338 339 340 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 338 def dropped_requests @dropped_requests ||= Set.new end |
#enqueue_reconnecting_locked(subscriptions) ⇒ Array<MCPClient::Subscription>
Returns the whole queue afterwards.
445 446 447 448 449 450 451 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 445 def enqueue_reconnecting_locked(subscriptions) queue = (@reconnecting_subscriptions ||= []) subscriptions.each do |subscription| queue << subscription unless queue.any? { |queued| queued.equal?(subscription) } end queue.dup end |
#enqueue_reconnecting_subscriptions(subscriptions) ⇒ Array<MCPClient::Subscription>
Put subscriptions on the queue the next session drains, at most once each.
Two paths write to this queue and they overlap: MCPClient::ServerStdio#cleanup moves the
open subscriptions onto it, and #defer_reestablished_attempt puts
back a hand-over whose write failed — and the second happens inside
the window the first leaves between taking the registry snapshot and
writing it here. An unguarded Array concated by one and <<ed by
the other is undefined in MRI: the same window can lose the entry, and
a subscription no session re-sends and no cleanup finds again is a
stream the spec says MUST be re-established, stranded with the host
never told. Scanning that Array with equal? while another thread
grows it does not make the append safe either — it only decided,
unreliably, whether to make a second one. So both paths come through
here, under the one lock that also guards the take, and membership is
by identity: a handle appears on the queue once, and one hand-over
goes out for it.
421 422 423 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 421 def enqueue_reconnecting_subscriptions(subscriptions) reconnecting_mutex.synchronize { enqueue_reconnecting_locked(subscriptions) } end |
#ensure_initialized ⇒ void
This method returns an undefined value.
Ensure the server process is started and initialized (handshake)
14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 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 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 14 def ensure_initialized # Read the retirement flag FIRST. A restart clears @initialized and # then the flag; a thread reading them the other way round could see # a stale handshake next to a cleared flag and skip the lock, then # register a request against a transport that is being replaced. return if !transport_retired? && @initialized @init_lock.synchronize do # The subprocess behind a completed handshake exited under it: # release its pipes and reader threads before connect overwrites # the handles, then negotiate again against the fresh process. release_retired_transport if transport_retired? return if @initialized begin # Ordinary requests are refused while the replacement is being # negotiated (see #send_request): the restart clears the # retirement before it negotiates, so the generation alone would # judge a half-restarted transport current. @transport_lock.synchronize { @negotiating = true } # A process the host connected explicitly is negotiated, not # replaced: spawning again would orphan it. spawned = !live_process? connect if spawned # The record of the process that is now the session. Everything the # crash-loop bound needs is written on it, by its own lifecycle: see # {MCPClient::ServerStdio::ChildSession}. A process spawned here gets # a fresh record; one the host connected explicitly keeps the record # it already has, or gets its first. session = spawned || @session.nil? ? (@session = MCPClient::ServerStdio::ChildSession.new) : @session start_reader unless @reader_thread&.alive? start_stderr_reader unless @stderr_thread&.alive? negotiate_protocol rescue StandardError # A failed negotiation must not leave the subprocess, its pipes # and its reader threads behind. @initialized stays false, so the # next request runs connect again and overwrites @stdin/@stdout/ # @wait_thread — putting the first process permanently out of # cleanup's reach. release_transport raise ensure @transport_lock.synchronize { @negotiating = false } end @initialized = true # The record is passed rather than read back: a cleanup on another # thread can retire @session while this one is still handing the # subscriptions to the process it established, and the crash-loop # bound must be stamped on the process that actually received them. reopen_subscriptions(session) end end |
#ensure_session_ready ⇒ void
This method returns an undefined value.
69 70 71 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 69 def ensure_session_ready ensure_initialized end |
#fail_open_attempt(subscription, id, error, io: nil) ⇒ void
This method returns an undefined value.
Undo the listen attempt that just failed — but only when the subscription is this attempt's to undo.
A write can block long enough for the child to exit and for the restart that follows to take the subscription over. Two things can have happened by then, and neither is this attempt's to tear down:
- the restart already re-opened it under a newer id. The registry is keyed by listen id, so unregistering "the subscription" would delete the new registration and finishing it would close a stream the fresh process is serving. Naming the id the write went out with keeps this attempt to its own. It only says so in the log — unless that newer stream has itself already failed, in which case the caller must be told rather than handed a closed handle with no explanation.
- the subscription is one a session is being handed (MCPClient::Subscription#reestablishing?): it is a stream the spec requires to be re-sent, not one this caller asked for, so a write that failed because the process was gone (stdin closed under it, an EPIPE to a child that exited on sight, or a nested restart holding the init lock) leaves it for the next session instead of ending it. The question is asked of the subscription rather than of its state: taking the new listen id has already moved it from :reconnecting to :pending by the time the write raises, so the state says "being opened" for the very hand-over that is failing.
Which of those it is, and the transition that follows, are decided in one step on the subscription (MCPClient::Subscription#fail_attempt). Asked first and acted on afterwards, the answer went stale in between: a restart re-opened the subscription under a newer id and had it acknowledged after the ownership check had passed, and this attempt then finished the healthy replacement. Only the registration comes first — it is scoped to this attempt's id, so a newer one is never touched by it.
An attempt that failed before the subscription took its id at all
(the request could not be built) arrives with no id: nothing was
registered for it, and no newer attempt can have superseded it, so
the failure is the caller's — or the next session's — to hear about.
A failure that is the caller's own abandons the request the write may
have put on the pipe, and abandoning a request on stdio is a
cancellation naming its id (basic/transports/stdio "Cancellation"):
left unnamed, the server was serving listen(n) for a subscription
that had ended with the error, and nothing this client held could
cancel it any more. It is sent best effort, to the pipe the request
went to. A superseded attempt's request went to a process the restart
has replaced, and a deferred hand-over's to one on its way out —
neither is cancelled on a pipe that is gone.
197 198 199 200 201 202 203 204 205 206 207 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 197 def fail_open_attempt(subscription, id, error, io: nil) unregister_subscription_id(subscription, id) if id failure = subscription_failure(error) case subscription.fail_attempt(id, failure) when :superseded then fail_superseded_attempt(subscription, id, error) when :deferred then defer_reestablished_attempt(subscription, id, error) else cancel_outstanding_listens(subscription, io: io) if io && id raise failure end end |
#fail_reconnecting_subscriptions(error) ⇒ void
This method returns an undefined value.
End the subscriptions waiting for a process that is not coming back.
539 540 541 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 539 def fail_reconnecting_subscriptions(error) fail_subscriptions(take_reconnecting_subscriptions, error) end |
#fail_subscriptions(pending, error) ⇒ void
This method returns an undefined value.
546 547 548 549 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 546 def fail_subscriptions(pending, error) failure = subscription_failure(error) pending.each { |subscription| subscription.finish(gracefully: false, error: failure) } end |
#fail_superseded_attempt(subscription, id, error) ⇒ void
This method returns an undefined value.
A failure the subscription has already moved on from: harmless while the stream that replaced it stands, and the caller's answer when it does not.
248 249 250 251 252 253 254 255 256 257 258 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 248 def fail_superseded_attempt(subscription, id, error) replacement_error = subscription.closed? ? subscription.error : nil if replacement_error @logger.warn("subscriptions/listen #{id} failed (#{error.}) and the stream that replaced it " \ "failed too: #{sanitize_log_text(replacement_error.)}") raise replacement_error end @logger.debug("subscriptions/listen #{id} failed after the subscription was re-opened " \ "(#{error.}); the newer stream stands") end |
#fall_forward_to_modern(error) ⇒ Hash
Resume the modern path after a fallback handshake was refused by a modern server: the rejection settles the era, so server/discover is re-issued with a version the server named and initialize is never sent again.
805 806 807 808 809 810 811 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 805 def fall_forward_to_modern(error) version = select_protocol_version(error.supported) @logger.info('The server refused the initialize handshake and supports ' \ "#{error.supported.join(', ')}; it is a modern server — retrying server/discover with #{version}") @protocol_version = version perform_discover end |
#fall_forward_to_modern?(error) ⇒ Boolean
Whether a rejected initialize handshake should send this connection back to the modern path. Only a well-formed rejection counts — a bare -32022 from a legacy endpoint identifies nothing — and only one that names a version this client speaks, since the retry has to declare one. A host that configured protocol: :legacy asked for the 2025-11-25 handshake and gets the error instead.
793 794 795 796 797 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 793 def fall_forward_to_modern?(error) return false if @protocol_mode == :legacy error.modern_protocol_error? && !select_protocol_version(error.supported).nil? end |
#hand_over_to_established_process ⇒ void
This method returns an undefined value.
Hand the queue to a process that is already established.
#ensure_initialized re-sends the queue itself, but only when it negotiated the process: a host request that observed the exit while this teardown was still parking its subscriptions established the replacement, found the queue empty and returned, and the restart above then took the initialized fast path — leaving the subscriptions parked on a queue nothing was going to read. Whichever of the two finishes last drains what it finds here, so the re-send MCP 2026-07-28 basic/patterns/subscriptions requires after a stdio reconnect happens however the two threads interleave.
Under the initialization lock, like every other hand-over: the process must not be replaced underneath the writes, and a queue taken while a negotiation is in flight would be re-sent to the process that negotiation is replacing.
528 529 530 531 532 533 534 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 528 def hand_over_to_established_process @init_lock.synchronize do next unless @initialized reopen_subscriptions end end |
#identifies_modern_server?(msg) ⇒ Boolean
Whether a response, as it comes off the wire, could only have been written by a modern server: a result carrying a 2026-07-28 marker, or one of the spec-defined modern errors in its mandated shape.
731 732 733 734 735 736 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 731 def identifies_modern_server?(msg) return true if modern_discover_answer?(msg) return false unless msg.key?('error') MCPClient::Errors::ServerError.from_jsonrpc(msg['error']).modern_protocol_error? end |
#interpret_discover_answer(res) ⇒ Hash
Turn the probe's answer into a DiscoverResult, or into the failure that says what kind of server sent it.
The stdio fallback rule is keyed to the probe being answered with an
error, or not answered at all — never to a result. resultType does
not exist before 2026-07-28, so a result carrying it came from a
modern server even when the rest of it is unusable: treating that as
a legacy answer would pin a dual-era server to the 2025-11-25
handshake for the life of the process, and would make a modern-only
server fail to connect after it had already answered the probe. Such
an answer therefore fails the negotiation instead of falling back.
A result with no 2026-07-28 marker at all is still a legacy answer: a
permissive server answering an unknown method with some object.
700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 700 def interpret_discover_answer(res) modern_answer = modern_discover_answer?(res) # Named so the identity of an answer that fails to validate below is # not recorded: apply_discover_result records it once the whole # result has validated. result = process_jsonrpc_response(res, method: 'server/discover') reject_input_required_discover!(result) reject_task_result_discover!(result) unless discover_result?(result) raise invalid_discover_answer(modern_answer, 'answered without a DiscoverResult') end apply_discover_result(result) rescue MCPClient::Errors::InvalidResultError, MCPClient::Errors::InputRequiredError => e raise invalid_discover_answer(modern_answer, "answered without a DiscoverResult (#{e.})") end |
#invalid_discover_answer(modern_answer, message) ⇒ StandardError
Returns the failure the probe should propagate.
741 742 743 744 745 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 741 def invalid_discover_answer(modern_answer, ) return MCPClient::Errors::ServerError.new("server/discover was #{}") unless modern_answer MCPClient::Errors::ConnectionError.new("Server is modern but incompatible: server/discover was #{}") end |
#legacy_after_probe(error) ⇒ void
This method returns an undefined value.
Not a recognized modern error: a legacy server (or one that never answered).
636 637 638 639 640 641 642 643 644 645 646 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 636 def legacy_after_probe(error) e = error @protocol_version = nil if @protocol_mode == :modern raise MCPClient::Errors::ConnectionError, "Server did not answer server/discover as a modern MCP server (#{e.}); it is most likely " \ 'a legacy server expecting the initialize handshake. Use protocol: :auto or :legacy to allow that.' end @logger.debug("server/discover probe failed (#{e.class}); treating the server as legacy") end |
#live_process? ⇒ Boolean
Returns whether a subprocess is connected and still running.
329 330 331 332 333 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 329 def live_process? return false unless @stdin.respond_to?(:closed?) && !@stdin.closed? @wait_thread.respond_to?(:alive?) && @wait_thread.alive? end |
#modern_discover_answer?(res) ⇒ Boolean
Returns whether its result could only have come from a 2026-07-28 server.
719 720 721 722 723 724 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 719 def modern_discover_answer?(res) result = res.is_a?(Hash) ? res['result'] : nil return false unless result.is_a?(Hash) result.key?('resultType') || result.key?(:resultType) || discover_result?(result) end |
#negotiate_protocol ⇒ void
This method returns an undefined value.
Establish the server's protocol era (MCP 2026-07-28 basic/transports/stdio "Backward Compatibility"): probe with server/discover unless configured legacy-only, and fall back to the initialize handshake when the probe shows a legacy server.
557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 557 def negotiate_protocol return perform_initialize if @protocol_mode == :legacy return if probe_modern_server perform_initialize rescue StandardError # Nothing was negotiated. The probe only PROPOSES its version, and a # failure that never reached one of the classified outcomes — the # host's request_meta provider raising, say, so no probe was even # sent — would otherwise leave that proposal behind as a settled # modern era. The next attempt would then take a legacy server's # startup request for prohibited modern traffic and drop it, and a # server waiting for that response answers nothing: the recovery # deadlocks until it times out. @protocol_version = nil settle_era_probe raise end |
#next_id ⇒ Integer
Generate a new unique request ID and mark it as awaiting a response. Registering the id before the request is sent lets the reader thread distinguish expected responses from late/unsolicited ones.
817 818 819 820 821 822 823 824 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 817 def next_id @mutex.synchronize do id = @next_id @next_id += 1 @awaiting[id] = true id end end |
#open_subscription(subscription) ⇒ void
This method returns an undefined value.
Send the subscriptions/listen request for a subscription (a fresh id each time it is opened or re-opened).
77 78 79 80 81 82 83 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 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 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 77 def open_subscription(subscription) # The pipe this attempt writes to, taken before anything is recorded # about it and written to whatever happens next. # # Reading the transport's *current* stdin at the write instead let a # listen that was still pending when the process exited be written to # the process that replaced it: the teardown had already forgotten that # id (nothing written to a dead process is outstanding, and none of its # ids may be cancelled on its successor), so the replacement was # serving a second stream this client could no longer name — and the # restart's own listen was the only one `close` cancelled. Pinning the # pipe makes the bookkeeping follow the process actually written to: a # write that lands late goes to the pipe it was recorded against, and # once the teardown has closed that pipe it fails into the error paths # below instead. # Read with the generation it belongs to, so the two agree: the pipe # is what this attempt writes to, and the generation is what a # teardown compares its claim against when it decides whose # subscriptions to park. stdin, generation = @transport_lock.synchronize { [@stdin, @transport_generation] } # Whether the subscription has taken this attempt's id yet. A failure # before that — the request could not be built — is nobody's to have # superseded, and was filed as exactly that while the id it never # took was compared with the one it had (see {#fail_open_attempt}). taken = false id = next_id # No caller waits on this id: the response, if any, is the server's # graceful closure and is routed to the subscription itself. @mutex.synchronize { @awaiting.delete(id) } request = build_jsonrpc_request('subscriptions/listen', { 'notifications' => subscription.requested }, id) # A {Subscription#close} racing with a re-open must not leave the # server holding a subscription this client can no longer cancel: # taking the id and registering it happen under the subscription's own # lock (so a close that wins stops the re-open outright), and a close # that cancelled this id while the request was still going out is # named again below, once the server has seen the listen. taken = true return unless subscription.with_open_id(id, generation) { register_subscription(subscription) } # Recorded before the write, and whatever the write does: from here on # the server may be serving this listen, and {#cancel_subscription} # has to be able to name it even after a later request has taken the # subscription's own id (see {Subscription#record_outstanding_listen}). # It is not cancellable until the write has finished, though — a close # that named it while the pipe still held nothing would put # `cancelled(n)` on the wire ahead of `listen(n)`. So this attempt # marks it written and then cancels it itself, if that close has # happened by then. # # Both of those happen however the write ends. A write that raised may # still have put the request on the pipe, and the id is recorded for # exactly that reason; whether that request is then cancelled is # decided with the failure's verdict (see {#fail_open_attempt}). The # pipe is named on both, so the cancellation goes to the process the # request went to and to no other. subscription.record_outstanding_listen(id, stdin) begin send_request(request, io: stdin) ensure subscription.mark_listen_written(id) end cancel_outstanding_listens(subscription, io: stdin) if subscription.closed_by_client? rescue StandardError => e fail_open_attempt(subscription, taken ? id : nil, e, io: stdin) end |
#perform_discover ⇒ Hash
Send server/discover and apply the DiscoverResult. Bounded by discover_timeout rather than the general read timeout so a silent legacy server delays the fallback only briefly.
670 671 672 673 674 675 676 677 678 679 680 681 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 670 def perform_discover req_id = next_id req = build_registered_request('server/discover', {}, req_id) send_request(req) begin res = wait_response(req_id, timeout: @discover_timeout) rescue MCPClient::Errors::RequestTimeoutError send_cancellation_notification(req_id) raise end interpret_discover_answer(res) end |
#perform_initialize ⇒ void
This method returns an undefined value.
Handshake: send initialize request and initialized notification
750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 750 def perform_initialize # Initialize request init_id = next_id init_req = build_registered_request('initialize', initialization_params, init_id) send_request(init_req) res = wait_response(init_id) begin result = process_jsonrpc_response(res) || {} rescue MCPClient::Errors::UnsupportedProtocolVersionError => e # A modern-only server SHOULD name the versions it supports when # rejecting initialize (basic/versioning). When one of them is # mutual the era is settled after all — the fallback ran only # because the probe was too slow — so go back to server/discover # instead of ending the session. A legacy-only configuration has # opted out of the modern era and gets the error. return fall_forward_to_modern(e) if fall_forward_to_modern?(e) raise MCPClient::Errors::ConnectionError, "Initialize failed: #{e.} (server supports: #{e.supported.join(', ')})" rescue MCPClient::Errors::ServerError => e raise MCPClient::Errors::ConnectionError, "Initialize failed: #{e.}" end # Store negotiated protocol version, server info and capabilities. # Disconnects if the server negotiated a version we cannot speak. @protocol_version = validate_protocol_version!(result) @server_info = result['serverInfo'] @capabilities = result['capabilities'] @instructions = result['instructions'] # Send initialized notification notif = build_jsonrpc_notification('notifications/initialized', {}) @stdin.puts(notif.to_json) end |
#probe_modern_server ⇒ Boolean
Send the server/discover probe with this client's preferred modern version. Three outcomes, per the stdio backward-compatibility rules: a DiscoverResult (modern: select a version from supportedVersions), a recognized modern error such as UnsupportedProtocolVersionError (modern: retry with an advertised version, never fall back), or any other error / a timeout (legacy: fall back to initialize). The fallback is deliberately not keyed to one error code — legacy servers answer pre-initialize requests with implementation-defined errors.
587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 587 def probe_modern_server # The probe DECLARES this version; it does not establish it. Until the # answer arrives the era stays unknown, so an incoming server request # is still handled — a legacy server MAY ping during initialization and # the receiver MUST respond promptly, and a server waiting for that # response answers nothing until it arrives. @protocol_version = MCPClient::LATEST_PROTOCOL_VERSION begin_era_probe modern_confirmed = false begin perform_discover rescue MCPClient::Errors::UnsupportedProtocolVersionError => e raise unless e.modern_protocol_error? # A well-formed rejection settles the era: whatever the retried # probe does next, this server is modern and never gets initialize. modern_confirmed = true settle_era_probe retry_discover_with_advertised_version(e) end true rescue MCPClient::Errors::ConnectionError # A DiscoverResult (or advertised list) with no mutual version: the # server is modern but incompatible. Nothing was negotiated. @protocol_version = nil raise rescue MCPClient::Errors::ServerError, MCPClient::Errors::TransportError => e # A recognized modern error (-32020/-32021, or -32022 with no usable # version) identifies a modern server: surface it, never fall back. # Anything else — including a 2xx-style result that is not a # DiscoverResult — is a legacy server, unless the era was already # settled by a well-formed rejection. if modern_confirmed || e.modern_protocol_error_for_probe? @protocol_version = nil raise MCPClient::Errors::ConnectionError, "Server is modern but incompatible: #{e.}" end legacy_after_probe(e) false ensure # However the probe ended, it is no longer proposing anything. settle_era_probe end |
#queue_subscriptions_of_ended_process(subscriptions) ⇒ void
This method returns an undefined value.
Hand the subscriptions of a process that is being torn down to the next one, forgetting the listen ids that process was holding: nothing written to it is outstanding any more, and none of those ids may be cancelled on the process that replaces it (see MCPClient::Subscription#record_outstanding_listen).
Both steps happen under the lock a hand-over takes them off the queue under, so a session that is already re-sending them cannot have the ids it has just written forgotten by this teardown: it cannot reach its own writes until this has finished.
437 438 439 440 441 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 437 def queue_subscriptions_of_ended_process(subscriptions) reconnecting_mutex.synchronize do enqueue_reconnecting_locked(subscriptions).each(&:discard_outstanding_listens) end end |
#reconnecting_mutex ⇒ Mutex
Returns guards the queue of subscriptions waiting for a process (created by MCPClient::ServerStdio#initialize, so no two threads ever race to make it).
472 473 474 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 472 def reconnecting_mutex @reconnecting_mutex ||= Mutex.new end |
#reconnecting_subscriptions ⇒ Array<MCPClient::Subscription>
Returns the subscriptions on the queue that a process could still be re-sent to.
465 466 467 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 465 def reconnecting_subscriptions reconnecting_mutex.synchronize { (@reconnecting_subscriptions || []).select(&:reconnectable?) } end |
#release_retired_transport ⇒ void
This method returns an undefined value.
Discard a transport whose subprocess exited after a successful handshake, so the next negotiation starts from a clean slate rather than on top of the dead process's handles.
351 352 353 354 355 356 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 351 def release_retired_transport @logger.info('The MCP server subprocess exited; restarting it for this request') @initialized = false @transport_retired = false release_transport end |
#release_transport ⇒ void
This method returns an undefined value.
Tear down a transport that is not going to be used again — a handshake that never completed, or a subprocess that exited under a completed one. Failures are swallowed: the transport being unusable is often the reason it is being released, and the original error is the one worth raising.
364 365 366 367 368 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 364 def release_transport cleanup rescue StandardError => e @logger.debug("Releasing the stdio transport did not complete cleanly: #{e.}") end |
#reopen_refusal(carrier) ⇒ StandardError?
Why this session must not be given the open subscriptions, if it must not (see #reopen_subscriptions).
375 376 377 378 379 380 381 382 383 384 385 386 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 375 def reopen_refusal(carrier) unless modern? return MCPClient::Errors::CapabilityError.new( 'the re-established server process negotiated ' \ "#{protocol_version || 'no version'}, which cannot carry a subscriptions/listen stream" ) end return nil unless crash_looping?(carrier) @logger.error('MCP server process exited again right after it was given its subscriptions; closing them') MCPClient::Errors::TransportError.new('MCP server process exited again right after a restart') end |
#reopen_subscriptions(session = @session) ⇒ void
This method returns an undefined value.
After the process was re-established, re-send subscriptions/listen for every subscription the host still holds open ("the server holds no subscription state across reconnections").
This is the one place a session is handed the subscriptions, so it is the one place that decides whether to hand them over at all:
- a re-established session that turns out to be legacy cannot carry
them. On
protocol: :autothe restarted process may negotiate an older revision than the one that died, and MCPClient::ServerStdio#cleanup has already moved the open subscriptions out of the registry — returning would leave them :reconnecting for ever with the host never told. - neither can a session that would only continue a crash loop: if the last process these subscriptions were given died less than SUBSCRIPTION_RESTART_MIN_INTERVAL after receiving them, handing them over again would spawn the same corpse for ever. Deciding here rather than at the restart is what makes the bound hold: a process is re-established by whichever thread gets there first — the reader's restart or a host request — and only the re-send is common to both.
Either way the subscriptions end with the error, so the host learns
from closed?/error rather than waiting on a stream that is not
coming back.
287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 287 def reopen_subscriptions(session = @session) pending = take_reconnecting_subscriptions # A session handed nothing asks nothing, and must not spend the record # either: the subscriptions are still open on the session this one # replaced (a nested restart re-establishes the process between a # hand-over and the next exit), and the process that carried them is # still what the next hand-over has to be judged against. return if pending.empty? # The record of the process that last carried subscriptions answers # exactly one question — whether handing them over again would only # respawn the same corpse — and this is the moment it is asked. Asking # spends it, whatever the answer: the loop it recorded is either # broken here (these subscriptions are closed and never handed on) or # replaced below by the record of the process that takes them. Left # standing it outlived the loop it described, and the next hand-over — # of a subscription opened directly on the replacement, which then ran # healthily for hours — was refused for a crash it had no part in. carrier = @subscription_carrier @subscription_carrier = nil refusal = reopen_refusal(carrier) return fail_subscriptions(pending, refusal) if refusal # Stamped on the process before the writes go out: one that exits # while they are still going to it survived receiving them by no time # at all, which is what the next hand-over needs to know. session&. @subscription_carrier = session pending.each do |subscription| open_subscription(subscription) # A re-sent listen is a new request the replacement has to # acknowledge, so it carries the deadline the first one did: a # process that takes it and then says nothing is otherwise bounded # by nothing at all on stdio. rearm_acknowledgment_deadline(subscription) rescue StandardError => e @logger.warn("Could not re-establish subscription: #{e.}") end end |
#restart_for_open_subscriptions ⇒ void
This method returns an undefined value.
Restart the process the reader just watched exit, for the subscriptions the host still holds.
A subscription is a standing request the host does not repeat: while it only waits for notifications there is no RPC for #ensure_initialized to re-establish the process on, so leaving the restart to "the next request" leaves every subscription :reconnecting for ever, with the host neither notified nor served. Restarting is also what MCP 2026-07-28 stdio "Unexpected Termination" asks of a client, and re-sending the subscriptions afterwards is what this transport already promises. With no subscription open there is nothing standing, and the process stays lazily re-established on the next request.
A server that keeps exiting must not be respawned in a loop; that
bound is enforced where the subscriptions are handed over rather than
here (see #reopen_subscriptions), because the process is
re-established by whichever thread gets there first — this restart or
a host request that raced it — and only the hand-over is common to
both. A restart that fails outright ends them here instead. Either way
the host learns from closed?/error rather than waiting on a stream
that is never coming back.
499 500 501 502 503 504 505 506 507 508 509 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 499 def restart_for_open_subscriptions pending = reconnecting_subscriptions return if pending.empty? @logger.info("Re-establishing the server process for #{pending.size} open subscription(s)") ensure_initialized hand_over_to_established_process rescue StandardError => e @logger.warn("Could not re-establish the server process: #{e.}") fail_reconnecting_subscriptions(e) end |
#retry_discover_with_advertised_version(error) ⇒ void
This method returns an undefined value.
After UnsupportedProtocolVersionError, pick a mutually supported version from the error's advertised list and re-issue the probe.
652 653 654 655 656 657 658 659 660 661 662 663 664 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 652 def retry_discover_with_advertised_version(error) version = select_protocol_version(error.supported) unless version raise MCPClient::Errors::ConnectionError, "Server rejected protocol version #{@protocol_version} and supports only " \ "#{error.supported.join(', ')}, none of which this client speaks " \ "(modern versions supported: #{MCPClient::MODERN_PROTOCOL_VERSIONS.join(', ')})" end @logger.info("Server does not support #{@protocol_version}; retrying server/discover with #{version}") @protocol_version = version perform_discover end |
#rpc_notify(method, params = {}) ⇒ void
This method returns an undefined value.
Send a JSON-RPC notification (no response expected)
1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 1116 def rpc_notify(method, params = {}) ensure_initialized if suppressed_modern_notification?(method) @logger.debug("Not sending #{method}: removed in MCP #{protocol_version}") return end notif = build_jsonrpc_notification(method, params) begin @stdin.puts(notif.to_json) rescue StandardError => e # The same failure a request write reports, reported the same way: # a notification writes on its own, outside the request path's # check-and-write, and a dead pipe there reached the host as a raw # IOError. raise MCPClient::Errors::TransportError, "Failed to send JSONRPC notification: #{e.}" end end |
#rpc_request(method, params = {}, timeout: nil) ⇒ Object
Generic JSON-RPC request: send method with params and wait for result
974 975 976 977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 974 def rpc_request(method, params = {}, timeout: nil) freshly_probed = !@initialized || transport_retired? ensure_initialized if method == 'ping' && modern? # `ping` was removed in MCP 2026-07-28; the mandatory server/discover # request is the modern heartbeat. The probe that just established # the connection IS such a round trip, so answer from it rather than # paying for a second one. return @last_discover_result if freshly_probed && @last_discover_result method = 'server/discover' end # The multi round-trip resolver sits outside the per-attempt # recovery, so a retry that carries inputResponses/requestState keeps # them through transport retries, version renegotiation and the like. result = resolve_input_round_trips(method, params, timeout) do |attempt_params| # Every round is its own request, and the subprocess may have exited # between two of them — the host handler that gathers the input # waits for a person. MCP 2026-07-28 basic/transports/stdio # ("Unexpected Termination"): the client SHOULD restart a server # that terminated unexpectedly. The continuation is then re-issued # against the fresh process with the answers and the state it was # gathered for, rather than written to a dead pipe. ensure_initialized with_retry(method) do sent_version = nil begin send_request_and_wait(method, attempt_params, timeout) { |version| sent_version = version } rescue MCPClient::Errors::UnsupportedProtocolVersionError => e # MCP 2026-07-28 basic/versioning: "The client SHOULD select a # mutually supported version from the supported list and retry # the request". The server rejected the request before # processing it, so re-sending cannot duplicate a side effect. # Compared against the version THIS request declared, read back # from the request itself: a concurrent request may have moved # the transport on while this one was being built. version = select_protocol_version(e.supported) raise unless modern? && version && version != sent_version @logger.info("Server does not support protocol version #{sent_version}; " \ "retrying #{method} with #{version}") @protocol_version = version send_request_and_wait(method, attempt_params, timeout) end end 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: the transport may be the reason the request timed out in the first place.
1104 1105 1106 1107 1108 1109 1110 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 1104 def send_cancellation_notification(request_id) notif = build_jsonrpc_notification('notifications/cancelled', { 'requestId' => request_id, 'reason' => 'Request timed out' }) @stdin.puts(notif.to_json) rescue StandardError => e @logger.debug("Failed to send cancellation notification: #{e.}") end |
#send_if_current(req, generation) ⇒ Hash?
Write a registered request, unless the transport it was registered on has been replaced since. Replacement is judged by the transport generation, which every teardown and every spawn bumps: a restart in between has dropped the id from @awaiting, and the process the request was built for is gone.
852 853 854 855 856 857 858 859 860 861 862 863 864 865 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 852 def send_if_current(req, generation) return req unless send_request(req, generation) == :replaced @logger.debug("The transport was replaced before #{req['method']} was sent; re-issuing it") # Nothing was written, so nobody will ever wait on this id: drop it # from the teardown's record as well as from the awaiting table. The # waiter is what normally consumes that record, and an id no request # carries has no waiter — it would accumulate for the process's life. @mutex.synchronize do @awaiting.delete(req['id']) dropped_requests.delete(req['id']) end nil end |
#send_on_current_transport(method, params) {|req| ... } ⇒ Array(Integer, Hash)
Register, build and write a request on the transport that is current when it is written. Between registering the id and writing, the host's request_meta provider runs, and in that window the subprocess may exit and another thread restart it: the restart drops every outstanding id, so a request written afterwards would go out unregistered and its answer be discarded as unsolicited. Nothing has been sent when that is detected, so the request is simply rebuilt — after waiting for the restart to complete — and sent registered.
1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 1071 def send_on_current_transport(method, params) loop do # A subprocess that exited under the handshake is restarted here # (MCP 2026-07-28 stdio "Unexpected Termination": the client # SHOULD restart it) rather than written to. ensure_initialized if transport_retired? generation = @transport_generation req_id = next_id req = build_registered_request(method, params, req_id) yield req if block_given? return [req_id, req] if send_if_current(req, generation) ensure_initialized end end |
#send_request(req, generation = nil, io: @stdin) ⇒ Symbol
Send a JSON-RPC request. With a generation, the request is written only if that is still the transport's generation, and the check and the write are one step under the transport lock: a restart replaces the handles and bumps the generation under the same lock, so a request judged current cannot be written to the replacement process — where it would be executed unregistered and its answer discarded.
880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 880 def send_request(req, generation = nil, io: @stdin) @logger.debug("Sending JSONRPC request: #{(req)}") # A request pinned to a session that has since ended is not written at # all: its payload names something else in the replacement session. # The pin is read outside the transport lock (the two locks never # nest); a restart completing between this check and the write moves # the transport generation, which the locked check below catches. @mutex.synchronize { check_session_pin! } @transport_lock.synchronize do # A replacement whose negotiation has not completed is not current # either, whatever its generation says: an ordinary request written # to it would reach the process before its handshake. return :replaced if generation && (generation != @transport_generation || @negotiating) raise IOError, 'the server process is gone' unless io # The write goes to the pipe this request was recorded against, # never to whichever pipe the transport holds by now. io.puts(req.to_json) end :sent rescue MCPClient::Errors::SessionChangedError, MCPClient::Errors::TaskReplacedError # A refusal, not a failure: the pin (or the caller's own pre-write # guard, see {MCPClient::SessionPin#guarded_writes}) turned the write # down. Nothing was written, nothing will answer this id, and the # refusal keeps its type — a definite "this request is not to be # sent" must not reach the caller as an ambiguous transport failure. @mutex.synchronize { @awaiting.delete(req['id']) } if req.is_a?(Hash) && req['id'] raise rescue StandardError => e # A request that failed to send will never receive a response, so drop # its awaiting marker; otherwise a broken transport (e.g. the server # exited) would leak an entry per retry/attempt into @awaiting. @mutex.synchronize { @awaiting.delete(req['id']) } if req.is_a?(Hash) && req['id'] raise MCPClient::Errors::TransportError, "Failed to send JSONRPC request: #{e.}" end |
#send_request_and_wait(method, params, timeout) {|version| ... } ⇒ Object
One request/response exchange with its own JSON-RPC id.
1039 1040 1041 1042 1043 1044 1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 1039 def send_request_and_wait(method, params, timeout) # As late as a request pinned to a session can be held back: every # reconnect on the way here (ensure_initialized, a retry after the # child exited) has happened by now. check_session_pin! req_id, = send_on_current_transport(method, params) do |built| clear_response_received_at if respond_to?(:clear_response_received_at, true) yield declared_protocol_version(built) if block_given? end begin res = wait_response(req_id, timeout: timeout) rescue MCPClient::Errors::RequestTimeoutError # MCP lifecycle: on timeout the sender SHOULD issue a cancellation # notification for the abandoned request and stop waiting. send_cancellation_notification(req_id) if cancellable_request?(method, params) raise end process_jsonrpc_response(res, method: method) end |
#subscription_failure(error) ⇒ MCPClient::Errors::MCPError
Returns the error a subscription ends with.
211 212 213 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 211 def subscription_failure(error) error.is_a?(MCPClient::Errors::MCPError) ? error : MCPClient::Errors::TransportError.new(error.) end |
#take_reconnecting_subscriptions ⇒ Array<MCPClient::Subscription>
Take the subscriptions waiting for a process, leaving the queue empty.
455 456 457 458 459 460 461 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 455 def take_reconnecting_subscriptions reconnecting_mutex.synchronize do pending = (@reconnecting_subscriptions || []).select(&:reconnectable?) @reconnecting_subscriptions = [] pending end end |
#transport_retired? ⇒ Boolean
Returns whether the subprocess behind the handshake exited.
343 344 345 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 343 def transport_retired? @transport_retired end |
#wait_response(id, timeout: nil) ⇒ Hash
Wait for a response with the given request ID
920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938 939 940 941 942 943 944 945 946 947 948 949 950 951 952 953 954 955 |
# File 'lib/mcp_client/server_stdio/json_rpc_transport.rb', line 920 def wait_response(id, timeout: nil) deadline = Time.now + (timeout || @read_timeout) @mutex.synchronize do until @pending.key?(id) # The subprocess exited: no answer is coming, however long the # timeout. (An answer that arrived before it exited is above.) # A restart another caller completed meanwhile has cleared the # retirement again, but it recorded this request as dropped — # the durable sign that the transport it went out on is gone. break if @transport_retired || dropped_requests.include?(id) remaining = deadline - Time.now break if remaining <= 0 @cond.wait(@mutex, remaining) end # Remove the response and the awaiting marker on both success and # timeout so neither @pending nor @awaiting accumulates entries. msg = @pending.delete(id) transport_gone = @transport_retired || !dropped_requests.delete?(id).nil? arrival = (@response_arrivals ||= {}).delete(id) @awaiting.delete(id) if msg # The response's receipt time is the reader's, not this wake-up. note_response_received_at(arrival || monotonic_now) if respond_to?(:note_response_received_at, true) return msg end if transport_gone raise MCPClient::Errors::TransportError, "The MCP server subprocess exited before answering JSONRPC request id=#{id}" end raise MCPClient::Errors::RequestTimeoutError, "Timeout waiting for JSONRPC response id=#{id}" end end |