Class: Parse::Agent::MCPRackApp::ListeningStreamBody Private
- Inherits:
-
Object
- Object
- Parse::Agent::MCPRackApp::ListeningStreamBody
- 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
-
#close ⇒ Object
private
Terminate the stream: stop heartbeats, detach the listener, and tear down the session's LiveQuery subscriptions.
-
#each {|String| ... } ⇒ Object
private
Rack body interface — called once by the Rack server.
-
#initialize(manager, session_id, heartbeat_interval, logger) ⇒ ListeningStreamBody
constructor
private
A new instance of ListeningStreamBody.
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.
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.}" @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.
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 |