Module: MCPClient::SubscriptionSupport
- Included in:
- JsonRpcCommon
- Defined in:
- lib/mcp_client/subscription_support.rb
Overview
subscriptions/listen support shared by every transport (MCP 2026-07-28 basic/patterns/subscriptions): opening subscriptions, the registry keyed by listen request id, routing of tagged notifications, acknowledgment, graceful and abrupt closure, and the resources/subscribe mapping. Transports provide ensure_session_ready, open_subscription and cancel_subscription.
Constant Summary collapse
- DEFAULT_ACK_TIMEOUT =
Seconds to wait for a resource subscription acknowledgment on a transport without its own read timeout.
30- CONTROL_NOTIFICATIONS =
Notifications that are subscription bookkeeping rather than something a subscription's listeners are watching for.
%w[notifications/subscriptions/acknowledged notifications/cancelled].freeze
Instance Method Summary collapse
-
#await_acknowledgment_deadline(subscription, ack_timeout) ⇒ Thread?
Arrange for an unacknowledged listen to be given up on.
-
#close_subscription_gracefully(subscription, result) ⇒ void
End a subscription on the server's closing response — but only when the result is one the client recognizes, and only when it is a completion.
-
#confirm_resource_subscription(subscription, uri) ⇒ void
Wait for the acknowledgment and check that it really covers the URI: a server MAY acknowledge a subset of the filter it was sent.
- #deliver_subscription_notification(subscription, method, params) ⇒ void
-
#discard_mapped_resource_subscription(subscription, uri) ⇒ void
Give up on a stream that is mapped to a URI the server is not honouring on it — because it answered without the URI, ended, or never answered at all within the acknowledgment timeout.
-
#drop_unacknowledged_resource_subscriptions(subscription) ⇒ void
Close the streams whose mapped URI the server's latest acknowledgment left out.
- #ensure_modern_listen! ⇒ void
-
#expire_unacknowledged_subscription(subscription, request_id, timeout) ⇒ void
End a listen the server never acknowledged, and tell the server so — unless the subscription has moved on from that request, or the server answered it in the instant between the wait and this: the verdict and the closure are one step on the subscription, and whoever is waiting for the handle to settle is woken once the server has been told.
- #handle_server_cancellation(params) ⇒ void
-
#handle_subscription_acknowledgment(params) ⇒ void
Record what the server agreed to honour, and recheck the resource subscriptions this stream carries against it.
-
#handle_subscription_control(method, params) ⇒ void
Subscription bookkeeping carried by a notification: the server's acknowledgment of a listen request, and its teardown of one.
-
#handle_subscription_response(message) ⇒ MCPClient::Subscription?
Handle a JSON-RPC response addressed to a listen request: a result is the server's graceful closure, an error a failed subscription.
-
#listen(notifications:, ack_timeout: nil) {|method, params| ... } ⇒ MCPClient::Subscription
Open a long-lived notification stream.
-
#live_resource_subscription(uri) ⇒ MCPClient::Subscription?
Its stream, while the server is currently honouring it for that URI — see MCPClient::Subscription#watching_resource?.
-
#notify_host(method, params) ⇒ void
Hand a notification to the host's callback, surviving whatever it does with it (see #route_notification).
-
#open_listen(filter, ack_timeout:, initial_deadline: true) {|method, params| ... } ⇒ MCPClient::Subscription
Open a listen stream for a normalized filter.
-
#open_resource_subscription(uri) ⇒ MCPClient::Subscription
The acknowledged subscription.
-
#rearm_acknowledgment_deadline(subscription) ⇒ Thread?
Put the same deadline on a listen request a transport has just re-issued: an HTTP stream re-opened after a drop, or a stdio subscription re-sent to the process that replaced the one it was on.
-
#recheck_mapped_resource_subscription(subscription, uri) ⇒ void
Check the acknowledgment that stands now that the URI is mapped.
- #register_subscription(subscription) ⇒ void
-
#report_declined_subscription_types(subscription) ⇒ void
"The client SHOULD check the acknowledged filter against what it requested and handle any unsupported types gracefully" (basic/patterns/subscriptions).
-
#require_resource_watch(subscription, uri) ⇒ void
Wait for the server's word on this URI to stand, and demand that it watch it.
-
#resource_subscription_mutex(uri) ⇒ Mutex
The lock that serializes opening the stream for one URI.
- #resource_subscription_mutexes ⇒ Hash{String => Mutex}
-
#resource_subscriptions ⇒ Hash{String => MCPClient::Subscription}
Listen streams opened by subscribe_resource.
-
#route_notification(method, params) ⇒ void
Route an incoming notification, in this order: subscription bookkeeping (acknowledgment, server-side teardown), then transport and host cache invalidation, then the delivery to the owning subscription's listeners, and last the host's
on_notificationcallback (so hosts still see subscription-delivered notifications exactly like request-scoped ones). -
#settled_resource_subscription(uri) ⇒ MCPClient::Subscription?
The stream already mapped to this URI, once the server's word on it stands — or nil when there is none to reuse.
-
#subscribe_resource_via_listen(uri) ⇒ MCPClient::Subscription
Open resource-update subscriptions the modern way: one listen stream per URI (resources/subscribe was replaced by subscriptions/listen.resourceSubscriptions).
-
#subscription_ack_timeout ⇒ Numeric
Seconds allowed for a resource subscription acknowledgment.
- #subscription_by_id(id) ⇒ MCPClient::Subscription?
-
#subscription_delivery_target(method, params) ⇒ MCPClient::Subscription?
The subscription a notification is delivered to, resolved from the payload before the host's callback is given the chance to edit it (see #route_notification).
-
#subscription_for_notification(params) ⇒ MCPClient::Subscription?
The subscription a notification belongs to, from its io.modelcontextprotocol/subscriptionId.
-
#subscriptions ⇒ Hash{Integer, String => MCPClient::Subscription}
The subscriptions this transport has opened, keyed by the String form of their listen request id.
-
#subscriptions_mutex ⇒ Mutex
Guards the subscription registry.
-
#unmap_resource_subscription(subscription, uri) ⇒ void
Drop a URI's mapping, but only while it still names this stream.
- #unregister_subscription(subscription) ⇒ void
-
#unregister_subscription_id(subscription, id) ⇒ void
Drop the registration a particular listen id made, and only that one: a subscription re-opened under a newer id (by a reconnect, or by a stdio restart racing a blocked write) is registered under that newer id, and the older attempt must not delete it.
-
#unsubscribe_resource_via_listen(uri) ⇒ MCPClient::Subscription?
The subscription that was closed, if any.
Instance Method Details
#await_acknowledgment_deadline(subscription, ack_timeout) ⇒ Thread?
Arrange for an unacknowledged listen to be given up on.
Started only once the request is on its way, so nothing is cancelled before it exists; it waits on the subscription's own settling signal, so an acknowledgment (or any other end) retires it at once rather than leaving a thread asleep for the whole deadline. Every re-issued request gets one too (#rearm_acknowledgment_deadline).
The deadline is the request's, not the handle's: the id it is set on is taken here, and only that request is expired by it (MCPClient::Subscription#expire_unanswered). Waiting on the mutable handle instead expired whatever request it was on by the time the wait returned — a first request's timer closed the replacement a restart had issued since, naming the replacement and a deadline it had not missed.
126 127 128 129 130 131 132 133 134 135 136 137 138 |
# File 'lib/mcp_client/subscription_support.rb', line 126 def await_acknowledgment_deadline(subscription, ack_timeout) timeout = ack_timeout.nil? ? subscription_ack_timeout : ack_timeout return nil unless timeout.is_a?(Numeric) && timeout.positive? request_id = subscription.id Thread.new do Thread.current.name = 'MCP-listen-ack' Thread.current.report_on_exception = false next if subscription.wait_until_settled(timeout) expire_unacknowledged_subscription(subscription, request_id, timeout) end end |
#close_subscription_gracefully(subscription, result) ⇒ void
This method returns an undefined value.
End a subscription on the server's closing response — but only when the result is one the client recognizes, and only when it is a completion. Every other response goes through JsonRpcCommon#validate_result_type!; skipping it here would make a missing, scalar or unknown-resultType result indistinguishable from a clean close.
Recognized is not enough on its own. input_required is a resultType
this client accepts — on tools/call, resources/read and prompts/get,
the three requests a server may answer with one
(basic/patterns/mrtr "Supported Requests"). subscriptions/listen is not
among them, and the whole meaning of input_required is that the
request has not completed, so reporting one as a graceful closure told
the host the server had finished with a stream it had not.
437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 |
# File 'lib/mcp_client/subscription_support.rb', line 437 def close_subscription_gracefully(subscription, result) validate_result_type!(result) type = MCPClient::JsonRpcCommon.result_type(result) unless type == 'complete' raise MCPClient::Errors::InvalidResultError, "Invalid result: resultType #{type.inspect} does not close a subscription; " \ "input_required is only valid for #{MCPClient::JsonRpcCommon::MRTR_METHODS.join(', ')}, " \ 'not subscriptions/listen' end @logger.debug("Server closed subscription #{subscription.id} gracefully") subscription.finish(gracefully: true) rescue MCPClient::Errors::InvalidResultError => e @logger.warn("subscriptions/listen #{subscription.id} closed with an invalid result: #{e.}") subscription.finish(gracefully: false, error: e) end |
#confirm_resource_subscription(subscription, uri) ⇒ void
This method returns an undefined value.
Wait for the acknowledgment and check that it really covers the URI: a server MAY acknowledge a subset of the filter it was sent.
642 643 644 |
# File 'lib/mcp_client/subscription_support.rb', line 642 def confirm_resource_subscription(subscription, uri) require_resource_watch(subscription, uri) end |
#deliver_subscription_notification(subscription, method, params) ⇒ void
This method returns an undefined value.
397 398 399 |
# File 'lib/mcp_client/subscription_support.rb', line 397 def deliver_subscription_notification(subscription, method, params) subscription&.deliver(method, params) end |
#discard_mapped_resource_subscription(subscription, uri) ⇒ void
This method returns an undefined value.
Give up on a stream that is mapped to a URI the server is not honouring on it — because it answered without the URI, ended, or never answered at all within the acknowledgment timeout.
Dropping the mapping is not enough. A stream that is merely
:reconnecting is still reconnectable, so #subscribe_resource_via_listen
would open a replacement beside it and the discarded one could come back
and deliver the same updates a second time — while
#unsubscribe_resource_via_listen, which looks for the stream through
the very mapping that was just dropped, could no longer find or cancel
it. Closing it is also what the caller is entitled to: on stdio it sends
the notifications/cancelled the spec requires of a client that stops
reading a stream, and on Streamable HTTP it closes the response stream.
The mapping is dropped first, and the stream is closed only once no other URI still names it: a stream that is a live watch for a second resource is not this URI's to end.
528 529 530 531 532 533 534 535 536 |
# File 'lib/mcp_client/subscription_support.rb', line 528 def discard_mapped_resource_subscription(subscription, uri) unmap_resource_subscription(subscription, uri) still_mapped = subscriptions_mutex.synchronize do resource_subscriptions.any? { |_mapped_uri, sub| sub.equal?(subscription) } end return if still_mapped subscription.close end |
#drop_unacknowledged_resource_subscriptions(subscription) ⇒ void
This method returns an undefined value.
Close the streams whose mapped URI the server's latest acknowledgment left out.
Every acknowledgment is checked, not just the first: a stream re-opened
after an HTTP drop or a stdio restart is a new listen request, the
server holds no subscription state across it, and it MAY acknowledge a
smaller subset the second time. #confirm_resource_subscription only
guards the acknowledgment the subscriber waited for, so without this a
narrowed re-acknowledgment left resource_subscriptions mapping a URI
to a stream that no longer carried it, and
#live_resource_subscription kept reporting success for a resource
nothing was watching.
552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 |
# File 'lib/mcp_client/subscription_support.rb', line 552 def drop_unacknowledged_resource_subscriptions(subscription) missing = subscription.unacknowledged_resource_uris return if missing.empty? mapped = subscriptions_mutex.synchronize do resource_subscriptions.select { |uri, sub| sub.equal?(subscription) && missing.include?(uri) }.keys end return if mapped.empty? @logger.warn("Server re-acknowledged subscription #{subscription.id} without " \ "#{mapped.map { |uri| sanitize_log_text(uri) }.join(', ')}; closing it so the resource " \ 'subscription no longer reports a watch the server is not honouring') # Closing drops the mapping itself (the transport's cancel_subscription # clears every URI pointing at this stream), so a later subscribe_resource # opens a fresh stream and raises if that one is refused too. subscription.close end |
#ensure_modern_listen! ⇒ void
This method returns an undefined value.
45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 |
# File 'lib/mcp_client/subscription_support.rb', line 45 def ensure_modern_listen! # A transport with no listen stream of its own (the deprecated SSE # transport, a host adapter written against the older interface) cannot # serve one however the session was negotiated. unless respond_to?(:open_subscription, true) raise MCPClient::Errors::CapabilityError, "#{self.class.name} does not support subscriptions/listen (an MCP 2026-07-28 stdio or " \ 'Streamable HTTP transport is required)' end ensure_session_ready return if modern? raise MCPClient::Errors::CapabilityError, 'subscriptions/listen requires an MCP 2026-07-28 server; this server negotiated ' \ "#{protocol_version || 'no version'} (use resources/subscribe and server notifications instead)" end |
#expire_unacknowledged_subscription(subscription, request_id, timeout) ⇒ void
This method returns an undefined value.
End a listen the server never acknowledged, and tell the server so — unless the subscription has moved on from that request, or the server answered it in the instant between the wait and this: the verdict and the closure are one step on the subscription, and whoever is waiting for the handle to settle is woken once the server has been told.
149 150 151 152 153 154 155 156 157 158 159 160 161 |
# File 'lib/mcp_client/subscription_support.rb', line 149 def expire_unacknowledged_subscription(subscription, request_id, timeout) error = MCPClient::Errors::RequestTimeoutError.new( "subscriptions/listen #{request_id} was not acknowledged within #{timeout}s" ) subscription.expire_unanswered(request_id, error) do @logger.warn("subscriptions/listen #{request_id} was not acknowledged within #{timeout}s; cancelling it") # The handle is already closed, so this is the cancellation alone: the # notifications/cancelled on stdio, the closed response stream on HTTP. cancel_subscription(subscription) rescue StandardError => e @logger.debug("Cancelling an unacknowledged subscription raised #{e.class}: #{e.}") end end |
#handle_server_cancellation(params) ⇒ void
This method returns an undefined value.
355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 |
# File 'lib/mcp_client/subscription_support.rb', line 355 def handle_server_cancellation(params) return unless params.is_a?(Hash) # "Malformed notifications MAY be ignored": a reason that is not a # string is one, and so is a request id of another type than the one # this client issued — the registry is keyed by the id's text, so the # type is checked on the subscription found (basic/patterns/cancellation # "Error handling"). Neither may end a stream the server still serves. reason = params['reason'] return unless reason.nil? || reason.is_a?(String) subscription = subscription_by_id(params['requestId']) return unless subscription && subscription.id == params['requestId'] @logger.info("Server cancelled subscription #{subscription.id}: " \ "#{sanitize_log_text(reason || 'no reason given')}") unregister_subscription(subscription) subscription.finish(gracefully: false, reason: reason) end |
#handle_subscription_acknowledgment(params) ⇒ void
This method returns an undefined value.
Record what the server agreed to honour, and recheck the resource subscriptions this stream carries against it.
323 324 325 326 327 328 329 330 |
# File 'lib/mcp_client/subscription_support.rb', line 323 def handle_subscription_acknowledgment(params) subscription = subscription_for_notification(params) return @logger.debug('Acknowledgment for an unknown subscription ignored') unless subscription subscription.acknowledge(params['notifications']) report_declined_subscription_types(subscription) drop_unacknowledged_resource_subscriptions(subscription) end |
#handle_subscription_control(method, params) ⇒ void
This method returns an undefined value.
Subscription bookkeeping carried by a notification: the server's acknowledgment of a listen request, and its teardown of one.
308 309 310 311 312 313 314 315 316 317 |
# File 'lib/mcp_client/subscription_support.rb', line 308 def handle_subscription_control(method, params) case method when 'notifications/subscriptions/acknowledged' handle_subscription_acknowledgment(params) when 'notifications/cancelled' # Servers MUST send notifications/cancelled only to tear down a # subscriptions/listen stream (basic/patterns/cancellation). handle_server_cancellation(params) end end |
#handle_subscription_response(message) ⇒ MCPClient::Subscription?
Handle a JSON-RPC response addressed to a listen request: a result is the server's graceful closure, an error a failed subscription.
405 406 407 408 409 410 411 412 413 414 415 416 417 418 |
# File 'lib/mcp_client/subscription_support.rb', line 405 def handle_subscription_response() subscription = subscription_by_id(['id']) return nil unless subscription unregister_subscription(subscription) if ['error'] error = MCPClient::Errors::ServerError.from_jsonrpc(['error']) @logger.warn("subscriptions/listen #{subscription.id} failed: #{sanitize_log_text(error.)}") subscription.finish(gracefully: false, error: error) else close_subscription_gracefully(subscription, ['result']) end subscription end |
#listen(notifications:, ack_timeout: nil) {|method, params| ... } ⇒ MCPClient::Subscription
Open a long-lived notification stream. Modern servers only: the legacy transports keep resources/subscribe and the HTTP GET stream.
The request itself is meant to outlive every other one this client
sends — its response is the server's closing of the stream — so the
deadline the lifecycle asks for ("implementations SHOULD establish
timeouts for all sent requests", basic/patterns/cancellation "Timeouts")
is on the acknowledgment rather than on the response: a server MUST
acknowledge a listen before it sends anything on it, so a listen that
has not been acknowledged is a request nothing is happening on. One that
expires is cancelled the way that section requires, and the handle is
closed carrying the timeout, instead of staying :pending for the life
of the process with nothing to tell the host why.
37 38 39 40 41 |
# File 'lib/mcp_client/subscription_support.rb', line 37 def listen(notifications:, ack_timeout: nil, &listener) filter = MCPClient::Subscription.normalize_filter(notifications) ensure_modern_listen! open_listen(filter, ack_timeout: ack_timeout, &listener) end |
#live_resource_subscription(uri) ⇒ MCPClient::Subscription?
Returns its stream, while the server is currently honouring it for that URI — see MCPClient::Subscription#watching_resource?.
574 575 576 577 |
# File 'lib/mcp_client/subscription_support.rb', line 574 def live_resource_subscription(uri) existing = subscriptions_mutex.synchronize { resource_subscriptions[uri] } existing if existing&.watching_resource?(uri) end |
#notify_host(method, params) ⇒ void
This method returns an undefined value.
Hand a notification to the host's callback, surviving whatever it does with it (see #route_notification).
297 298 299 300 301 |
# File 'lib/mcp_client/subscription_support.rb', line 297 def notify_host(method, params) @notification_callback&.call(method, params) rescue StandardError => e @logger.warn("Notification callback error for #{sanitize_log_text(method)}: #{sanitize_log_text(e.)}") end |
#open_listen(filter, ack_timeout:, initial_deadline: true) {|method, params| ... } ⇒ MCPClient::Subscription
Open a listen stream for a normalized filter.
The deadline on the first request and the deadline on the requests a
reconnect or a restart re-issues are two settings, because one caller
wants them apart: subscribe_resource waits for the first
acknowledgment itself and reports its absence as its own failure, so
it wants no watchdog racing that wait — but the re-issued requests are
ones nobody is waiting on, and those it wants bounded like any other
(see #open_resource_subscription).
79 80 81 82 83 84 |
# File 'lib/mcp_client/subscription_support.rb', line 79 def open_listen(filter, ack_timeout:, initial_deadline: true, &listener) subscription = MCPClient::Subscription.new(server: self, requested: filter, ack_timeout: ack_timeout, &listener) open_subscription(subscription) await_acknowledgment_deadline(subscription, ack_timeout) if initial_deadline subscription end |
#open_resource_subscription(uri) ⇒ MCPClient::Subscription
Returns the acknowledged subscription.
581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 |
# File 'lib/mcp_client/subscription_support.rb', line 581 def open_resource_subscription(uri) # No watchdog on the first request: this caller waits for the # acknowledgment itself, on the same timeout, and reports a stream that # never arrives as its own failure rather than through a handle # something else closed. # # The requests a restart or a reconnect re-issues are another matter. # Nobody waits on those — the host is waiting for updates — so a # replacement the server accepted and then never acknowledged was, with # `ack_timeout: false`, pending for ever: the mapped stream was only # discarded the next time the URI was asked about, and a host that never # asked again was never told. Those requests are bounded by the # transport's own timeout, like any other listen's. ensure_modern_listen! subscription = open_listen(MCPClient::Subscription.normalize_filter('resourceSubscriptions' => [uri]), ack_timeout: subscription_ack_timeout, initial_deadline: false) begin confirm_resource_subscription(subscription, uri) subscriptions_mutex.synchronize { resource_subscriptions[uri] = subscription } recheck_mapped_resource_subscription(subscription, uri) rescue StandardError # Nothing watches a stream the caller could not use. unmap_resource_subscription(subscription, uri) subscription.close raise end subscription end |
#rearm_acknowledgment_deadline(subscription) ⇒ Thread?
Put the same deadline on a listen request a transport has just re-issued: an HTTP stream re-opened after a drop, or a stdio subscription re-sent to the process that replaced the one it was on.
Each of those is a new JSON-RPC request, for which the server holds no
subscription state and which it has to acknowledge afresh, so the
"implementations SHOULD establish timeouts for all sent requests"
(basic/patterns/cancellation "Timeouts") that bounded the first bounds it
too. It used to bound only the first: the watchdog listen started
retires at the first acknowledgment, and a replacement the server
accepted and then never acknowledged left the handle :pending with
nothing to tell the host why — indefinitely on stdio, and for as long as
the peer kept sending SSE comments on Streamable HTTP.
Watchdogs do not pile up behind a stream that keeps dropping: each one is bound to the request it was armed for, and retires as soon as the subscription has moved on to a newer one.
105 106 107 |
# File 'lib/mcp_client/subscription_support.rb', line 105 def rearm_acknowledgment_deadline(subscription) await_acknowledgment_deadline(subscription, subscription.ack_timeout) end |
#recheck_mapped_resource_subscription(subscription, uri) ⇒ void
This method returns an undefined value.
Check the acknowledgment that stands now that the URI is mapped.
#drop_unacknowledged_resource_subscriptions can only see a stream through the mapping, and the mapping is written after the acknowledgment the subscriber waited for. A stream that dropped and was re-opened in between is acknowledged afresh and MAY be granted more narrowly, so without this second look a re-acknowledgment that landed in that window was stored as a live watch nothing was honouring.
622 623 624 |
# File 'lib/mcp_client/subscription_support.rb', line 622 def recheck_mapped_resource_subscription(subscription, uri) require_resource_watch(subscription, uri) end |
#register_subscription(subscription) ⇒ void
This method returns an undefined value.
190 191 192 |
# File 'lib/mcp_client/subscription_support.rb', line 190 def register_subscription(subscription) subscriptions_mutex.synchronize { subscriptions[subscription.id] = subscription } end |
#report_declined_subscription_types(subscription) ⇒ void
This method returns an undefined value.
"The client SHOULD check the acknowledged filter against what it
requested and handle any unsupported types gracefully"
(basic/patterns/subscriptions). Gracefully, for the notification types
of a plain listen, is: the stream stays up for what was granted,
MCPClient::Subscription#unsupported names the rest, and the host is
told — a host that opted in to prompt-list changes and was granted
tool-list changes alone would otherwise keep waiting for notifications
that are never coming, with the handle :active and nothing said. The
resource URIs are held to more than a log line
(#drop_unacknowledged_resource_subscriptions): a watch the server
declined is one nothing is watching.
345 346 347 348 349 350 351 |
# File 'lib/mcp_client/subscription_support.rb', line 345 def report_declined_subscription_types(subscription) declined = subscription.unsupported return if declined.empty? @logger.warn("Server acknowledged subscription #{subscription.id} without #{declined.join(', ')}; " \ 'no notifications of those types will arrive on it') end |
#require_resource_watch(subscription, uri) ⇒ void
This method returns an undefined value.
Wait for the server's word on this URI to stand, and demand that it watch it.
The word that stands is the acknowledgment of the listen request the stream is currently on. A re-open clears it — the server holds no subscription state across one and has to grant the filter again — so a replacement in flight is waited for rather than read as an answer. Reading it as one is what used to report a watch nobody had made: a subscription with no acknowledgment has no unacknowledged URIs either, so "the URI is not missing from the acknowledgment" passed for "the server is watching it".
661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 |
# File 'lib/mcp_client/subscription_support.rb', line 661 def require_resource_watch(subscription, uri) case subscription.await_resource_watch(uri, subscription_ack_timeout) when :watching nil when :not_watching raise MCPClient::Errors::ResourceReadError, "the server acknowledged the subscription without '#{uri}'" when :closed raise subscription.error if subscription.error raise MCPClient::Errors::ResourceReadError, "the server closed the subscription for '#{uri}'" else raise MCPClient::Errors::ResourceReadError, "timed out after #{subscription_ack_timeout}s waiting for the server to acknowledge '#{uri}'" end end |
#resource_subscription_mutex(uri) ⇒ Mutex
The lock that serializes opening the stream for one URI. Kept for the life of the transport: it is one Mutex per URI the host subscribed to.
688 689 690 |
# File 'lib/mcp_client/subscription_support.rb', line 688 def resource_subscription_mutex(uri) subscriptions_mutex.synchronize { resource_subscription_mutexes[uri] ||= Mutex.new } end |
#resource_subscription_mutexes ⇒ Hash{String => Mutex}
693 694 695 |
# File 'lib/mcp_client/subscription_support.rb', line 693 def resource_subscription_mutexes @resource_subscription_mutexes ||= {} end |
#resource_subscriptions ⇒ Hash{String => MCPClient::Subscription}
Returns listen streams opened by subscribe_resource.
711 712 713 |
# File 'lib/mcp_client/subscription_support.rb', line 711 def resource_subscriptions @resource_subscriptions ||= {} end |
#route_notification(method, params) ⇒ void
This method returns an undefined value.
Route an incoming notification, in this order:
subscription bookkeeping (acknowledgment, server-side teardown), then
transport and host cache invalidation, then the delivery to the owning
subscription's listeners, and last the host's on_notification callback
(so hosts still see subscription-delivered notifications exactly like
request-scoped ones).
The invalidations come first on purpose, the transport's and the host's
alike. A listener runs on the subscription's own dispatcher thread, so
queuing its delivery makes the notification visible at once — and a
listener that reacts to a list_changed notification by calling a cached
list method (client.list_tools, say) would then read the very entry the
notification says is stale. Dropping the caches before the delivery is
queued makes "the caches are already invalid when a listener sees the
notification" a guarantee instead of a race the scheduler usually wins.
That is why the host's invalidation has a hook of its own
(MCPClient::ServerBase#on_cache_invalidation) rather than riding on the
host callback below: while it did, the guarantee held only for the
transport's own caches and the host's were dropped after the delivery.
The host callback comes last, because it is the only step that can block. It is host code driven by the peer and it runs on whatever thread is routing — on stdio the process's sole stdout reader — so a callback that issues a synchronous request of its own waits there for a response only that reader can deliver. Round 3 moved the subscription's listeners off that thread for exactly this reason; running the callback ahead of them put the queueing back behind it, and a host handler that blocked stalled a delivery the dispatcher would otherwise have made at once. Queueing costs nothing to move: #deliver_subscription_notification hands the notification to the dispatcher rather than to the listeners, so nothing host-supplied runs before the callback either way.
Being last, the callback can prevent nothing. An exception escaping it
used to take the notification down with it — the subscription's
listeners never saw something the host's own handler had already been
told about — and, on stdio, the transport's reader thread with it; it is
now logged and routing is over anyway. Nor can it drop or redirect a
delivery by editing the payload it is handed: that is the very hash the
delivery was routed by, and by the time the callback can touch it the
subscription has already been resolved and the entry queued, so
deleting or rewriting _meta changes nothing about where it went.
272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 |
# File 'lib/mcp_client/subscription_support.rb', line 272 def route_notification(method, params) handle_subscription_control(method, params) # notifications/message is the Logging utility, Deprecated as a whole in # 2026-07-28 (SEP-2577). The notice belongs here rather than in # MCPClient::Client: a host that registered on_notification on the # transport itself receives log messages without a Client ever existing. warn_logging_deprecated if method == 'notifications/message' invalidate_cache_for_notification(method, params) # The host's caches go with the transport's, on their own hook rather # than on the host callback below: that callback is deliberately last — # it is the step that may block — and a client whose invalidation rode on # it dropped its entries only after the delivery had been queued, so a # listener could read the very list the notification says is stale. Only # the invalidation is moved ahead; everything else the host does with a # notification is still behind the delivery. notify_cache_invalidation(method, params) deliver_subscription_notification(subscription_delivery_target(method, params), method, params) notify_host(method, params) end |
#settled_resource_subscription(uri) ⇒ MCPClient::Subscription?
The stream already mapped to this URI, once the server's word on it stands — or nil when there is none to reuse.
A mapped stream is not a watch merely for being open, and it is not one merely for having been granted the URI once. After an HTTP connection drops or a stdio process restarts, the request that replaces it is a new listen the server holds no state for: it may be rejected, or acknowledged without this URI, and until it is answered nothing has been granted. Reporting success from that state — which used to happen for every handle that was not closed, and then for every handle whose old acknowledgment was still on record, so for the whole of an HTTP backoff or a stdio handshake — tells the subscriber about a watch nobody has made yet. So this asks whether the server is watching the URI now (MCPClient::Subscription#await_live_resource_watch), waiting out a stream that is between listen attempts rather than reading what the last one was granted, and a stream that comes back without the URI — or does not come back at all — stops being this URI's stream and is closed with it (see #discard_mapped_resource_subscription).
498 499 500 501 502 503 504 505 506 |
# File 'lib/mcp_client/subscription_support.rb', line 498 def settled_resource_subscription(uri) mapped = subscriptions_mutex.synchronize { resource_subscriptions[uri] } return nil unless mapped return mapped if mapped.watching_resource?(uri) return mapped if mapped.await_live_resource_watch(uri, subscription_ack_timeout) == :watching discard_mapped_resource_subscription(mapped, uri) nil end |
#subscribe_resource_via_listen(uri) ⇒ MCPClient::Subscription
Open resource-update subscriptions the modern way: one listen stream per URI (resources/subscribe was replaced by subscriptions/listen.resourceSubscriptions). Blocks until the server acknowledges the stream, so the caller learns about a rejection the way it did from resources/subscribe.
463 464 465 466 467 468 469 470 471 472 473 474 475 476 |
# File 'lib/mcp_client/subscription_support.rb', line 463 def subscribe_resource_via_listen(uri) existing = live_resource_subscription(uri) return existing if existing # One stream per URI: two threads subscribing to the same resource must # not open two, or the second registration would hide the first and # unsubscribe_resource would close only one of them. resource_subscription_mutex(uri).synchronize do existing = settled_resource_subscription(uri) next existing if existing open_resource_subscription(uri) end end |
#subscription_ack_timeout ⇒ Numeric
Returns seconds allowed for a resource subscription acknowledgment.
679 680 681 682 |
# File 'lib/mcp_client/subscription_support.rb', line 679 def subscription_ack_timeout timeout = defined?(@read_timeout) ? @read_timeout : nil timeout.is_a?(Numeric) && timeout.positive? ? timeout : DEFAULT_ACK_TIMEOUT end |
#subscription_by_id(id) ⇒ MCPClient::Subscription?
182 183 184 185 186 |
# File 'lib/mcp_client/subscription_support.rb', line 182 def subscription_by_id(id) return nil if id.nil? subscriptions_mutex.synchronize { subscriptions[id] } end |
#subscription_delivery_target(method, params) ⇒ MCPClient::Subscription?
The subscription a notification is delivered to, resolved from the payload before the host's callback is given the chance to edit it (see #route_notification).
381 382 383 384 385 386 387 388 389 390 |
# File 'lib/mcp_client/subscription_support.rb', line 381 def subscription_delivery_target(method, params) return nil if CONTROL_NOTIFICATIONS.include?(method) = params.is_a?(Hash) ? params['_meta'] : nil return nil unless .is_a?(Hash) && .key?(MCPClient::JsonRpcCommon::META_SUBSCRIPTION_ID) subscription = subscription_by_id([MCPClient::JsonRpcCommon::META_SUBSCRIPTION_ID]) @logger.debug("Notification #{sanitize_log_text(method)} for an unknown subscription ignored") unless subscription subscription end |
#subscription_for_notification(params) ⇒ MCPClient::Subscription?
The subscription a notification belongs to, from its io.modelcontextprotocol/subscriptionId.
217 218 219 220 221 222 |
# File 'lib/mcp_client/subscription_support.rb', line 217 def subscription_for_notification(params) = params.is_a?(Hash) ? params['_meta'] : nil return nil unless .is_a?(Hash) && .key?(MCPClient::JsonRpcCommon::META_SUBSCRIPTION_ID) subscription_by_id([MCPClient::JsonRpcCommon::META_SUBSCRIPTION_ID]) end |
#subscriptions ⇒ Hash{Integer, String => MCPClient::Subscription}
The subscriptions this transport has opened, keyed by the String form of their listen request id. Keyed by the JSON-RPC id the listen went out with, exactly as it was sent. A JSON-RPC id of another type is another request's id — the cancellation path has always compared them exactly, and the acknowledgment and delivery paths do too — so a peer that tags a message with "9" does not reach the subscription listening on 9.
171 172 173 |
# File 'lib/mcp_client/subscription_support.rb', line 171 def subscriptions @subscriptions ||= {} end |
#subscriptions_mutex ⇒ Mutex
Returns guards the subscription registry.
176 177 178 |
# File 'lib/mcp_client/subscription_support.rb', line 176 def subscriptions_mutex @subscriptions_mutex ||= Mutex.new end |
#unmap_resource_subscription(subscription, uri) ⇒ void
This method returns an undefined value.
Drop a URI's mapping, but only while it still names this stream.
630 631 632 633 634 |
# File 'lib/mcp_client/subscription_support.rb', line 630 def unmap_resource_subscription(subscription, uri) subscriptions_mutex.synchronize do resource_subscriptions.delete(uri) if resource_subscriptions[uri].equal?(subscription) end end |
#unregister_subscription(subscription) ⇒ void
This method returns an undefined value.
196 197 198 |
# File 'lib/mcp_client/subscription_support.rb', line 196 def unregister_subscription(subscription) subscriptions_mutex.synchronize { subscriptions.delete(subscription.id) } end |
#unregister_subscription_id(subscription, id) ⇒ void
This method returns an undefined value.
Drop the registration a particular listen id made, and only that one: a subscription re-opened under a newer id (by a reconnect, or by a stdio restart racing a blocked write) is registered under that newer id, and the older attempt must not delete it.
207 208 209 210 211 |
# File 'lib/mcp_client/subscription_support.rb', line 207 def unregister_subscription_id(subscription, id) subscriptions_mutex.synchronize do subscriptions.delete(id) if subscriptions[id].equal?(subscription) end end |
#unsubscribe_resource_via_listen(uri) ⇒ MCPClient::Subscription?
Returns the subscription that was closed, if any.
699 700 701 702 703 704 705 706 707 708 |
# File 'lib/mcp_client/subscription_support.rb', line 699 def unsubscribe_resource_via_listen(uri) # Behind the same per-URI lock as subscribing, so an unsubscribe cannot # slip between a subscribe's acknowledgment and its registration and # leave the stream running. resource_subscription_mutex(uri).synchronize do subscription = subscriptions_mutex.synchronize { resource_subscriptions.delete(uri) } subscription&.close subscription end end |