Class: MCPClient::Subscription::NotificationDispatcher

Inherits:
Object
  • Object
show all
Defined in:
lib/mcp_client/subscription/notification_dispatcher.rb

Overview

One subscription's notification dispatcher: the queue the transport reader fills, the thread the listeners run on, and the policy that keeps the queue bounded.

Listeners never run on the transport's reader. On stdio that is the single stdout reader, so a listener reacting to an update with a request of its own (re-reading the resource that changed, say) would otherwise wait for a response only the thread it is blocking could deliver — and every other message would wait with it. Enqueuing therefore never blocks and never waits for a listener.

The queue is filled by the peer and drained by the host, so it needs a ceiling, and overflow has to discard something. There are two ceilings, because a queue bounded by count alone is not bounded in memory: the method name and params of every queued notification are retained until its listener has run, and the peer chooses how big they are. So there is a count ceiling (MAX_PENDING_NOTIFICATIONS) and a byte budget (MAX_PENDING_NOTIFICATION_BYTES).

Two rules keep the queue honest, and they are the same rule seen from each end:

  1. Every queued notification is charged exactly what it retains — its method name as well as its params, since the entry keeps both. #pending_bytes is the sum of the queue, always, because the queue and its totals are only ever changed together (see #push and #discard).
  2. Every eviction removes an entry whose removal relieves the pressure that caused it. Each pressure has its own candidates: the byte budget admits only the entries charged against it, the count ceiling admits every entry (each is one of the count), and the oversized slot admits only its occupant. So overflow always makes progress and a signal is never spent on pressure that discarding it cannot relieve.

Earlier revisions decided these two by rules that disagreed — a payload exempt from the charge but not from eviction — and the queue could throw away the only notice of a resource and still be over budget.

Which of the candidates goes is chosen by identity — the notification's method together with the resource URI or task id it names — never by arrival order alone: one stream can carry a mixed filter, and dropping the oldest entry would throw away the only queued update for a quiet resource to keep newer ones for a busy one, with nothing left to tell the listener to re-read the quiet one. Since every MCP notification is a "look again" signal about state the host re-reads for itself, a second notice of the same thing is redundant and the first notice of a thing is not. So overflow gives up, in order of preference:

  1. the oldest candidate of the same identity as the arriving notification — the listener still gets the newest word on it;
  2. otherwise the oldest candidate of whichever identity has the most queued, so nothing loses its only notice while something else has a spare;
  3. only when every candidate names a different thing, the oldest — the queue is then full of distinct signals and one must go.

A notification whose payload is larger than the whole byte budget is not charged against it. It is held in a slot of its own instead, and there is only ever one such slot: a second oversized payload takes it from the first, which is the only thing that ever displaces one. So such a payload is neither lost for being large nor able to displace what the budget holds — nothing else is charged to its slot, and discarding it would free nothing the budget is short of — while what the queue retains stays within the budget plus one peer-sized payload.

Defined Under Namespace

Classes: Queued

Instance Method Summary collapse

Constructor Details

#initialize(owner) ⇒ NotificationDispatcher

Returns a new instance of NotificationDispatcher.

Parameters:



79
80
81
82
83
84
85
86
87
88
89
90
# File 'lib/mcp_client/subscription/notification_dispatcher.rb', line 79

def initialize(owner)
  @owner = owner
  @mutex = Mutex.new
  @ready = ConditionVariable.new
  @buffer = []
  @bytes = 0
  @oversized_bytes = 0
  @dropped = 0
  @warned_about_drops = false
  @stopped = false
  start_thread
end

Instance Method Details

#deliver(listeners, method, params) ⇒ void

This method returns an undefined value.

Queue one notification for the listeners, making room for it first.

Parameters:

  • listeners (Array<Proc>) —

    the listeners to run

  • method (String) —

    notification method

  • params (Hash, nil) —

    notification params



114
115
116
117
118
119
120
121
122
123
124
125
126
# File 'lib/mcp_client/subscription/notification_dispatcher.rb', line 114

def deliver(listeners, method, params)
  # Measured before the lock is taken: sizing a large payload must not
  # hold up the reader thread that is delivering the next one.
  bytes = payload_bytesize(method, params)
  entry = Queued.new(identity(method, params), listeners, method, params, bytes, bytes > byte_capacity)
  @mutex.synchronize do
    return if @stopped

    make_room(entry)
    push(entry)
    @ready.signal
  end
end

#dropped ⇒ Integer

Returns notifications discarded because the listeners could not keep up with the peer.

Returns:

  • (Integer) —

    notifications discarded because the listeners could not keep up with the peer



105
106
107
# File 'lib/mcp_client/subscription/notification_dispatcher.rb', line 105

def dropped
  @mutex.synchronize { @dropped }
end

#pending ⇒ Integer

Returns notifications waiting for the listeners.

Returns:

  • (Integer) —

    notifications waiting for the listeners



93
94
95
# File 'lib/mcp_client/subscription/notification_dispatcher.rb', line 93

def pending
  @mutex.synchronize { @buffer.size }
end

#pending_bytes ⇒ Integer

Returns bytes retained by the notifications waiting for the listeners.

Returns:

  • (Integer) —

    bytes retained by the notifications waiting for the listeners



99
100
101
# File 'lib/mcp_client/subscription/notification_dispatcher.rb', line 99

def pending_bytes
  @mutex.synchronize { @bytes }
end

#stop ⇒ void

This method returns an undefined value.

End the dispatcher after everything already queued has been delivered.



130
131
132
133
134
135
# File 'lib/mcp_client/subscription/notification_dispatcher.rb', line 130

def stop
  @mutex.synchronize do
    @stopped = true
    @ready.broadcast
  end
end