Class: MCPClient::ServerStdio
- Inherits:
-
ServerBase
- Object
- ServerBase
- MCPClient::ServerStdio
- Includes:
- RequestMetaScope, JsonRpcTransport
- Defined in:
- lib/mcp_client/server_stdio.rb,
lib/mcp_client/server_stdio/child_session.rb,
lib/mcp_client/server_stdio/json_rpc_transport.rb
Overview
JSON-RPC implementation of MCP server over stdio.
Defined Under Namespace
Modules: JsonRpcTransport Classes: ChildSession, TornDownTransport
Constant Summary collapse
- READ_TIMEOUT =
Timeout in seconds for responses
15- SHUTDOWN_GRACE_PERIOD =
Grace period in seconds allowed at each stage of the shutdown sequence (after closing stdin, then after SIGTERM) before escalating further, per MCP 2025-11-25 basic/lifecycle.mdx (Shutdown / stdio): close stdin, wait for the server to exit, send SIGTERM, then SIGKILL if it still runs.
2- SUBSCRIPTION_RESTART_MIN_INTERVAL =
Seconds a process must last after the open subscriptions were re-sent to it — counted from the moment it received them, not from the moment it was spawned — before another unexpected exit counts as a crash rather than a crash loop (see JsonRpcTransport#reopen_subscriptions).
5- STDERR_READ_CHUNK_SIZE =
Chunk size (bytes) used when draining the subprocess stderr pipe
8192- STDERR_MAX_LINE_SIZE =
Maximum bytes buffered for a single unterminated stderr line before it is flushed. Bounds memory when a server writes to stderr without newlines (e.g. progress output using carriage returns).
64 * 1024
- PROTOCOL_MODES =
How the server's protocol era is established (MCP 2026-07-28 basic/transports/stdio "Backward Compatibility"):
- :auto probe with server/discover, fall back to initialize on any non-modern error or timeout (dual-era client, the default)
- :modern probe with server/discover and fail if the server is legacy
- :legacy skip the probe and run the initialize handshake
%i[auto modern legacy].freeze
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
-
#capabilities ⇒ Hash?
readonly
Server capabilities from the initialize response.
-
#command ⇒ String, Array
readonly
The command used to launch the server.
-
#discover_timeout ⇒ Numeric
readonly
Seconds allowed for the server/discover probe.
-
#env ⇒ Object
readonly
Returns the value of attribute env.
-
#protocol_mode ⇒ Symbol
readonly
The configured protocol mode (:auto, :modern or :legacy).
-
#server_info ⇒ Hash?
readonly
Server info from the initialize response.
Attributes included from JsonRpcCommon
#request_meta, #send_client_info
Attributes inherited from ServerBase
#instructions, #logger, #name, #read_timeout
Instance Method Summary collapse
-
#call_tool(tool_name, parameters) ⇒ Object
Call a tool with the given parameters.
-
#cancel_outstanding_listens(subscription, io: @stdin) ⇒ void
Tell the server the client has stopped reading every listen request it wrote for this subscription on the process this pipe belongs to.
-
#cancel_subscription(subscription) ⇒ void
Cancel a subscription: on stdio there is no per-request stream to close, so the client sends notifications/cancelled referencing the subscriptions/listen request id.
-
#claim_transport(generation) ⇒ TornDownTransport?
Take the process of a transport generation off the transport, for one teardown to dismantle.
-
#claimed_subscription?(subscription, claimed) ⇒ Boolean
Whether a registered subscription belongs to the process a teardown claimed.
-
#cleanup ⇒ void
Clean up the server connection Closes all stdio handles and terminates any running processes and threads following the MCP 2025-11-25 stdio shutdown sequence (basic/lifecycle.mdx): close stdin, wait for the server to exit, send SIGTERM if it does not exit within a reasonable time, then SIGKILL if it still does not exit.
-
#collect_prompt_pages(page) ⇒ Array<MCPClient::Prompt>
Collect every page of prompts/list, recording what each page was answered under so the cache can bind the combined list to it.
-
#collect_tool_pages(page) ⇒ Array<MCPClient::Tool>
Collect every page of tools/list, recording what each page was answered under so the cache can bind the combined list to it.
-
#complete(ref:, argument:, context: nil) ⇒ Hash
Request completion suggestions from the server (MCP 2025-06-18).
-
#connect ⇒ Boolean
Connect to the MCP server by launching the command process via stdin/stdout.
-
#current_session?(session) ⇒ Boolean
Whether the process a reader was started for is still the live one.
-
#ending_session? ⇒ Boolean
Whether #cleanup ends a session: only a 2025-11-25 handshake opens one, and only once it completed.
-
#flush_stderr_lines(buffer) ⇒ void
Emit and remove all newline-terminated lines from the stderr buffer.
-
#flush_stderr_overflow(buffer) ⇒ void
Flush an unterminated stderr fragment that has grown past the size cap, so a newline-less stderr stream cannot buffer without bound.
-
#forget_torn_down_transport(claimed) ⇒ void
The last word of a teardown, whatever else went wrong: the record of the process it claimed is ended, and — only if no process has been established since the claim — the transport forgets its session and its handshake, and wakes the requests that will never be answered.
-
#get_prompt(prompt_name, parameters) ⇒ Object
Get a prompt with the given parameters.
-
#handle_elicitation_create(request_id, params) ⇒ void
Handle elicitation/create request from server (MCP 2025-06-18).
-
#handle_line(line) ⇒ void
Handle a line of output from the stdio server Parses JSON-RPC messages and adds them to pending responses.
-
#handle_ping(request_id) ⇒ void
Handle a server-initiated ping request (MCP ping utility) The receiver MUST respond promptly with an empty result.
-
#handle_reader_eof(session, generation = @transport_generation) ⇒ void
EOF on the process's stdout: unless the client closed its stdin, the server exited on its own.
-
#handle_roots_list(request_id, params) ⇒ void
Handle roots/list request from server (MCP 2025-06-18).
-
#handle_sampling_create_message(request_id, params) ⇒ void
Handle sampling/createMessage request from server (MCP 2025-11-25).
-
#handle_server_exit(session = @session, generation = @transport_generation) ⇒ void
The server process ended on its own (its stdout reached EOF).
-
#handle_server_request(msg) ⇒ void
Handle incoming JSON-RPC request from server (MCP 2025-06-18).
-
#initialize(command:, retries: 0, retry_backoff: 1, read_timeout: READ_TIMEOUT, name: nil, logger: nil, env: {}, protocol: :auto, discover_timeout: nil) ⇒ ServerStdio
constructor
Initialize a new ServerStdio instance.
-
#list_prompts ⇒ Array<MCPClient::Prompt>
List all prompts available from the MCP server.
-
#list_resource_templates(cursor: nil) ⇒ Hash
List all resource templates available from the MCP server.
-
#list_resources(cursor: nil) ⇒ Hash
List all resources available from the MCP server.
-
#list_tools ⇒ Array<MCPClient::Tool>
List all tools available from the MCP server.
-
#live_transport?(generation, session) ⇒ Boolean
Whether the process a reader was started for is still the live one: its transport generation has not been claimed for teardown, and its record is the current session.
-
#log_level=(level) ⇒ Hash
deprecated
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.
-
#modern_peer? ⇒ Boolean
Whether a server-initiated request is prohibited traffic.
-
#on_elicitation_request(&block) ⇒ void
Register a callback for elicitation requests (MCP 2025-06-18).
-
#on_roots_list_request(&block) ⇒ void
deprecated
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.
-
#on_sampling_request(&block) ⇒ void
deprecated
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.
-
#park_open_subscriptions(claimed) ⇒ void
Subscriptions do not survive the process: keep the ones the host still wants so they are re-sent once the process is re-established (basic/patterns/subscriptions "Graceful Closure").
-
#pin_pipe_encodings ⇒ void
Pin the subprocess pipe encodings to UTF-8 instead of inheriting the process locale (Encoding.default_external).
-
#read_resource(uri) ⇒ Array<MCPClient::ResourceContent>
Read a resource by its URI.
-
#record_response(id, msg, arrived) ⇒ void
Queue a response for the caller waiting on it.
-
#retire_transport(generation) ⇒ void
The subprocess closed its stdout: it has exited (or its pipes were dropped), so no response will ever arrive on this transport again.
-
#send_elicitation_response(request_id, result) ⇒ void
Send elicitation response back to server (MCP 2025-06-18).
-
#send_error_response(request_id, code, message) ⇒ void
Send error response back to server (MCP 2025-06-18).
-
#send_message(message) ⇒ void
Send a JSON-RPC message to the server.
-
#send_roots_list_response(request_id, result) ⇒ void
Send roots/list response back to server (MCP 2025-06-18).
-
#send_sampling_response(request_id, result) ⇒ void
Send sampling response back to server (MCP 2025-11-25).
-
#send_subscription_cancellation(id, io: @stdin) ⇒ void
Tell the server the client closed a subscriptions/listen request.
-
#signal_server_process(signal, wait_thread = @wait_thread) ⇒ void
Send a signal to the server process, tolerating a process that has already exited or cannot be signalled.
-
#spawn_server_process ⇒ Array
The stdin, stdout, stderr and wait thread of the spawned process.
-
#start_reader ⇒ Thread
Spawn a reader thread to collect JSON-RPC responses.
-
#start_stderr_reader ⇒ Thread
Spawn a thread to continuously drain the subprocess stderr.
-
#subscribe_resource(uri) ⇒ Boolean
Subscribe to resource updates.
-
#teardown_transport(generation) ⇒ void
Tear down the process of one transport generation.
-
#terminate_server_process(wait_thread = @wait_thread) ⇒ void
Terminate the spawned server process per the MCP 2025-11-25 stdio shutdown sequence (basic/lifecycle.mdx): stdin has already been closed, so wait for the process to exit on its own; if it does not exit within the grace period send SIGTERM, wait again, and finally send SIGKILL.
-
#unsubscribe_resource(uri) ⇒ Boolean
Unsubscribe from resource updates.
Methods included from JsonRpcTransport
#build_registered_request, #call_tool_streaming, #crash_looping?, #declared_protocol_version, #defer_reestablished_attempt, #discover_result?, #dropped_requests, #enqueue_reconnecting_locked, #enqueue_reconnecting_subscriptions, #ensure_initialized, #ensure_session_ready, #fail_open_attempt, #fail_reconnecting_subscriptions, #fail_subscriptions, #fail_superseded_attempt, #fall_forward_to_modern, #fall_forward_to_modern?, #hand_over_to_established_process, #identifies_modern_server?, #interpret_discover_answer, #invalid_discover_answer, #legacy_after_probe, #live_process?, #modern_discover_answer?, #negotiate_protocol, #next_id, #open_subscription, #perform_discover, #perform_initialize, #probe_modern_server, #queue_subscriptions_of_ended_process, #reconnecting_mutex, #reconnecting_subscriptions, #release_retired_transport, #release_transport, #reopen_refusal, #reopen_subscriptions, #restart_for_open_subscriptions, #retry_discover_with_advertised_version, #rpc_notify, #rpc_request, #send_cancellation_notification, #send_if_current, #send_on_current_transport, #send_request, #send_request_and_wait, #subscription_failure, #take_reconnecting_subscriptions, #transport_retired?, #wait_response
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
#call_tool_streaming, #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(command:, retries: 0, retry_backoff: 1, read_timeout: READ_TIMEOUT, name: nil, logger: nil, env: {}, protocol: :auto, discover_timeout: nil) ⇒ ServerStdio
Initialize a new ServerStdio instance
73 74 75 76 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 |
# File 'lib/mcp_client/server_stdio.rb', line 73 def initialize(command:, retries: 0, retry_backoff: 1, read_timeout: READ_TIMEOUT, name: nil, logger: nil, env: {}, protocol: :auto, discover_timeout: nil) super(name: name) unless PROTOCOL_MODES.include?(protocol) raise ArgumentError, "protocol must be one of #{PROTOCOL_MODES.inspect}, got #{protocol.inspect}" end @protocol_mode = protocol @discover_timeout = discover_timeout || read_timeout @command_array = command.is_a?(Array) ? command : nil @command = command.is_a?(Array) ? command.join(' ') : command @mutex = Mutex.new @cond = ConditionVariable.new # Serializes process start and protocol negotiation: two threads must # not each spawn the server or run the probe and the fallback on the # same pipes. Reentrant because negotiation waits on @mutex inside. @init_lock = Monitor.new @next_id = 1 @pending = {} # Ids of requests awaiting a response; used to drop late/unsolicited # responses so @pending cannot grow without bound on a long-lived session @awaiting = {} @initialized = false @server_info = nil @capabilities = nil initialize_logger(logger) @max_retries = retries @retry_backoff = retry_backoff @read_timeout = read_timeout @env = env || {} @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 @reader_thread = nil @stderr_thread = nil # Bumped whenever a subprocess is spawned or torn down, so a reader # thread only ever speaks for the transport it was started for. @transport_generation = 0 # Guards the pair (subprocess handles, generation) so that judging # whether a request's transport is still current and writing it are # one step: a restart cannot slip in between (see # JsonRpcTransport#send_request). @transport_lock = Mutex.new @negotiating = false @transport_retired = false @modern_answer_received = false # The record of the live child process, and of the one the open # subscriptions were last re-sent to — the crash-loop bound is read from # the latter (see ChildSession and JsonRpcTransport#reopen_subscriptions). @session = nil @subscription_carrier = nil # The subscriptions waiting for a process, and the lock the two paths # that write them share: a `cleanup` moving the open subscriptions onto # the queue overlaps a hand-over whose listen write failed putting one # back (JsonRpcTransport#enqueue_reconnecting_subscriptions). Made here # so no two threads ever race to make it. @reconnecting_mutex = Mutex.new @reconnecting_subscriptions = [] end |
Instance Attribute Details
#capabilities ⇒ Hash? (readonly)
Server capabilities from the initialize response
139 140 141 |
# File 'lib/mcp_client/server_stdio.rb', line 139 def capabilities @capabilities end |
#command ⇒ String, Array (readonly)
Returns the command used to launch the server.
26 27 28 |
# File 'lib/mcp_client/server_stdio.rb', line 26 def command @command end |
#discover_timeout ⇒ Numeric (readonly)
Returns seconds allowed for the server/discover probe.
145 146 147 |
# File 'lib/mcp_client/server_stdio.rb', line 145 def discover_timeout @discover_timeout end |
#env ⇒ Object (readonly)
Returns the value of attribute env.
26 |
# File 'lib/mcp_client/server_stdio.rb', line 26 attr_reader :command, :env |
#protocol_mode ⇒ Symbol (readonly)
Returns the configured protocol mode (:auto, :modern or :legacy).
142 143 144 |
# File 'lib/mcp_client/server_stdio.rb', line 142 def protocol_mode @protocol_mode end |
#server_info ⇒ Hash? (readonly)
Server info from the initialize response
135 136 137 |
# File 'lib/mcp_client/server_stdio.rb', line 135 def server_info @server_info end |
Instance Method Details
#call_tool(tool_name, parameters) ⇒ Object
Call a tool with the given parameters
710 711 712 713 714 715 716 717 718 719 720 721 |
# File 'lib/mcp_client/server_stdio.rb', line 710 def call_tool(tool_name, parameters) ensure_initialized rpc_request('tools/call', build_named_request_params(tool_name, parameters)) rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ToolCallError, "Error calling tool '#{tool_name}': #{e.}" rescue StandardError => e raise MCPClient::Errors::ToolCallError, "Error calling tool '#{tool_name}': #{e.}" end |
#cancel_outstanding_listens(subscription, io: @stdin) ⇒ void
This method returns an undefined value.
Tell the server the client has stopped reading every listen request it wrote for this subscription on the process this pipe belongs to.
Not just the one the subscription is on: a second listen written for it on one process — a hand-over the queue duplicated, say — leaves the server serving the first stream, and naming only the newest id left that one open with the client no longer able to refer to it.
The subscription's current id is not named on top of those. It used to be, so that a handle closed between taking an id and writing its request was cancelled anyway — but that cancellation named a request the server had not been sent, and reached the pipe ahead of it. Every id this client wrote is recorded (MCPClient::Subscription#record_outstanding_listen), so a written id is here already; an id that is only assigned is the opener's to cancel, once it has actually written it.
Ids written to a pipe other than this one are left where they are: the process reading this one was never sent those requests, and naming them on it would cancel requests it has never seen (MCPClient::Subscription#take_outstanding_listens).
1339 1340 1341 |
# File 'lib/mcp_client/server_stdio.rb', line 1339 def cancel_outstanding_listens(subscription, io: @stdin) subscription.take_outstanding_listens(io).each { |id| send_subscription_cancellation(id, io: io) } end |
#cancel_subscription(subscription) ⇒ void
This method returns an undefined value.
Cancel a subscription: on stdio there is no per-request stream to close, so the client sends notifications/cancelled referencing the subscriptions/listen request id.
1304 1305 1306 1307 1308 1309 1310 1311 1312 |
# File 'lib/mcp_client/server_stdio.rb', line 1304 def cancel_subscription(subscription) # Closed first: a re-open in flight holds the subscription's lock until # it has taken its new id, so the id unregistered and cancelled below is # the one the server was actually sent. subscription.finish(by_client: true) unregister_subscription(subscription) subscriptions_mutex.synchronize { resource_subscriptions.delete_if { |_uri, sub| sub.equal?(subscription) } } cancel_outstanding_listens(subscription) end |
#claim_transport(generation) ⇒ TornDownTransport?
Take the process of a transport generation off the transport, for one teardown to dismantle. Under the transport lock: the handles and the generation change together, and closing stdin here is what makes a request judged current either already written or never written.
1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 |
# File 'lib/mcp_client/server_stdio.rb', line 1116 def claim_transport(generation) @transport_lock.synchronize do return nil unless @stdin && generation == @transport_generation # Past this point the reader threads speak for a transport that is # being dismantled on purpose: their EOF must not retire whatever # replaces it. A 2025-11-25 handshake opened a session that ends # with the process. A stateless 2026-07-28 peer holds none: a task it # created outlives the connection exactly as it does over sessionless # HTTP, so the process that replaces this one is asked about it # rather than the task being written off as gone with a session that # never existed. bump_session_epoch if ending_session? @transport_generation += 1 claimed = TornDownTransport.new(generation: @transport_generation, stdin: @stdin, stdout: @stdout, stderr: @stderr, wait_thread: @wait_thread, reader_thread: @reader_thread, stderr_thread: @stderr_thread, session: @session) @stdin.close unless @stdin.closed? @stdin = @stdout = @stderr = @wait_thread = @reader_thread = @stderr_thread = nil # The ids still outstanding went out on the process just claimed and # will never be answered. They are recorded as dropped here, at the # claim, so their waiters fail on that record — promptly, woken here — # even when a replacement is established before this teardown # finishes and {#forget_torn_down_transport} therefore leaves the # replacement's bookkeeping alone. @mutex.synchronize do dropped_requests.merge(@awaiting.keys) @cond.broadcast end claimed end end |
#claimed_subscription?(subscription, claimed) ⇒ Boolean
Whether a registered subscription belongs to the process a teardown claimed. The claim's generation is the one the transport moved to, so the process it took is everything below it; a subscription opened on the replacement carries a higher one. One that never recorded a generation is treated as the claim's, which is what a registry entry was before the stamp existed.
1191 1192 1193 1194 |
# File 'lib/mcp_client/server_stdio.rb', line 1191 def claimed_subscription?(subscription, claimed) generation = subscription.open_generation generation.nil? || generation < claimed.generation end |
#cleanup ⇒ void
This method returns an undefined value.
Clean up the server connection Closes all stdio handles and terminates any running processes and threads following the MCP 2025-11-25 stdio shutdown sequence (basic/lifecycle.mdx): close stdin, wait for the server to exit, send SIGTERM if it does not exit within a reasonable time, then SIGKILL if it still does not exit.
Tears down the process the transport holds now: a cleanup that
finds the process already claimed by another teardown — the reader of a
process that exited, dismantling it while the host got here — has
nothing to do, and never touches a replacement the host established in
the meantime (see #teardown_transport).
1054 1055 1056 1057 1058 1059 1060 1061 |
# File 'lib/mcp_client/server_stdio.rb', line 1054 def cleanup # 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. forget_transport_thread_state teardown_transport(@transport_lock.synchronize { @transport_generation }) end |
#collect_prompt_pages(page) ⇒ Array<MCPClient::Prompt>
Collect every page of prompts/list, recording what each page was answered under so the cache can bind the combined list to it.
468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 |
# File 'lib/mcp_client/server_stdio.rb', line 468 def collect_prompt_pages(page) pages = [] received_ats = [] fingerprints = [] epoch = cache_epoch(:prompts) prompts = collect_paginated('prompts') do |cursor| params = {} params['cursor'] = cursor if cursor started = monotonic_now page[:cursor] = cursor page_result = fetch_list_page(:prompts, cursor) { rpc_request('prompts/list', params) } || {} result = require_complete_result!(page_result, 'prompts/list') pages << result received_ats << response_received_at(since: started) fingerprints << request_params_fingerprint prompts = (result['prompts'] || []).map { |td| MCPClient::Prompt.from_json(td, server: self) } [prompts, result['nextCursor']] end record_list_cache_hint('prompts/list', pages, received_ats, params: fingerprints, epoch: epoch) prompts end |
#collect_tool_pages(page) ⇒ Array<MCPClient::Tool>
Collect every page of tools/list, recording what each page was answered under so the cache can bind the combined list to it.
682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 |
# File 'lib/mcp_client/server_stdio.rb', line 682 def collect_tool_pages(page) pages = [] received_ats = [] fingerprints = [] epoch = cache_epoch(:tools) tools = collect_paginated('tools') do |cursor| params = {} params['cursor'] = cursor if cursor started = monotonic_now page[:cursor] = cursor page_result = fetch_list_page(:tools, cursor) { rpc_request('tools/list', params) } || {} result = require_complete_result!(page_result, 'tools/list') pages << result received_ats << response_received_at(since: started) fingerprints << request_params_fingerprint tools = (result['tools'] || []).map { |td| MCPClient::Tool.from_json(td, server: self) } [tools, result['nextCursor']] end record_list_cache_hint('tools/list', pages, received_ats, params: fingerprints, epoch: epoch) tools end |
#complete(ref:, argument:, context: nil) ⇒ Hash
Request completion suggestions from the server (MCP 2025-06-18)
729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 |
# File 'lib/mcp_client/server_stdio.rb', line 729 def complete(ref:, argument:, context: nil) ensure_initialized 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::CapabilityError raise rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ServerError, "Error requesting completion: #{e.}" rescue StandardError => e raise MCPClient::Errors::ServerError, "Error requesting completion: #{e.}" end |
#connect ⇒ Boolean
Connect to the MCP server by launching the command process via stdin/stdout
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 |
# File 'lib/mcp_client/server_stdio.rb', line 150 def connect handles = spawn_server_process # The handles and the generation change together: a request judged # current against the old generation must not find the new stdin. @transport_lock.synchronize do @stdin, @stdout, @stderr, @wait_thread = handles @transport_generation += 1 # A fresh process is not the one that exited: a restart made for the # open subscriptions must not be torn down again by the next request. @transport_retired = false # A fresh process has said nothing yet: what the previous one wrote # identifies nothing about this one, and what was negotiated WITH it # binds nothing here. The era is per process (stdio "Backward # Compatibility"), and the replacement's reader starts before the # probe proposes anything: a 2025-11-25 replacement that pings at # startup would otherwise be judged by the dead process's era, its # ping dropped, and both the probe and the handshake left waiting on # a server that answers nothing until its pong arrives. @modern_answer_received = false @protocol_version = nil end pin_pipe_encodings true rescue StandardError => e raise MCPClient::Errors::ConnectionError, "Failed to connect to MCP server: #{e.}" end |
#current_session?(session) ⇒ Boolean
Whether the process a reader was started for is still the live one. A reader whose process has already been torn down and replaced must not tear down its successor.
297 298 299 |
# File 'lib/mcp_client/server_stdio.rb', line 297 def current_session?(session) session.nil? || @session.equal?(session) end |
#ending_session? ⇒ Boolean
Whether #cleanup ends a session: only a 2025-11-25 handshake opens one, and only once it completed. A stateless 2026-07-28 peer, a process that never got through its handshake and a probe still in flight leave nothing session-scoped behind — task ids and their bookkeeping stay what they are for the process that comes next.
1038 1039 1040 |
# File 'lib/mcp_client/server_stdio.rb', line 1038 def ending_session? @initialized && !modern_peer? end |
#flush_stderr_lines(buffer) ⇒ void
This method returns an undefined value.
Emit and remove all newline-terminated lines from the stderr buffer.
1014 1015 1016 1017 1018 1019 |
# File 'lib/mcp_client/server_stdio.rb', line 1014 def flush_stderr_lines(buffer) while (newline_index = buffer.index("\n")) line = buffer.slice!(0, newline_index + 1) @logger.debug("[stderr] #{line.chomp}") end end |
#flush_stderr_overflow(buffer) ⇒ void
This method returns an undefined value.
Flush an unterminated stderr fragment that has grown past the size cap, so a newline-less stderr stream cannot buffer without bound.
1025 1026 1027 1028 1029 1030 |
# File 'lib/mcp_client/server_stdio.rb', line 1025 def flush_stderr_overflow(buffer) return if buffer.bytesize <= STDERR_MAX_LINE_SIZE @logger.debug("[stderr] #{buffer}") buffer.clear end |
#forget_torn_down_transport(claimed) ⇒ void
This method returns an undefined value.
The last word of a teardown, whatever else went wrong: the record of the process it claimed is ended, and — only if no process has been established since the claim — the transport forgets its session and its handshake, and wakes the requests that will never be answered.
No further response can arrive on a transport that is being dismantled, so nothing is outstanding any more. Responses that already arrived are kept: they are answers this client received and has not handed to their caller yet, and a restart happening in that window must not turn a completed request into a timeout. Each one belongs to a caller that is about to take it out of the map, so keeping them cannot accumulate. Waiters are woken so a request that will never be answered re-checks its deadline rather than blocking on a reader thread that has been killed. A replacement established meanwhile owns whatever is outstanding now, and its handshake and session are its own.
1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 |
# File 'lib/mcp_client/server_stdio.rb', line 1213 def forget_torn_down_transport(claimed) # The process this session was is gone, whatever else went wrong above. # Its record outlives it: it is what the next session's re-send of the # open subscriptions asks about (JsonRpcTransport#reopen_subscriptions). claimed.session&.ended # Asked and acted on in one step, under the lock a replacement is # established through. Asking first and writing afterwards let the # answer go stale in between: a host request that established the # replacement in that window did so *after* both checks passed, and # this teardown then marked its outstanding requests dropped and left # it uninitialized with no session — a live process the transport could # no longer name. @transport_lock.synchronize do # Either signal says a process has been established since the claim: # `connect` moves the generation on, and the handshake that follows # records a new session. next if @transport_generation != claimed.generation next unless @session.equal?(claimed.session) @mutex.synchronize do # The ids still outstanding are recorded as dropped: their waiters # fail on that record, whenever they next run, rather than wait out # their timeouts because the restart cleared the retirement first. dropped_requests.merge(@awaiting.keys) @response_arrivals&.clear @awaiting.clear @cond.broadcast end @session = nil # Cached results belong to the process that just ended. clear_result_cache # The next request re-establishes the process and, on a modern # server, re-sends the subscriptions the host still holds. @initialized = false end end |
#get_prompt(prompt_name, parameters) ⇒ Object
Get a prompt with the given parameters
496 497 498 499 500 501 502 503 504 505 506 507 |
# File 'lib/mcp_client/server_stdio.rb', line 496 def get_prompt(prompt_name, parameters) ensure_initialized rpc_request('prompts/get', build_named_request_params(prompt_name, parameters)) rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::PromptGetError, "Error calling prompt '#{prompt_name}': #{e.}" rescue StandardError => e raise MCPClient::Errors::PromptGetError, "Error calling prompt '#{prompt_name}': #{e.}" end |
#handle_elicitation_create(request_id, params) ⇒ void
This method returns an undefined value.
Handle elicitation/create request from server (MCP 2025-06-18)
863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 |
# File 'lib/mcp_client/server_stdio.rb', line 863 def handle_elicitation_create(request_id, params) # Without a callback there is no user to interact with: answer with a # JSON-RPC error rather than fabricating a user "decline". unless @elicitation_request_callback @logger.warn('Received elicitation request but no callback registered') send_error_response(request_id, -32_601, 'Elicitation not supported: no handler configured') return end # Call the registered callback result = @elicitation_request_callback.call(request_id, params) # Send the response back to the server (echoing related-task _meta) send_elicitation_response(request_id, (result, params)) end |
#handle_line(line) ⇒ void
This method returns an undefined value.
Handle a line of output from the stdio server Parses JSON-RPC messages and adds them to pending responses
336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 |
# File 'lib/mcp_client/server_stdio.rb', line 336 def handle_line(line) # The response is dated from the arrival of its line, before it is # decoded: parsing time is not freshness. arrived = respond_to?(:monotonic_now, true) ? monotonic_now : nil msg = JSON.parse(line) @logger.debug("Received line: #{(msg)}") # A JSON-parseable line that is not an object cannot be a JSON-RPC # message; skip it rather than raising inside the reader thread unless msg.is_a?(Hash) @logger.debug("Skipping non-object JSON-RPC line: #{line.chomp}") return end # Dispatch JSON-RPC requests from server (has id AND method) - MCP 2025-06-18 if msg['method'] && msg.key?('id') if modern_peer? # MCP 2026-07-28 stdio: "The server MUST NOT write JSON-RPC requests # to stdout" and "The client MUST NOT write JSON-RPC responses" — # server-to-client interactions travel in InputRequiredResult. @logger.warn("Ignoring server-initiated request #{msg['method']}: " \ 'a modern MCP server MUST NOT write JSON-RPC requests to stdout') else handle_server_request(msg) end return end # Dispatch JSON-RPC notifications (no id, has method) if msg['method'] && !msg.key?('id') route_notification(msg['method'], msg['params']) return end # Handle standard JSON-RPC responses (has id, no method) id = msg['id'] return unless id # A response to a subscriptions/listen request ends that subscription # (no caller is waiting on it). return if handle_subscription_response(msg) record_response(id, msg, arrived) rescue JSON::ParserError, EncodingError # Skip non-JSONRPC or undecodable lines in the output stream so a single # bad line cannot kill the reader thread end |
#handle_ping(request_id) ⇒ void
This method returns an undefined value.
Handle a server-initiated ping request (MCP ping utility) The receiver MUST respond promptly with an empty result.
850 851 852 853 854 855 856 857 |
# File 'lib/mcp_client/server_stdio.rb', line 850 def handle_ping(request_id) response = { 'jsonrpc' => '2.0', 'id' => request_id, 'result' => {} } (response) end |
#handle_reader_eof(session, generation = @transport_generation) ⇒ void
This method returns an undefined value.
EOF on the process's stdout: unless the client closed its stdin, the server exited on its own.
Waiting out an initialization still in flight is what makes an exit
during one recoverable. The handling used to be skipped outright while
@initialized was false — which is exactly the state a process that
answered the discovery probe and then exited leaves behind. Nothing else
noticed: MCPClient::ServerStdio::JsonRpcTransport#ensure_initialized went on to mark the dead
connection initialized and re-send the open subscriptions to it, the
failed writes were deferred back onto the queue for "the next process",
and with this reader already gone there was no one left to establish one.
The subscriptions stayed :reconnecting for ever, with the host neither
served nor told.
The lock is that initialization finishing, and taking it cannot deadlock: the reader is never the thread inside it, and by the time it is waiting here it can deliver no further responses, so anything that thread is still waiting for is already bounded by its own timeout.
274 275 276 277 278 279 |
# File 'lib/mcp_client/server_stdio.rb', line 274 def handle_reader_eof(session, generation = @transport_generation) @init_lock.synchronize { nil } unless @initialized return unless live_transport?(generation, session) handle_server_exit(session, generation) end |
#handle_roots_list(request_id, params) ⇒ void
This method returns an undefined value.
Handle roots/list request from server (MCP 2025-06-18)
883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 |
# File 'lib/mcp_client/server_stdio.rb', line 883 def handle_roots_list(request_id, params) # If no callback is registered, return empty roots list unless @roots_list_request_callback @logger.debug('Received roots/list request but no callback registered, returning empty list') send_roots_list_response(request_id, { 'roots' => [] }) return end # Call the registered callback result = @roots_list_request_callback.call(request_id, params) # Serving a roots/list answer that carries a root means this host # declared, and is using, the deprecated Roots capability (SEP-2577) — # with or without a Client. An empty answer is not use of it. warn_roots_deprecated(result) # Send the response back to the server (echoing related-task _meta) send_roots_list_response(request_id, (result, params)) end |
#handle_sampling_create_message(request_id, params) ⇒ void
This method returns an undefined value.
Handle sampling/createMessage request from server (MCP 2025-11-25)
906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 |
# File 'lib/mcp_client/server_stdio.rb', line 906 def (request_id, params) # If no callback is registered, return error unless @sampling_request_callback @logger.warn('Received sampling request but no callback registered, returning error') # sampling.mdx § Error Handling reserves -1 for "User rejected sampling # request"; a capability this client never declared is an unsupported # method (-32601, Method not found), as Client#handle_sampling_request answers. send_error_response(request_id, -32_601, 'Sampling not supported') return end # Sampling, and the includeContext values it may carry, are deprecated # (SEP-2577, SEP-2596) — with or without a Client. warn_sampling_deprecated(params) return if refused_undeclared_sampling_tools?(request_id, params) # Call the registered callback result = @sampling_request_callback.call(request_id, params) # Send the response back to the server (echoing related-task _meta) send_sampling_response(request_id, (result, params)) end |
#handle_server_exit(session = @session, generation = @transport_generation) ⇒ void
This method returns an undefined value.
The server process ended on its own (its stdout reached EOF). MCP 2026-07-28 stdio "Unexpected Termination": the client SHOULD restart it; in-flight requests are lost and subscriptions must be re-established. Marking the session uninitialized makes the next request spawn a fresh process and re-open live subscriptions — and a host that is only waiting for notifications makes no such request, so an open subscription restarts the process here instead.
1265 1266 1267 1268 1269 1270 1271 1272 1273 1274 1275 1276 1277 1278 |
# File 'lib/mcp_client/server_stdio.rb', line 1265 def handle_server_exit(session = @session, generation = @transport_generation) return unless live_transport?(generation, session) @logger.warn('MCP server process ended unexpectedly') # Stamped before the teardown, on the record the teardown retires: this # is the one path that knows the process was not asked to go. session&.exited_unexpectedly # Retired before the teardown bumps the generation past this reader's, # so the exit stays observable the way any other unexpected exit is: # the next request releases the dead handles and negotiates again. retire_transport(generation) teardown_transport(generation) restart_for_open_subscriptions end |
#handle_server_request(msg) ⇒ void
This method returns an undefined value.
Handle incoming JSON-RPC request from server (MCP 2025-06-18)
818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 |
# File 'lib/mcp_client/server_stdio.rb', line 818 def handle_server_request(msg) request_id = msg['id'] method = msg['method'] params = msg['params'] || {} @logger.debug("Received server request: #{method} (id: #{request_id})") case method when 'ping' handle_ping(request_id) when 'elicitation/create' handle_elicitation_create(request_id, params) when 'roots/list' handle_roots_list(request_id, params) when 'sampling/createMessage' (request_id, params) else # Unknown request method, send error response send_error_response(request_id, -32_601, "Method not found: #{method}") end rescue StandardError => e # The exception message is host-internal (file paths, connection # strings, library internals): log it locally, answer the peer with a # constant message, matching the SSE and Streamable HTTP transports. @logger.error("Error handling server request: #{e.}") send_error_response(request_id, -32_603, 'Internal error') end |
#list_prompts ⇒ Array<MCPClient::Prompt>
List all prompts available from the MCP server
443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 |
# File 'lib/mcp_client/server_stdio.rb', line 443 def list_prompts cached = hinted_list_value(:prompts) return cached if cached ensure_initialized # A cursor the server rejects restarts the list from its first page, # exactly as it does on the HTTP transports (MCP pagination). page = { cursor: nil } prompts = restarting_rejected_cursor('prompts', page) { collect_prompt_pages(page) } attach_list_value(:prompts, prompts) prompts rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::PromptGetError, "Error listing prompts: #{e.}" rescue StandardError => e raise MCPClient::Errors::PromptGetError, "Error listing prompts: #{e.}" end |
#list_resource_templates(cursor: nil) ⇒ Hash
List all resource templates available from the MCP server
565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 |
# File 'lib/mcp_client/server_stdio.rb', line 565 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 ensure_initialized 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 { |td| MCPClient::ResourceTemplate.from_json(td, server: self) } templates_result = { 'resourceTemplates' => templates, 'nextCursor' => result['nextCursor'] } attach_list_value(:templates, templates_result) unless cursor templates_result rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ResourceReadError, "Error listing resource templates: #{e.}" rescue StandardError => e raise MCPClient::Errors::ResourceReadError, "Error listing resource templates: #{e.}" end |
#list_resources(cursor: nil) ⇒ Hash
List all resources available from the MCP server
514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 |
# File 'lib/mcp_client/server_stdio.rb', line 514 def list_resources(cursor: nil) cached = cursor ? nil : hinted_list_value(:resources) return cached if cached ensure_initialized 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 { |td| MCPClient::Resource.from_json(td, server: self) } resources_result = { 'resources' => resources, 'nextCursor' => result['nextCursor'] } attach_list_value(:resources, resources_result) unless cursor resources_result rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ResourceReadError, "Error listing resources: #{e.}" rescue StandardError => e raise MCPClient::Errors::ResourceReadError, "Error listing resources: #{e.}" end |
#list_tools ⇒ Array<MCPClient::Tool>
List all tools available from the MCP server
654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 |
# File 'lib/mcp_client/server_stdio.rb', line 654 def list_tools # MCP 2026-07-28 caching: a list the server put a positive ttlMs on is # served here while it is still fresh, so a host reaching for the # transport directly does not re-list on every call. cached = hinted_list_value(:tools) return cached if cached ensure_initialized # A cursor the server rejects restarts the list from its first page, # exactly as it does on the HTTP transports (MCP pagination). page = { cursor: nil } tools = restarting_rejected_cursor('tools', page) { collect_tool_pages(page) } attach_list_value(:tools, tools) tools rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ToolCallError, "Error listing tools: #{e.}" rescue StandardError => e raise MCPClient::Errors::ToolCallError, "Error listing tools: #{e.}" end |
#live_transport?(generation, session) ⇒ Boolean
Whether the process a reader was started for is still the live one: its transport generation has not been claimed for teardown, and its record is the current session. Read under the transport lock, since a teardown claims the generation under it.
288 289 290 |
# File 'lib/mcp_client/server_stdio.rb', line 288 def live_transport?(generation, session) @transport_lock.synchronize { !@stdin.nil? && generation == @transport_generation } && current_session?(session) end |
#log_level=(level) ⇒ Hash
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)
757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 |
# File 'lib/mcp_client/server_stdio.rb', line 757 def log_level=(level) MCPClient::Deprecations.warn(:logging, @logger) ensure_initialized # MCP 2026-07-28 removed logging/setLevel: the level is declared per # request in _meta["io.modelcontextprotocol/logLevel"], so store it # and let every subsequent request carry it. 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::CapabilityError, ArgumentError raise rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ServerError, "Error setting log level: #{e.}" rescue StandardError => e raise MCPClient::Errors::ServerError, "Error setting log level: #{e.}" end |
#modern_peer? ⇒ Boolean
Whether a server-initiated request is prohibited traffic.
Judged by the ESTABLISHED era, not by protocol_version: during the server/discover probe the latter is only the version this client proposed. A legacy server MAY ping while the probe is unanswered and then wait for the response before doing anything else, so treating its request as prohibited modern traffic deadlocks the negotiation.
Two exceptions. A client configured protocol: :modern has already ruled out the legacy fallback that the accommodation exists for: it will never speak legacy, so it never runs a host callback for a server request nor writes the response back — not even while its own probe is still in flight. And once the reader has seen an answer only a modern server could have written, the server is modern whatever the negotiating thread has got round to applying.
435 436 437 |
# File 'lib/mcp_client/server_stdio.rb', line 435 def modern_peer? protocol_era == :modern || @protocol_mode == :modern || @modern_answer_received end |
#on_elicitation_request(&block) ⇒ void
This method returns an undefined value.
Register a callback for elicitation requests (MCP 2025-06-18)
785 786 787 |
# File 'lib/mcp_client/server_stdio.rb', line 785 def on_elicitation_request(&block) @elicitation_request_callback = block end |
#on_roots_list_request(&block) ⇒ void
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)
799 800 801 |
# File 'lib/mcp_client/server_stdio.rb', line 799 def on_roots_list_request(&block) @roots_list_request_callback = block end |
#on_sampling_request(&block) ⇒ void
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)
811 812 813 |
# File 'lib/mcp_client/server_stdio.rb', line 811 def on_sampling_request(&block) @sampling_request_callback = block end |
#park_open_subscriptions(claimed) ⇒ void
This method returns an undefined value.
Subscriptions do not survive the process: keep the ones the host still wants so they are re-sent once the process is re-established (basic/patterns/subscriptions "Graceful Closure"). They are moved to the pending list outside the registry lock, since a subscription being opened holds its own lock while taking that one.
Only the claimed process's, the way the pipe teardown is. Draining the
registry took whatever it held at that moment, and a teardown that got
here after a host request had established the replacement parked the
streams that replacement was already serving: their outstanding listen
ids were discarded with the dead process's, so close had nothing left
to cancel and the server went on serving a stream this client could no
longer name. A subscription carries the generation its listen went out
on (MCPClient::Subscription#with_open_id), which is the same question
the pipe asks, asked of the registry.
Enqueued through the one lock a deferred hand-over writes under too, since a listen write failing on the process being torn down lands in the window between the snapshot and the write (JsonRpcTransport#queue_subscriptions_of_ended_process, which also forgets the listen ids this process was holding).
1172 1173 1174 1175 1176 1177 1178 1179 1180 |
# File 'lib/mcp_client/server_stdio.rb', line 1172 def park_open_subscriptions(claimed) open_subscriptions = subscriptions_mutex.synchronize do mine, theirs = subscriptions.values.partition { |subscription| claimed_subscription?(subscription, claimed) } subscriptions.keep_if { |_, subscription| theirs.include?(subscription) } mine end open_subscriptions.each(&:mark_reconnecting) queue_subscriptions_of_ended_process(open_subscriptions.select(&:reconnectable?)) end |
#pin_pipe_encodings ⇒ void
This method returns an undefined value.
Pin the subprocess pipe encodings to UTF-8 instead of inheriting the process locale (Encoding.default_external). JSON-RPC messages MUST be UTF-8 encoded (MCP 2025-11-25 basic/transports.mdx); under a non-UTF-8 locale (e.g. LANG=C) a valid UTF-8 message would otherwise fail to decode and kill the reader thread. The server MAY also write UTF-8 to stderr, so that pipe is pinned as well.
197 198 199 200 201 |
# File 'lib/mcp_client/server_stdio.rb', line 197 def pin_pipe_encodings [@stdin, @stdout, @stderr].each do |io| io&.set_encoding(Encoding::UTF_8) end end |
#read_resource(uri) ⇒ Array<MCPClient::ResourceContent>
Read a resource by its URI
544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 |
# File 'lib/mcp_client/server_stdio.rb', line 544 def read_resource(uri) ensure_initialized # A null result reaches the shared guard as-is: it is a malformed # response, not an empty resource. 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.}" rescue MCPClient::Errors::TransportError raise rescue StandardError => e raise MCPClient::Errors::ResourceReadError, "Error reading resource '#{uri}': #{e.}" end |
#record_response(id, msg, arrived) ⇒ void
This method returns an undefined value.
Queue a response for the caller waiting on it.
388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 |
# File 'lib/mcp_client/server_stdio.rb', line 388 def record_response(id, msg, arrived) @mutex.synchronize do # Only retain a response that corresponds to an outstanding request. # Late responses (arriving after the caller timed out) and unsolicited # responses are dropped so @pending cannot grow without bound. if @awaiting.key?(id) # The answer is recorded as identifying the peer BEFORE it is # queued and before the next line is read: the thread waiting for # it may not run until after the server has written its next line, # and if that line is a request a modern server MUST NOT have # written, it must already be known as prohibited traffic — a # legacy accommodation is only owed while the probe is unanswered. # Only an OUTSTANDING request's answer says anything, though: the # response to a probe that timed out (and was cancelled) SHOULD be # ignored, and the session it fell back to is a 2025-11-25 one # whose server requests are still owed their responses. And only # while the era is being negotiated: an answer on an established # 2025-11-25 session renegotiates nothing, however modern its # shape, and that session's ping, roots, sampling and elicitation # requests stay owed their responses. @modern_answer_received = true if era_probe_in_flight? && identifies_modern_server?(msg) @pending[id] = msg # Dated from arrival: the waiter may wake much later. (@response_arrivals ||= {})[id] = arrived || monotonic_now if respond_to?(:monotonic_now, true) @cond.broadcast else @logger.debug("Discarding response for unknown or expired request id=#{id}") end end end |
#retire_transport(generation) ⇒ void
This method returns an undefined value.
The subprocess closed its stdout: it has exited (or its pipes were dropped), so no response will ever arrive on this transport again. MCP 2026-07-28 basic/transports/stdio ("Unexpected Termination") says a client SHOULD restart a server that terminated unexpectedly, so retire the handshake rather than let it describe a process that no longer exists: the next request releases these handles and negotiates again against a fresh subprocess.
Nothing is restarted or replayed from here. A request that was in flight may already have been executed server-side, so it fails as it would on any other broken transport; only the session is recoverable. A deliberate shutdown bumps the generation first, so its own reader reaching EOF is not mistaken for an unexpected exit.
238 239 240 241 242 243 244 245 246 247 248 |
# File 'lib/mcp_client/server_stdio.rb', line 238 def retire_transport(generation) return unless generation == @transport_generation # Waiters are woken: a request in flight on this transport will never # be answered, and should fail now rather than wait out its timeout. @mutex.synchronize do @transport_retired = true @cond.broadcast end @logger.debug('Server stdout closed; the transport will be re-established on the next request') end |
#send_elicitation_response(request_id, result) ⇒ void
This method returns an undefined value.
Send elicitation response back to server (MCP 2025-06-18)
965 966 967 968 969 970 971 972 973 974 975 976 977 978 979 980 |
# File 'lib/mcp_client/server_stdio.rb', line 965 def send_elicitation_response(request_id, result) # Error-shaped results become JSON-RPC error responses (e.g. -32602 for # an undeclared elicitation mode), mirroring the sampling error path. if result.is_a?(Hash) && result['error'] send_error_response(request_id, result['error']['code'] || -32_603, result['error']['message'] || 'Elicitation error') return end response = { 'jsonrpc' => '2.0', 'id' => request_id, 'result' => result } (response) end |
#send_error_response(request_id, code, message) ⇒ void
This method returns an undefined value.
Send error response back to server (MCP 2025-06-18)
987 988 989 990 991 992 993 994 995 996 997 |
# File 'lib/mcp_client/server_stdio.rb', line 987 def send_error_response(request_id, code, ) response = { 'jsonrpc' => '2.0', 'id' => request_id, 'error' => { 'code' => code, 'message' => } } (response) end |
#send_message(message) ⇒ void
This method returns an undefined value.
Send a JSON-RPC message to the server
1002 1003 1004 1005 1006 1007 1008 1009 |
# File 'lib/mcp_client/server_stdio.rb', line 1002 def () json = JSON.generate() @stdin.puts(json) @stdin.flush @logger.debug("Sent message: #{()}") rescue StandardError => e @logger.error("Error sending message: #{e.}") end |
#send_roots_list_response(request_id, result) ⇒ void
This method returns an undefined value.
Send roots/list response back to server (MCP 2025-06-18)
933 934 935 936 937 938 939 940 |
# File 'lib/mcp_client/server_stdio.rb', line 933 def send_roots_list_response(request_id, result) response = { 'jsonrpc' => '2.0', 'id' => request_id, 'result' => result } (response) end |
#send_sampling_response(request_id, result) ⇒ void
This method returns an undefined value.
Send sampling response back to server (MCP 2025-11-25)
946 947 948 949 950 951 952 953 954 955 956 957 958 959 |
# File 'lib/mcp_client/server_stdio.rb', line 946 def send_sampling_response(request_id, result) # Check if result contains an error if result.is_a?(Hash) && result['error'] send_error_response(request_id, result['error']['code'] || -1, result['error']['message'] || 'Sampling error') return end response = { 'jsonrpc' => '2.0', 'id' => request_id, 'result' => result } (response) end |
#send_subscription_cancellation(id, io: @stdin) ⇒ void
This method returns an undefined value.
Tell the server the client closed a subscriptions/listen request.
1347 1348 1349 1350 1351 1352 1353 1354 1355 |
# File 'lib/mcp_client/server_stdio.rb', line 1347 def send_subscription_cancellation(id, io: @stdin) return unless io && id notif = build_jsonrpc_notification('notifications/cancelled', { 'requestId' => id, 'reason' => 'Client closed subscription' }) io.puts(notif.to_json) rescue StandardError => e @logger.debug("Failed to send subscription cancellation: #{e.}") end |
#signal_server_process(signal, wait_thread = @wait_thread) ⇒ void
This method returns an undefined value.
Send a signal to the server process, tolerating a process that has already exited or cannot be signalled.
1361 1362 1363 1364 1365 |
# File 'lib/mcp_client/server_stdio.rb', line 1361 def signal_server_process(signal, wait_thread = @wait_thread) Process.kill(signal, wait_thread.pid) rescue Errno::ESRCH, Errno::EPERM => e @logger.debug("Could not send SIG#{signal} to server process: #{e.class}") end |
#spawn_server_process ⇒ Array
Returns the stdin, stdout, stderr and wait thread of the spawned process.
178 179 180 181 182 183 184 185 186 187 188 |
# File 'lib/mcp_client/server_stdio.rb', line 178 def spawn_server_process if @command_array return Open3.popen3(@env, *@command_array) if @env.any? Open3.popen3(*@command_array) elsif @env.any? Open3.popen3(@env, @command) else Open3.popen3(@command) end end |
#start_reader ⇒ Thread
Spawn a reader thread to collect JSON-RPC responses
205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 |
# File 'lib/mcp_client/server_stdio.rb', line 205 def start_reader generation = @transport_generation # The record of the process this reader belongs to, so an EOF can be # told from the EOF of a process that has since been replaced. session = @session stdout = @stdout @reader_thread = Thread.new do stdout.each_line do |line| handle_line(line) end handle_reader_eof(session, generation) rescue StandardError # Reader thread aborted unexpectedly ensure retire_transport(generation) end end |
#start_stderr_reader ⇒ Thread
Spawn a thread to continuously drain the subprocess stderr.
The child's stderr pipe has a fixed OS buffer (typically 64KB). If it is never read, a server that logs verbosely to stderr eventually blocks on write once the buffer fills, which stalls the whole subprocess (it stops producing stdout / reading stdin) and deadlocks the client. Draining stderr keeps the pipe empty; lines are surfaced at debug level.
Reads happen in bounded chunks (not IO#each_line) so that a server which writes to stderr without newline delimiters cannot make the client buffer a single "line" without limit: any pending fragment larger than STDERR_MAX_LINE_SIZE is flushed rather than retained.
314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 |
# File 'lib/mcp_client/server_stdio.rb', line 314 def start_stderr_reader stderr = @stderr @stderr_thread = Thread.new do buffer = +'' loop do buffer << stderr.readpartial(STDERR_READ_CHUNK_SIZE) flush_stderr_lines(buffer) flush_stderr_overflow(buffer) end rescue IOError # EOFError (a subclass of IOError) on EOF, or IOError on close; # emit any trailing partial line before exiting @logger.debug("[stderr] #{buffer.chomp}") if buffer && !buffer.empty? rescue StandardError # reader aborted unexpectedly; nothing actionable end end |
#subscribe_resource(uri) ⇒ Boolean
Subscribe to resource updates
599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 |
# File 'lib/mcp_client/server_stdio.rb', line 599 def subscribe_resource(uri) ensure_initialized 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::CapabilityError raise rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ResourceReadError, "Error subscribing to resource '#{uri}': #{e.}" rescue StandardError => e raise MCPClient::Errors::ResourceReadError, "Error subscribing to resource '#{uri}': #{e.}" end |
#teardown_transport(generation) ⇒ void
This method returns an undefined value.
Tear down the process of one transport generation.
The teardown claims the process under the transport lock — the handles
come off the transport and the generation moves on, so a request judged
current is written before this or not at all, and a second teardown of
the same process (the host's cleanup racing the reader's own, or the
reverse) finds nothing to claim. Everything after the claim works on the
claimed handles alone: a cleanup that ran ahead of it may already have
let the host establish a replacement, and a reader whose teardown
finishes only now used to clear that replacement's session, handles,
reader and handshake on its way out — leaving a live, re-sent
subscription registered on a transport that had just forgotten its
process.
1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 |
# File 'lib/mcp_client/server_stdio.rb', line 1087 def teardown_transport(generation) claimed = claim_transport(generation) return unless claimed begin park_open_subscriptions(claimed) terminate_server_process(claimed.wait_thread) claimed.stdout.close unless claimed.stdout.nil? || claimed.stdout.closed? claimed.stderr.close unless claimed.stderr.nil? || claimed.stderr.closed? # The reader calls this itself when the process exits on its own, and a # thread that kills itself here would abandon the rest of the shutdown — # including the restart that a live subscription depends on. It is at # EOF by then and returns on its own. claimed.reader_thread&.kill unless claimed.reader_thread.equal?(Thread.current) claimed.stderr_thread&.kill rescue StandardError # Clean up resources during unexpected termination ensure forget_torn_down_transport(claimed) end end |
#terminate_server_process(wait_thread = @wait_thread) ⇒ void
Terminate the spawned server process per the MCP 2025-11-25 stdio shutdown sequence (basic/lifecycle.mdx): stdin has already been closed, so wait for the process to exit on its own; if it does not exit within the grace period send SIGTERM, wait again, and finally send SIGKILL.
1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 |
# File 'lib/mcp_client/server_stdio.rb', line 1288 def terminate_server_process(wait_thread = @wait_thread) return unless wait_thread return if wait_thread.join(SHUTDOWN_GRACE_PERIOD) signal_server_process('TERM', wait_thread) return if wait_thread.join(SHUTDOWN_GRACE_PERIOD) signal_server_process('KILL', wait_thread) wait_thread.join(SHUTDOWN_GRACE_PERIOD) end |
#unsubscribe_resource(uri) ⇒ Boolean
Unsubscribe from resource updates
628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 |
# File 'lib/mcp_client/server_stdio.rb', line 628 def unsubscribe_resource(uri) ensure_initialized 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::CapabilityError raise rescue MCPClient::Errors::ServerError => e # 2026-07-28 protocol errors carry actionable data (requiredCapabilities, # supported versions); keep them intact instead of wrapping. raise if e.protocol_error? raise MCPClient::Errors::ResourceReadError, "Error unsubscribing from resource '#{uri}': #{e.}" rescue StandardError => e raise MCPClient::Errors::ResourceReadError, "Error unsubscribing from resource '#{uri}': #{e.}" end |