Class: Nanobot::Bus::MessageBus

Inherits:
Object
  • Object
show all
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

Instance Method Summary collapse

Constructor Details

#initialize(logger: nil) ⇒ MessageBus

Returns a new instance of MessageBus.

Parameters:

  • logger (Logger, nil) (defaults to: nil) —

    optional logger instance



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)

Parameters:

  • timeout (Numeric, nil) (defaults to: nil) —

    timeout in seconds, nil for blocking

Returns:



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)

Parameters:

Raises:

  • (ArgumentError)


26
27
28
29
30
31
# File 'lib/nanobot/bus/message_bus.rb', line 26

def publish_inbound(message)
  raise ArgumentError, 'Message must be an InboundMessage' unless message.is_a?(InboundMessage)

  @inbound_queue.push(message)
  @logger.debug "Published inbound message from #{message.channel}:#{message.chat_id}"
end

#publish_outbound(message) ⇒ Object

Publish an outbound message (from agent to channels)

Parameters:

Raises:

  • (ArgumentError)


52
53
54
55
56
57
# File 'lib/nanobot/bus/message_bus.rb', line 52

def publish_outbound(message)
  raise ArgumentError, 'Message must be an OutboundMessage' unless message.is_a?(OutboundMessage)

  @outbound_queue.push(message)
  @logger.debug "Published outbound message to #{message.channel}:#{message.chat_id}"
end

#queue_sizes ⇒ Hash{Symbol => Integer}

Get queue sizes for monitoring

Returns:

  • (Hash{Symbol => Integer}) —

    :inbound and :outbound counts



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

Returns:

  • (Boolean)


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

Parameters:

  • channel (String) —

    channel name

  • callback (Proc) —

    callback to invoke with OutboundMessage

Raises:

  • (ArgumentError)


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