Class: MCPClient::Subscription

Inherits:
Object
  • Object
show all
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

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(server:, requested:, ack_timeout: nil) {|method, params| ... } ⇒ Subscription

Returns a new instance of Subscription.

Parameters:

  • server (MCPClient::ServerBase) —

    owning transport

  • requested (Hash) —

    normalized filter

  • ack_timeout (Numeric, false, nil) (defaults to: nil) —

    the acknowledgment deadline the host asked for, kept for the requests re-issued later (see MCPClient::SubscriptionSupport#rearm_acknowledgment_deadline). Set here rather than assigned afterwards: a transport may re-open the stream before listen has returned the handle.

Yields:

  • (method, params) —

    optional listener for notifications on this 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.

Returns:

  • (Numeric, false, nil) —

    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.

Returns:

  • (Hash, nil) —

    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.

Returns:

  • (String, nil) —

    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.

Returns:



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).

Returns:

  • (Integer, String, nil) —

    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).

Returns:

  • (Hash) —

    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.

Returns:



137
138
139
# File 'lib/mcp_client/subscription.rb', line 137

def server
  @server
end

#state ⇒ Symbol (readonly)

Returns :pending, :active, :reconnecting or :closed.

Returns:

  • (Symbol) —

    :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.

Parameters:

  • value (Object)

Returns:

  • (Object) —

    frozen, sharing nothing mutable with the original



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.

Returns:

  • (Hash{String => String}) —

    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.

Returns:

  • (Hash{String => Symbol}) —

    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.

Parameters:

  • filter (Hash) —

    the notification filter

Returns:

  • (Hash) —

    camelCase String keys, frozen

Raises:

  • (ArgumentError) —

    on an unknown key or a mistyped value



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.

Parameters:

  • name (String) —

    the camelCase wire name of the field

  • type (Symbol) —

    :boolean or :string_array

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

    a snake_case spelling to accept for it

Raises:

  • (ArgumentError) —

    on a published field, an unknown type, or a name already registered with another type



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.

Parameters:

  • filter (Hash, nil) —

    the acknowledged SubscriptionFilter



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.

Returns:

  • (Boolean) —

    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.

Parameters:

  • id (Integer, String) —

    the listen request id



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.

Parameters:

  • uri (String) —

    the resource URI

  • timeout (Numeric) —

    seconds to wait for the stream to be granted

Returns:

  • (Symbol) —

    :watching, :not_watching (granted without the URI), :closed, or :timeout



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.

Parameters:

  • uri (String) —

    the resource URI

  • timeout (Numeric) —

    seconds to wait for an answer

Returns:

  • (Symbol) —

    :watching, :not_watching (answered without the URI), :closed, or :timeout



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).

Returns:



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

Returns:

  • (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.

Returns:

  • (Boolean) —

    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.

Returns:

  • (Boolean) —

    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).

Returns:

  • (Integer)


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.

Parameters:

Yields:

  • after the subscription has been ended, before the waiters wake

Returns:

  • (Symbol, nil) —

    :expired when the subscription was ended here; nil when the request was no longer this one's to expire



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.

Parameters:

  • id (Integer, String, nil) —

    the listen id the attempt sent under; nil for an attempt that failed before it took one (the request could not be built), which no newer attempt can have superseded

  • error (MCPClient::Errors::MCPError) —

    why it failed

Returns:

  • (Symbol) —

    :superseded, :deferred (it is :reconnecting again) or :failed (ended with the error — or already closed)



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.

Parameters:

  • id (Integer, String) —

    the listen request id



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.

Yields:

  • (method, params)

Returns:

  • (self)


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.

Parameters:

  • id (Integer, String) —

    a listen request id

Returns:

  • (Boolean)


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.

Returns:

  • (Integer, nil)


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).

Returns:

  • (Integer)


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.

Returns:

  • (Integer) —

    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.

Returns:

  • (Boolean) —

    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.

Returns:

  • (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



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.

Parameters:

  • id (Integer, String) —

    the listen request id

  • io (IO, nil) (defaults to: nil) —

    the pipe the request is being written to; nil leaves the id cancellable wherever the caller is cancelling



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.

Returns:

  • (Boolean)


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.

Parameters:

  • io (IO, nil) (defaults to: nil) —

    cancel only what was written to this pipe; nil takes every written id, and an id recorded against no pipe is taken by any caller

Returns:

  • (Array) —

    the recorded ids that have been written, oldest first



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.

Returns:

  • (Array<String>) —

    empty until acknowledged



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.

Returns:

  • (Array<String>) —

    empty until acknowledged



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.

Parameters:

  • timeout (Numeric) —

    seconds to wait

Returns:

  • (Symbol, nil) —

    :active or :closed, nil while it is still pending



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.

Parameters:

  • uri (String) —

    the resource URI

Returns:

  • (Boolean)


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.

Parameters:

  • id (Integer, String) —

    the new listen request id

Yields:

  • runs while the id cannot change underneath it

Returns:

  • (Boolean) —

    false when the host had already closed it



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