Class: Nanobot::Bus::MessageBus
- Inherits:
-
Object
- Object
- Nanobot::Bus::MessageBus
- Defined in:
- lib/nanobot/bus/message_bus.rb
Overview
MessageBus provides a thread-safe queue-based message routing system that decouples channels from the agent core
Instance Attribute Summary collapse
-
#logger ⇒ Object
readonly
Returns the value of attribute logger.
Instance Method Summary collapse
-
#consume_inbound(timeout: nil) ⇒ InboundMessage?
Consume an inbound message (agent reads from channels).
-
#initialize(logger: nil) ⇒ MessageBus
constructor
A new instance of MessageBus.
-
#publish_inbound(message) ⇒ Object
Publish an inbound message (from channel to agent).
-
#publish_outbound(message) ⇒ Object
Publish an outbound message (from agent to channels).
-
#queue_sizes ⇒ Hash{Symbol => Integer}
Get queue sizes for monitoring.
-
#running? ⇒ Boolean
Check if the message bus is running.
-
#start_dispatch ⇒ Object
Start the outbound dispatcher thread.
-
#stop ⇒ Object
Stop the message bus.
-
#subscribe_outbound(channel, &callback) ⇒ Object
Subscribe to outbound messages for a specific channel.
Constructor Details
#initialize(logger: nil) ⇒ MessageBus
Returns a new instance of MessageBus.
14 15 16 17 18 19 20 21 22 |
# File 'lib/nanobot/bus/message_bus.rb', line 14 def initialize(logger: nil) @inbound_queue = Queue.new @outbound_queue = Queue.new @outbound_subscribers = Hash.new { |h, k| h[k] = [] } @running = false @logger = logger || Logger.new(IO::NULL) @dispatch_thread = nil @mutex = Mutex.new end |
Instance Attribute Details
#logger ⇒ Object (readonly)
Returns the value of attribute logger.
11 12 13 |
# File 'lib/nanobot/bus/message_bus.rb', line 11 def logger @logger end |
Instance Method Details
#consume_inbound(timeout: nil) ⇒ InboundMessage?
Consume an inbound message (agent reads from channels)
36 37 38 39 40 41 42 43 44 45 46 47 48 |
# File 'lib/nanobot/bus/message_bus.rb', line 36 def consume_inbound(timeout: nil) if timeout begin Timeout.timeout(timeout) do @inbound_queue.pop end rescue Timeout::Error nil end else @inbound_queue.pop end end |
#publish_inbound(message) ⇒ Object
Publish an inbound message (from channel to agent)
26 27 28 29 30 31 |
# File 'lib/nanobot/bus/message_bus.rb', line 26 def publish_inbound() raise ArgumentError, 'Message must be an InboundMessage' unless .is_a?(InboundMessage) @inbound_queue.push() @logger.debug "Published inbound message from #{message.channel}:#{message.chat_id}" end |
#publish_outbound(message) ⇒ Object
Publish an outbound message (from agent to channels)
52 53 54 55 56 57 |
# File 'lib/nanobot/bus/message_bus.rb', line 52 def publish_outbound() raise ArgumentError, 'Message must be an OutboundMessage' unless .is_a?(OutboundMessage) @outbound_queue.push() @logger.debug "Published outbound message to #{message.channel}:#{message.chat_id}" end |
#queue_sizes ⇒ Hash{Symbol => Integer}
Get queue sizes for monitoring
100 101 102 103 104 105 |
# File 'lib/nanobot/bus/message_bus.rb', line 100 def queue_sizes { inbound: @inbound_queue.size, outbound: @outbound_queue.size } end |
#running? ⇒ Boolean
Check if the message bus is running
94 95 96 |
# File 'lib/nanobot/bus/message_bus.rb', line 94 def running? @running end |
#start_dispatch ⇒ Object
Start the outbound dispatcher thread
72 73 74 75 76 77 78 79 80 81 82 |
# File 'lib/nanobot/bus/message_bus.rb', line 72 def start_dispatch @mutex.synchronize do return if @running @running = true end @dispatch_thread = Thread.new do dispatch_loop end @logger.info 'Message bus dispatch started' end |
#stop ⇒ Object
Stop the message bus
85 86 87 88 89 90 |
# File 'lib/nanobot/bus/message_bus.rb', line 85 def stop @running = false @outbound_queue.push(nil) # Unblock the dispatch thread @dispatch_thread&.join(5) # Wait up to 5 seconds @logger.info 'Message bus stopped' end |
#subscribe_outbound(channel, &callback) ⇒ Object
Subscribe to outbound messages for a specific channel
62 63 64 65 66 67 68 69 |
# File 'lib/nanobot/bus/message_bus.rb', line 62 def subscribe_outbound(channel, &callback) raise ArgumentError, 'Block required' unless block_given? @mutex.synchronize do @outbound_subscribers[channel] << callback end @logger.debug "Subscribed to outbound messages for channel: #{channel}" end |