Class: SolidMCP::PubSub
- Inherits:
-
Object
- Object
- SolidMCP::PubSub
- Defined in:
- lib/solid_mcp/pub_sub.rb
Instance Method Summary collapse
-
#broadcast(session_id, event_type, data) ⇒ Object
Broadcast a message to a session (uses MessageWriter for batching).
-
#initialize(options = {}) ⇒ PubSub
constructor
A new instance of PubSub.
-
#shutdown ⇒ Object
Shutdown all listeners.
-
#subscribe(session_id, &block) ⇒ Object
Subscribe to messages for a specific session.
-
#unsubscribe(session_id) ⇒ Object
Unsubscribe from a session.
Constructor Details
#initialize(options = {}) ⇒ PubSub
Returns a new instance of PubSub.
8 9 10 11 12 |
# File 'lib/solid_mcp/pub_sub.rb', line 8 def initialize( = {}) @options = @subscriptions = Concurrent::Map.new @listeners = Concurrent::Map.new end |
Instance Method Details
#broadcast(session_id, event_type, data) ⇒ Object
Broadcast a message to a session (uses MessageWriter for batching)
31 32 33 |
# File 'lib/solid_mcp/pub_sub.rb', line 31 def broadcast(session_id, event_type, data) MessageWriter.instance.enqueue(session_id, event_type, data) end |
#shutdown ⇒ Object
Shutdown all listeners
36 37 38 39 40 41 42 |
# File 'lib/solid_mcp/pub_sub.rb', line 36 def shutdown @listeners.each do |_, listener| listener.stop end @listeners.clear MessageWriter.instance.shutdown end |
#subscribe(session_id, &block) ⇒ Object
Subscribe to messages for a specific session
15 16 17 18 19 20 21 22 |
# File 'lib/solid_mcp/pub_sub.rb', line 15 def subscribe(session_id, &block) # Atomically get or create callbacks array callbacks = @subscriptions.compute_if_absent(session_id) { Concurrent::Array.new } callbacks << block # Start a listener for this session if not already running ensure_listener_for(session_id) end |
#unsubscribe(session_id) ⇒ Object
Unsubscribe from a session
25 26 27 28 |
# File 'lib/solid_mcp/pub_sub.rb', line 25 def unsubscribe(session_id) @subscriptions.delete(session_id) stop_listener_for(session_id) end |