Class: Parse::Agent::MCPRackApp::ListeningStreamBody Private

Inherits:
Object
  • Object
show all
Defined in:
lib/parse/agent/mcp_rack_app.rb

Overview

This class is part of a private API. You should avoid using this class if possible, as it may be removed or be changed in the future.

Rack body for the long-lived GET listening stream that carries notifications/resources/updated to a subscribing client.

On #each it registers a delivery callback with the Parse::Agent::MCPSubscriptions::Manager keyed by the session id, then blocks reading from an internal queue and yields SSE-formatted notification events as they are published by the LiveQuery bridge. A periodic SSE comment heartbeat keeps the connection warm and surfaces a dead socket as a write error so the Rack server invokes #close.

#close detaches the listener and tears down every LiveQuery subscription bound to the session — so a dropped stream leaves no LiveQuery sockets behind. Re-opening the stream requires the client to re-issue its resources/subscribe calls (subscriptions do not survive a listening-stream disconnect in this single-process implementation).

The publish callback runs on a LiveQuery dispatcher / debounce thread and only pushes to the thread-safe queue; all yields happen on the Rack I/O thread driving #each, mirroring SSEBody's threading model.

Constant Summary collapse

DONE =

This constant is part of a private API. You should avoid using this constant if possible, as it may be removed or be changed in the future.

:__listening_done__

Instance Method Summary collapse

Constructor Details

#initialize(manager, session_id, heartbeat_interval, logger) ⇒ ListeningStreamBody

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 a new instance of ListeningStreamBody.

Parameters:



1788
1789
1790
1791
1792
1793
1794
1795
1796
1797
1798
# File 'lib/parse/agent/mcp_rack_app.rb', line 1788

def initialize(manager, session_id, heartbeat_interval, logger)
  @manager = manager
  @session_id = session_id
  @heartbeat_interval = heartbeat_interval
  @logger = logger
  @queue = Queue.new
  @heartbeat = nil
  @closed = false
  @counted = false
  @close_mutex = Mutex.new
end

Instance Method Details

#close ⇒ 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.

Terminate the stream: stop heartbeats, detach the listener, and tear down the session's LiveQuery subscriptions. Idempotent.



1828
1829
1830
1831
1832
1833
1834
1835
1836
1837
1838
1839
1840
1841
1842
1843
1844
1845
# File 'lib/parse/agent/mcp_rack_app.rb', line 1828

def close
  @close_mutex.synchronize do
    return if @closed
    @closed = true
  end
  # Balance the #each increment exactly once (close is idempotent via
  # @closed, and only #each sets @counted).
  MCPRackApp.adjust_listening_stream_count(-1) if @counted
  @heartbeat&.kill
  @heartbeat = nil
  begin
    @manager.detach_listener(@session_id)
  rescue StandardError => e
    line = "[Parse::Agent::MCPRackApp::ListeningStreamBody] detach error: #{e.class}: #{e.message}"
    @logger ? @logger.warn(line) : warn(line)
  end
  @queue << DONE rescue nil
end

#each {|String| ... } ⇒ 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.

Rack body interface — called once by the Rack server.

Yields:

  • (String) —

    SSE-formatted event / comment strings.



1802
1803
1804
1805
1806
1807
1808
1809
1810
1811
1812
1813
1814
1815
1816
1817
1818
1819
1820
1821
1822
1823
1824
# File 'lib/parse/agent/mcp_rack_app.rb', line 1802

def each
  queue = @queue
  # Count this stream against the concurrent-listening-stream soft cap.
  # Incrementing here (in #each, not the constructor) means a body the
  # Rack server never iterates — or a client that disconnects before
  # iteration — never inflates the counter; the matching decrement is in
  # #close, which #each's `ensure` always runs.
  MCPRackApp.adjust_listening_stream_count(1)
  @counted = true
  @manager.attach_listener(@session_id) do |notification|
    queue << format_event(notification)
  end
  # Initial comment flushes response headers and confirms the stream.
  yield ": connected\n\n"
  start_heartbeat
  loop do
    msg = @queue.pop
    break if msg == DONE
    yield msg
  end
ensure
  close
end