Class: MCPClient::Subscription
- Inherits:
-
Object
- Object
- MCPClient::Subscription
- Defined in:
- lib/mcp_client/subscription.rb,
lib/mcp_client/subscription/notification_dispatcher.rb
Overview
A long-lived notification subscription opened with subscriptions/listen
(MCP 2026-07-28 basic/patterns/subscriptions).
The subscription is identified by the JSON-RPC id of its listen request;
every notification delivered on it carries that id in
_meta["io.modelcontextprotocol/subscriptionId"]. It starts :pending,
becomes :active when the server acknowledges it (with the subset of
notification types it agreed to honour), and ends :closed — gracefully
when the server answers the listen request, otherwise on a transport
drop, a server notifications/cancelled, an error, or #close.
Defined Under Namespace
Classes: NotificationDispatcher
Constant Summary collapse
- FILTER_FIELDS =
The published SubscriptionFilter's fields and their value types (MCP 2026-07-28 schema, SubscriptionFilter). Nothing beyond these four is a core filter field: an extension that defines one of its own — the tasks extension's
taskIds, say — registers it with register_filter_field. { 'toolsListChanged' => :boolean, 'promptsListChanged' => :boolean, 'resourcesListChanged' => :boolean, 'resourceSubscriptions' => :string_array }.freeze
- FILTER_ALIASES =
snake_case spellings accepted for the published filter fields
{ 'tools_list_changed' => 'toolsListChanged', 'prompts_list_changed' => 'promptsListChanged', 'resources_list_changed' => 'resourcesListChanged', 'resource_subscriptions' => 'resourceSubscriptions' }.freeze
- FILTER_VALUE_TYPES =
The value types a filter field may have
i[boolean string_array].freeze
- STATES =
i[pending active reconnecting closed].freeze
- MAX_PENDING_NOTIFICATIONS =
Ceiling on the notifications waiting for this subscription's listeners. The queue is filled by the peer and drained by the host, so a chatty server and a listener that does real work (re-reading the resource that changed, say) would otherwise grow it without bound.
A full queue discards by identity rather than by arrival order, so overflow costs a listener a repeated notice of the same resource or task and never its only notice of one of them; see NotificationDispatcher for the policy. Blocking the transport reader instead would reinstate the deadlock the dispatcher exists to prevent — the reader would wait for a listener that is waiting for a response only that reader can deliver.
1024- MAX_PENDING_NOTIFICATION_BYTES =
Ceiling on the bytes the queued notifications retain, because a count is not a memory bound: the method name and params of every queued notification are held until its listener has run, one Streamable HTTP listen event may approach HttpTransportBase::ListenStream::LISTEN_MAX_BUFFER_BYTES and a stdio line has no inbound limit at all — so a peer facing a slow listener could put tens of gigabytes behind a nominally bounded queue.
This changes when overflow starts, not what it discards: whichever ceiling the arriving notification would breach, the entry that goes is still chosen by identity, and only ever an entry whose removal relieves the breach. A notification larger than the whole budget is not charged against it and is held in a slot of its own, of which there is only ever one — so the retained total is the budget plus at worst one peer-sized payload rather than MAX_PENDING_NOTIFICATIONS of them, and no signal is lost to its size or displaced by one.
8 * 1024 * 1024
- SETTLED_STATES =
The states #wait_until_settled waits for: the server has answered the listen request one way or the other.
i[active closed].freeze
Instance Attribute Summary collapse
-
#ack_timeout ⇒ Numeric, ...
readonly
The acknowledgment deadline the host asked for; every listen request re-issued for this subscription is bounded by it, not only the first.
-
#acknowledged ⇒ Hash?
readonly
The filter the server agreed to honour, once acknowledged.
-
#close_reason ⇒ String?
readonly
The reason a server-side teardown gave, if any.
-
#error ⇒ MCPClient::Errors::MCPError?
readonly
Why the subscription failed, if it did.
-
#id ⇒ Integer, ...
readonly
The JSON-RPC id of the listen request (nil before it is sent).
-
#requested ⇒ Hash
readonly
The requested SubscriptionFilter (camelCase keys).
-
#server ⇒ MCPClient::ServerBase
readonly
The transport that owns the subscription.
-
#state ⇒ Symbol
readonly
:pending, :active, :reconnecting or :closed.
Class Method Summary collapse
-
.deep_frozen_copy(value) ⇒ Object
A detached, deeply frozen copy of a parsed JSON value.
-
.filter_aliases ⇒ Hash{String => String}
Every accepted snake_case spelling.
-
.filter_fields ⇒ Hash{String => Symbol}
Every accepted filter field — the published four and the registered extension fields — with its type.
-
.normalize_filter(filter) ⇒ Hash
Normalize and validate a SubscriptionFilter given with String or Symbol, camelCase or snake_case keys.
-
.register_filter_field(name, type, alias_name: nil) ⇒ void
Register a filter field an extension defines beyond the published SubscriptionFilter, so Subscription.normalize_filter accepts it (and its snake_case spelling) with the value type the extension gives it.
Instance Method Summary collapse
-
#acknowledge(filter) ⇒ void
private
Record what the server agreed to honour.
-
#active? ⇒ Boolean
Whether the server acknowledged it and it is still open.
-
#assign_id(id) ⇒ void
private
Put the subscription on a listen id, without the registration and the write #with_open_id holds its lock across.
-
#await_live_resource_watch(uri, timeout) ⇒ Symbol
Block until this subscription is a running watch of the URI, and report whether it is.
-
#await_resource_watch(uri, timeout) ⇒ Symbol
Block until the server has answered the listen request this subscription is waiting on, and report what it said about this URI.
-
#close ⇒ MCPClient::Subscription?
Cancel the subscription: the transport closes the stream (HTTP) or sends notifications/cancelled (stdio).
- #closed? ⇒ Boolean
-
#closed_by_client? ⇒ Boolean
Whether #close ended it.
-
#closed_gracefully? ⇒ Boolean
Whether the server ended it with a response to the listen request.
- #deliver(method, params) ⇒ Object private
-
#discard_outstanding_listens ⇒ void
private
Forget the recorded listen ids without cancelling them: the session they were written to is gone, so nothing is outstanding and none of them must be cancelled on the session that replaces it.
-
#dropped_notifications ⇒ Integer
Notifications dropped because the listeners could not keep up with the server (see MAX_PENDING_NOTIFICATIONS and MAX_PENDING_NOTIFICATION_BYTES).
-
#expire_unanswered(id, error) { ... } ⇒ Symbol?
private
End the request a deadline was set on, if the subscription is still on that request and the server has still not answered it — in one step.
-
#fail_attempt(id, error) ⇒ Symbol
private
Undo the listen attempt that failed — or find that it is no longer this attempt's to undo — in one step.
- #finish(gracefully: false, by_client: false, error: nil, reason: nil) ⇒ Object private
-
#initialize(server:, requested:, ack_timeout: nil) {|method, params| ... } ⇒ Subscription
constructor
A new instance of Subscription.
- #inspect ⇒ Object
-
#mark_listen_written(id) ⇒ void
private
The write of a listen request has finished — sent, or raised having sent who knows how much.
- #mark_reconnecting ⇒ Object private
-
#on_notification {|method, params| ... } ⇒ self
Add a listener for notifications delivered on this subscription.
-
#open_as?(id) ⇒ Boolean
private
Whether this subscription is still the stream a given listen id opened.
-
#open_generation ⇒ Integer?
private
The transport generation the listen this subscription is on went out under, or nil on a transport that does not number its processes.
-
#pending_notification_bytes ⇒ Integer
Bytes retained by the notifications queued for this subscription's listeners, measured as the JSON the peer sent for them — their method names as well as their params (see MAX_PENDING_NOTIFICATION_BYTES).
-
#pending_notifications ⇒ Integer
Notifications queued for this subscription's listeners.
-
#reconnectable? ⇒ Boolean
private
Whether the subscription should be re-established after a reconnect.
-
#reconnecting? ⇒ Boolean
Whether it is waiting for a transport to re-establish it — the stream dropped, or the stdio process it was on exited, and the transport that noticed has queued it for the next session.
-
#record_outstanding_listen(id, io = nil) ⇒ void
private
Record a listen request the transport has written for this subscription on the session it is on.
-
#reestablishing? ⇒ Boolean
private
Whether a transport is handing this subscription to a new session: it was #mark_reconnectinged and no server has acknowledged it since.
-
#take_outstanding_listens(io = nil) ⇒ Array
private
The listen ids the server may still be serving for this subscription, leaving none behind: the caller is cancelling them.
-
#unacknowledged_resource_uris ⇒ Array<String>
Requested resource URIs the server did not agree to watch.
-
#unsupported ⇒ Array<String>
Requested notification types the server did not agree to honour.
-
#wait_until_settled(timeout) ⇒ Symbol?
Block until the server has settled the subscription: acknowledged it, or ended it with a response, an error or a cancellation.
-
#watching_resource?(uri) ⇒ Boolean
Whether this stream is, right now, an active acknowledged watch of a resource: the server granted that URI and the stream it granted it on is the one still running.
-
#with_open_id(id, generation = nil) { ... } ⇒ Boolean
private
Take a fresh listen id and register/send the request under this subscription's own lock, so a concurrent #close either wins outright (nothing is sent) or waits and then cancels the id that was sent.
Constructor Details
#initialize(server:, requested:, ack_timeout: nil) {|method, params| ... } ⇒ Subscription
Returns a new instance of Subscription.
192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 |
# File 'lib/mcp_client/subscription.rb', line 192 def initialize(server:, requested:, ack_timeout: nil, &listener) @server = server @requested = requested @ack_timeout = ack_timeout @listeners = [] @listeners << listener if listener @state = :pending @mutex = Mutex.new @settled = ConditionVariable.new @id = nil @acknowledged = nil @error = nil @close_reason = nil @closed_gracefully = false @closed_by_client = false @dispatcher = nil # Whether the server has answered the listen request this subscription # is on — or the last one it was on, while no replacement has gone out. # Written where those two things happen ({#acknowledge} and the two # methods that take a new listen id), never inferred from how far a # transport has got through a reconnect. @answered = false # Whether a transport is handing this subscription to a new session: # see {#reestablishing?}. @reestablishing = false # The listen ids the transport has written for this subscription, and # not yet cancelled, each paired with the pipe it was written to: see # {#record_outstanding_listen}. @outstanding_listens = [] # Of those, the ones whose write has not finished yet. They are not # cancellable: see {#take_outstanding_listens}. @unwritten_listens = [] # The transport generation this subscription's listen went out on: see # {#with_open_id}. @open_generation = nil end |
Instance Attribute Details
#ack_timeout ⇒ Numeric, ... (readonly)
Returns the acknowledgment deadline the host asked for; every listen request re-issued for this subscription is bounded by it, not only the first.
147 148 149 |
# File 'lib/mcp_client/subscription.rb', line 147 def ack_timeout @ack_timeout end |
#acknowledged ⇒ Hash? (readonly)
Returns the filter the server agreed to honour, once acknowledged.
135 136 137 |
# File 'lib/mcp_client/subscription.rb', line 135 def acknowledged @acknowledged end |
#close_reason ⇒ String? (readonly)
Returns the reason a server-side teardown gave, if any.
143 144 145 |
# File 'lib/mcp_client/subscription.rb', line 143 def close_reason @close_reason end |
#error ⇒ MCPClient::Errors::MCPError? (readonly)
Returns why the subscription failed, if it did.
141 142 143 |
# File 'lib/mcp_client/subscription.rb', line 141 def error @error end |
#id ⇒ Integer, ... (readonly)
Returns the JSON-RPC id of the listen request (nil before it is sent).
131 132 133 |
# File 'lib/mcp_client/subscription.rb', line 131 def id @id end |
#requested ⇒ Hash (readonly)
Returns the requested SubscriptionFilter (camelCase keys).
133 134 135 |
# File 'lib/mcp_client/subscription.rb', line 133 def requested @requested end |
#server ⇒ MCPClient::ServerBase (readonly)
Returns the transport that owns the subscription.
137 138 139 |
# File 'lib/mcp_client/subscription.rb', line 137 def server @server end |
#state ⇒ Symbol (readonly)
Returns :pending, :active, :reconnecting or :closed.
139 140 141 |
# File 'lib/mcp_client/subscription.rb', line 139 def state @state end |
Class Method Details
.deep_frozen_copy(value) ⇒ Object
A detached, deeply frozen copy of a parsed JSON value.
690 691 692 693 694 695 696 697 |
# File 'lib/mcp_client/subscription.rb', line 690 def self.deep_frozen_copy(value) case value when Hash then value.to_h { |key, item| [deep_frozen_copy(key), deep_frozen_copy(item)] }.freeze when Array then value.map { |item| deep_frozen_copy(item) }.freeze when String then value.dup.freeze else value end end |
.filter_aliases ⇒ Hash{String => String}
Returns every accepted snake_case spelling.
75 76 77 |
# File 'lib/mcp_client/subscription.rb', line 75 def filter_aliases FILTER_ALIASES.merge(extension_filter_fields[:aliases]).freeze end |
.filter_fields ⇒ Hash{String => Symbol}
Returns every accepted filter field — the published four and the registered extension fields — with its type.
70 71 72 |
# File 'lib/mcp_client/subscription.rb', line 70 def filter_fields FILTER_FIELDS.merge(extension_filter_fields[:fields]).freeze end |
.normalize_filter(filter) ⇒ Hash
Normalize and validate a SubscriptionFilter given with String or Symbol, camelCase or snake_case keys.
The result is detached from the caller and frozen. The filter is not
serialized once and forgotten: Streamable HTTP builds the listen request
on the stream's own thread, after listen has returned, and every
reconnect builds it again — so an array the caller kept a reference to
would let a later << or a mutated String change the request that goes
out, or change what a re-opened stream asks for.
161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 |
# File 'lib/mcp_client/subscription.rb', line 161 def self.normalize_filter(filter) raise ArgumentError, 'notifications must be a Hash (SubscriptionFilter)' unless filter.is_a?(Hash) fields = filter_fields aliases = filter_aliases filter.to_h do |key, value| name = key.to_s name = aliases.fetch(name, name) type = fields[name] raise ArgumentError, "Unknown subscription filter field #{key.inspect}" unless type case type when :boolean raise ArgumentError, "#{name} must be true or false" unless [true, false].include?(value) when :string_array raise ArgumentError, "#{name} must be an array of strings" unless value.is_a?(Array) && value.all?(String) value = value.map { |item| item.dup.freeze }.freeze end [name, value] end.freeze end |
.register_filter_field(name, type, alias_name: nil) ⇒ void
This method returns an undefined value.
Register a filter field an extension defines beyond the published SubscriptionFilter, so normalize_filter accepts it (and its snake_case spelling) with the value type the extension gives it. A published field cannot be redefined; registering the same extension field twice with the same type is a no-op.
51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 |
# File 'lib/mcp_client/subscription.rb', line 51 def register_filter_field(name, type, alias_name: nil) name = name.to_s raise ArgumentError, "#{name} is a published SubscriptionFilter field" if FILTER_FIELDS.key?(name) raise ArgumentError, 'a filter field is boolean or string_array' unless FILTER_VALUE_TYPES.include?(type) extension_filter_fields_mutex.synchronize do fields = @extension_filter_fields || { fields: {}, aliases: {} } registered = fields[:fields][name] raise ArgumentError, "#{name} is already registered as #{registered}" if registered && registered != type fields[:fields][name] = type fields[:aliases][alias_name.to_s] = name if alias_name @extension_filter_fields = fields end nil end |
Instance Method Details
#acknowledge(filter) ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Record what the server agreed to honour.
The filter is copied and frozen through and through, arrays and strings
included. The hash it arrives in is the peer's, parsed from the
acknowledgment notification, and that same hash is handed to the host's
on_notification callback and to this subscription's own listeners — so
host code that edits it in place would otherwise be rewriting this
subscription's record of what the server granted. Adding a URI the
acknowledgment left out is enough to make a waiting subscribe_resource
report a watch that does not exist.
673 674 675 676 677 678 679 680 681 682 683 684 685 |
# File 'lib/mcp_client/subscription.rb', line 673 def acknowledge(filter) detached = Subscription.deep_frozen_copy(filter.is_a?(Hash) ? filter : {}) @mutex.synchronize do return if @state == :closed @acknowledged = detached @state = :active @answered = true # A stream the server has granted is no longer one being handed over. @reestablishing = false @settled.broadcast end end |
#active? ⇒ Boolean
Returns whether the server acknowledged it and it is still open.
262 263 264 |
# File 'lib/mcp_client/subscription.rb', line 262 def active? @mutex.synchronize { @state == :active } end |
#assign_id(id) ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Put the subscription on a listen id, without the registration and the write #with_open_id holds its lock across. No transport opens or re-opens a stream this way — #with_open_id is the step every one of them takes — so a guarantee about opening or reconnecting is one this method cannot stand in for.
416 417 418 419 420 421 422 423 424 425 426 |
# File 'lib/mcp_client/subscription.rb', line 416 def assign_id(id) @mutex.synchronize do @id = id @acknowledged = nil # A request the server has not seen is a question it has not answered. @answered = false # A subscription the host closed stays closed, whatever a racing # reconnect does. @state = :pending unless @state == :closed end end |
#await_live_resource_watch(uri, timeout) ⇒ Symbol
Block until this subscription is a running watch of the URI, and report whether it is.
This is the other question, and the one a subscribe_resource looking
for a stream to reuse asks: not "what did the server say" but "is the
server watching this, now". Only a running stream answers it. An
acknowledgment left on record by a stream that has dropped is not a
grant — no server-side subscription exists between listen attempts, the
request that replaces it is a new one the server may reject or
acknowledge more narrowly, and reading the old record as the current
grant reported a watch for the whole of an HTTP backoff or a stdio
handshake. So a stream between attempts is waited for instead.
373 374 375 |
# File 'lib/mcp_client/subscription.rb', line 373 def await_live_resource_watch(uri, timeout) await_watch(uri, timeout) { @state == :active } end |
#await_resource_watch(uri, timeout) ⇒ Symbol
Block until the server has answered the listen request this subscription is waiting on, and report what it said about this URI.
This is the question the subscriber that opened the stream asks: it is waiting for the answer to its own listen request, and it gets it. A connection that merely drops does not unask that question and does not unanswer it — the server did grant the filter, and making the caller wait out its whole acknowledgment timeout for an answer it had already been given would be the transport's problem told as the subscriber's. A replacement request that has actually gone out is another matter: the server holds no subscription state across one and has to grant the filter again, so the answer to it is waited for rather than assumed — which is why #with_open_id and #assign_id unanswer it explicitly, instead of that turning on whether the reconnect has got as far as taking an id.
Waiting is also the right answer for a request in flight with nothing granted yet, which used to read as success: a subscription with no acknowledgment has no unacknowledged URIs either.
353 354 355 |
# File 'lib/mcp_client/subscription.rb', line 353 def await_resource_watch(uri, timeout) await_watch(uri, timeout) { @answered } end |
#close ⇒ MCPClient::Subscription?
Cancel the subscription: the transport closes the stream (HTTP) or sends notifications/cancelled (stdio).
399 400 401 402 403 404 |
# File 'lib/mcp_client/subscription.rb', line 399 def close return nil if closed? @server.cancel_subscription(self) self end |
#closed? ⇒ Boolean
267 268 269 |
# File 'lib/mcp_client/subscription.rb', line 267 def closed? @mutex.synchronize { @state == :closed } end |
#closed_by_client? ⇒ Boolean
Returns whether #close ended it.
284 285 286 |
# File 'lib/mcp_client/subscription.rb', line 284 def closed_by_client? @mutex.synchronize { @closed_by_client } end |
#closed_gracefully? ⇒ Boolean
Returns whether the server ended it with a response to the listen request.
279 280 281 |
# File 'lib/mcp_client/subscription.rb', line 279 def closed_gracefully? @mutex.synchronize { @closed_gracefully } end |
#deliver(method, params) ⇒ Object
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
700 701 702 703 704 705 706 707 708 709 710 711 |
# File 'lib/mcp_client/subscription.rb', line 700 def deliver(method, params) @mutex.synchronize do return if @state == :closed || @listeners.empty? # Queued while still open, so a closure that follows cannot swallow a # notification that had already arrived, and the stop that ends the # dispatcher can never overtake this one. Enqueuing never waits for # the listeners: the transport reader must stay free to deliver the # responses a listener's own requests are waiting for. (@dispatcher ||= NotificationDispatcher.new(self)).deliver(@listeners.dup, method, params) end end |
#discard_outstanding_listens ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Forget the recorded listen ids without cancelling them: the session they were written to is gone, so nothing is outstanding and none of them must be cancelled on the session that replaces it.
553 554 555 556 557 558 |
# File 'lib/mcp_client/subscription.rb', line 553 def discard_outstanding_listens @mutex.synchronize do @outstanding_listens = [] @unwritten_listens = [] end end |
#dropped_notifications ⇒ Integer
Notifications dropped because the listeners could not keep up with the server (see MAX_PENDING_NOTIFICATIONS and MAX_PENDING_NOTIFICATION_BYTES).
248 249 250 251 |
# File 'lib/mcp_client/subscription.rb', line 248 def dropped_notifications dispatcher = @mutex.synchronize { @dispatcher } dispatcher ? dispatcher.dropped : 0 end |
#expire_unanswered(id, error) { ... } ⇒ Symbol?
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
End the request a deadline was set on, if the subscription is still on that request and the server has still not answered it — in one step.
The watchdog that waits out the deadline cannot decide this for itself: by the time its wait returns, a restart may have replaced the request with a newer one (a new id, with a deadline of its own, which the older request's timer must not spend), and an acknowledgment may have landed in the instant between the wait and the verdict, which is the server's answer and is kept. Anyone waiting for the subscription to settle is woken only once the block — the cancellation the transport sends for the expired request — has run: a host that sees the handle settle on a timeout sees a server that has already been told, rather than one the watchdog is still writing to. The block runs outside the lock, since telling the server takes the ids recorded on this subscription.
638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 |
# File 'lib/mcp_client/subscription.rb', line 638 def expire_unanswered(id, error) @mutex.synchronize do return nil if @id != id || @answered || SETTLED_STATES.include?(@state) # A deadline bounds the request that is out, not the subscription: a # transport waiting to send the next listen (an HTTP reconnect inside # its backoff, a stdio restart still spawning) has nothing in flight # for this deadline to expire, and ending the handle here would also # mark it closed by the client — unreconnectable, so the re-send that # basic/patterns/subscriptions requires after a reconnect never # happens. The next attempt arms a deadline of its own with its id. return nil if @state == :reconnecting close_locked(by_client: true, error: error, announce: false) end begin yield if block_given? ensure @mutex.synchronize { @settled.broadcast } end :expired end |
#fail_attempt(id, error) ⇒ Symbol
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Undo the listen attempt that failed — or find that it is no longer this attempt's to undo — in one step.
#open_as? used to answer the first half of that question on its own, and a restart could re-open the subscription under a newer id and have it acknowledged between the answer and the transition it guarded: the older attempt then finished the very stream the fresh process was serving, with nothing left to cancel it. Asking and acting under one hold of the lock is what makes the answer good for the transition it decides.
The three answers are the three things a failed attempt can be: superseded by a newer attempt, which owns the subscription now; a hand-over to a new session that could not be written (#reestablishing?), which goes back to waiting for the next one; or the caller's own request, which ends with the error. A hand-over is deferred only when it is the process that could not be written to: that process is on its way out, and the next one drains the queue. An attempt that failed before it took an id at all (the request could not be built) says nothing about the process, which stays up and healthy — nothing would ever drain the queue on it, and the previous acknowledgment would keep the watchdog from expiring it: the subscription stayed :reconnecting for ever with the host never told. So that failure is the subscription's own, and it ends with it.
602 603 604 605 606 607 608 609 610 611 612 613 614 615 |
# File 'lib/mcp_client/subscription.rb', line 602 def fail_attempt(id, error) @mutex.synchronize do return :superseded if id && @id != id return :failed if @state == :closed if @reestablishing && id @state = :reconnecting return :deferred end close_locked(error: error) :failed end end |
#finish(gracefully: false, by_client: false, error: nil, reason: nil) ⇒ Object
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
740 741 742 |
# File 'lib/mcp_client/subscription.rb', line 740 def finish(gracefully: false, by_client: false, error: nil, reason: nil) @mutex.synchronize { close_locked(gracefully: gracefully, by_client: by_client, error: error, reason: reason) } end |
#inspect ⇒ Object
750 751 752 |
# File 'lib/mcp_client/subscription.rb', line 750 def inspect "#<MCPClient::Subscription id=#{@id.inspect} state=#{@state} requested=#{@requested.keys.join(',')}>" end |
#mark_listen_written(id) ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
The write of a listen request has finished — sent, or raised having sent who knows how much. Either way the id may now be cancelled.
An id the session that carried it has since discarded (#discard_outstanding_listens) stays discarded: nothing written to a process that is gone is outstanding, and a late write that lands on its closed pipe must not put the id back.
514 515 516 |
# File 'lib/mcp_client/subscription.rb', line 514 def mark_listen_written(id) @mutex.synchronize { @unwritten_listens.delete(id) } end |
#mark_reconnecting ⇒ Object
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
714 715 716 717 718 719 720 721 |
# File 'lib/mcp_client/subscription.rb', line 714 def mark_reconnecting @mutex.synchronize do next if @state == :closed @state = :reconnecting @reestablishing = true end end |
#on_notification {|method, params| ... } ⇒ self
Add a listener for notifications delivered on this subscription.
256 257 258 259 |
# File 'lib/mcp_client/subscription.rb', line 256 def on_notification(&block) @mutex.synchronize { @listeners << block } self end |
#open_as?(id) ⇒ Boolean
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Whether this subscription is still the stream a given listen id opened. A transport that fails an attempt asks before undoing it: a restart racing a blocked write may already have re-opened the subscription under a newer id, and that stream is not the older attempt's to tear down.
567 568 569 |
# File 'lib/mcp_client/subscription.rb', line 567 def open_as?(id) @mutex.synchronize { @id == id } end |
#open_generation ⇒ Integer?
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
The transport generation the listen this subscription is on went out under, or nil on a transport that does not number its processes.
462 463 464 |
# File 'lib/mcp_client/subscription.rb', line 462 def open_generation @mutex.synchronize { @open_generation } end |
#pending_notification_bytes ⇒ Integer
Bytes retained by the notifications queued for this subscription's listeners, measured as the JSON the peer sent for them — their method names as well as their params (see MAX_PENDING_NOTIFICATION_BYTES).
239 240 241 242 |
# File 'lib/mcp_client/subscription.rb', line 239 def pending_notification_bytes dispatcher = @mutex.synchronize { @dispatcher } dispatcher ? dispatcher.pending_bytes : 0 end |
#pending_notifications ⇒ Integer
Returns notifications queued for this subscription's listeners.
230 231 232 233 |
# File 'lib/mcp_client/subscription.rb', line 230 def pending_notifications dispatcher = @mutex.synchronize { @dispatcher } dispatcher ? dispatcher.pending : 0 end |
#reconnectable? ⇒ Boolean
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Returns whether the subscription should be re-established after a reconnect.
746 747 748 |
# File 'lib/mcp_client/subscription.rb', line 746 def reconnectable? @mutex.synchronize { @state != :closed && !@closed_by_client } end |
#reconnecting? ⇒ Boolean
Returns whether it is waiting for a transport to re-establish it — the stream dropped, or the stdio process it was on exited, and the transport that noticed has queued it for the next session.
274 275 276 |
# File 'lib/mcp_client/subscription.rb', line 274 def reconnecting? @mutex.synchronize { @state == :reconnecting } end |
#record_outstanding_listen(id, io = nil) ⇒ void
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
This method returns an undefined value.
Record a listen request the transport has written for this subscription on the session it is on.
A cancellation has to name a request the server may be serving, and
that is not always the id the subscription happens to be on: a second
listen written for it on one session (a hand-over that queued it twice,
say) leaves the server holding the first stream, and
notifications/cancelled for the newest id alone would never close it —
the stream stays open until the server's own timeout, with the client
unable to name it again.
An attempt whose write raised is recorded too: the client cannot know how much of it the peer saw, and cancelling a request the server never received is ignored, while failing to cancel one it did receive is not. Recorded before the write for that reason, and marked written by #mark_listen_written whichever way the write ends.
The pipe it is written to is recorded with it, because forgetting the
ids of a process that is gone (#discard_outstanding_listens) cannot
reach an attempt that has not recorded its id yet. One paused here while
its process was torn down recorded afterwards, with nothing left to
forget it, and the close that followed named it on the process that
replaced it — a request that one had never been sent, while
"the cancelled request MUST have been previously issued"
(basic/patterns/cancellation). Recording the pipe makes the id
cancellable on that pipe alone, whenever it is recorded.
497 498 499 500 501 502 |
# File 'lib/mcp_client/subscription.rb', line 497 def record_outstanding_listen(id, io = nil) @mutex.synchronize do @outstanding_listens << [id, io] unless @outstanding_listens.any? { |(known, _)| known == id } @unwritten_listens << id unless @unwritten_listens.include?(id) end end |
#reestablishing? ⇒ Boolean
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Whether a transport is handing this subscription to a new session: it was #mark_reconnectinged and no server has acknowledged it since.
Unlike #reconnecting? this survives the :pending that taking the new listen id moves it to, which is the whole point: a transport whose re-send fails on the write has to tell a stream it is handing over — which MUST be re-sent, and belongs to the next session — from one it is opening for a caller, which is the caller's to hear about. The state alone cannot: by the time the write raises, the re-send has already moved it off :reconnecting.
735 736 737 |
# File 'lib/mcp_client/subscription.rb', line 735 def reestablishing? @mutex.synchronize { @reestablishing && @state != :closed } end |
#take_outstanding_listens(io = nil) ⇒ Array
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
The listen ids the server may still be serving for this subscription, leaving none behind: the caller is cancelling them.
An id whose write has not finished is not among them, however impatient
the caller: "the cancelled request MUST have been previously issued"
(basic/patterns/cancellation), and cancelling an id the pipe has not
carried yet put cancelled(n) on the wire ahead of listen(n). The
transport that is writing it cancels it itself once the write is done
and it finds the subscription closed — the one moment at which the
cancellation can name a request the server has actually been sent.
Nor is an id written to a different pipe among them: the process on this one was never sent that request (see #record_outstanding_listen). Those are left recorded rather than dropped, since the caller that pins a pipe is not always the one that will cancel on the pipe they went to.
538 539 540 541 542 543 544 545 546 |
# File 'lib/mcp_client/subscription.rb', line 538 def take_outstanding_listens(io = nil) @mutex.synchronize do taken, kept = @outstanding_listens.partition do |(id, recorded_io)| !@unwritten_listens.include?(id) && cancellable_on?(recorded_io, io) end @outstanding_listens = kept taken.map(&:first) end end |
#unacknowledged_resource_uris ⇒ Array<String>
Requested resource URIs the server did not agree to watch.
306 307 308 309 310 311 312 313 |
# File 'lib/mcp_client/subscription.rb', line 306 def unacknowledged_resource_uris ack = @mutex.synchronize { @acknowledged } return [] unless ack wanted = Array(@requested['resourceSubscriptions']) granted = ack['resourceSubscriptions'].is_a?(Array) ? ack['resourceSubscriptions'] : [] wanted - granted end |
#unsupported ⇒ Array<String>
Requested notification types the server did not agree to honour.
Support is read from the value the server acknowledged, not from the
mere presence of the field: an acknowledgment that names
resourceSubscriptions with none of the URIs it was sent has accepted
no resource subscription at all, and a flag acknowledged as false will
not be honoured either. A list the server granted in part counts as
supported; see #unacknowledged_resource_uris for the URIs it left out.
297 298 299 300 301 302 |
# File 'lib/mcp_client/subscription.rb', line 297 def unsupported ack = @mutex.synchronize { @acknowledged } return [] unless ack @requested.keys.reject { |field| granted?(@requested[field], ack[field]) } end |
#wait_until_settled(timeout) ⇒ Symbol?
Block until the server has settled the subscription: acknowledged it, or ended it with a response, an error or a cancellation.
381 382 383 384 385 386 387 388 389 390 391 392 393 394 |
# File 'lib/mcp_client/subscription.rb', line 381 def wait_until_settled(timeout) deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout @mutex.synchronize do loop do answer = settled_state return answer if answer remaining = deadline - Process.clock_gettime(Process::CLOCK_MONOTONIC) return nil if remaining <= 0 @settled.wait(@mutex, remaining) end end end |
#watching_resource?(uri) ⇒ Boolean
Whether this stream is, right now, an active acknowledged watch of a resource: the server granted that URI and the stream it granted it on is the one still running.
Being open is not enough, which is what a subscribe_resource looking
for a stream to reuse has to know: a stream between listen attempts is
serving nothing, and the request that replaces it is a new one the
server holds no state for — it may be rejected, or acknowledged more
narrowly.
326 327 328 |
# File 'lib/mcp_client/subscription.rb', line 326 def watching_resource?(uri) @mutex.synchronize { @state == :active && acknowledges_resource?(uri) } end |
#with_open_id(id, generation = nil) { ... } ⇒ Boolean
This method is part of a private API. You should avoid using this method if possible, as it may be removed or be changed in the future.
Take a fresh listen id and register/send the request under this subscription's own lock, so a concurrent #close either wins outright (nothing is sent) or waits and then cancels the id that was sent. A closed subscription is never re-opened.
436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 |
# File 'lib/mcp_client/subscription.rb', line 436 def with_open_id(id, generation = nil) @mutex.synchronize do return false if @state == :closed @id = id @acknowledged = nil # The request that went out is a new one: whatever the server said # about the last, it has not answered this. @answered = false @state = :pending # Stamped in the same step as the registration, so every registered # subscription names the process its listen went out on: a teardown # parks the ones belonging to the process it claimed and leaves the # replacement's alone (see {MCPClient::ServerStdio#park_open_subscriptions}). # A transport with no such notion (the HTTP ones, whose streams are # per-connection) passes none and is never asked. @open_generation = generation yield true end end |