Module: Legion::Extensions::Llm::Transport::FleetLane

Defined in:
lib/legion/extensions/llm/transport/fleet_lane.rb

Overview

Shared RabbitMQ live-work lane construction for provider fleet workers. The queue defaults live in ONE home — Llm.default_settings (10 U10: no inline literal defaults).

Class Method Summary collapse

Class Method Details

.build_queue_class(queue_name:, exchange_class:, routing_key: queue_name, base_queue_class: nil, settings: {}) ⇒ Object



39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
# File 'lib/legion/extensions/llm/transport/fleet_lane.rb', line 39

def build_queue_class(queue_name:, exchange_class:, routing_key: queue_name, base_queue_class: nil,
                      settings: {})
  parent = base_queue_class || legion_queue_class
  unless parent
    raise ArgumentError,
          'base_queue_class is required when Legion::Transport::Queue is not loaded'
  end

  options = queue_options(settings)
  Class.new(parent) do
    define_method(:queue_name) { queue_name }
    define_method(:queue_options) { options }
    define_method(:dlx_enabled) { false }
    define_method(:initialize) do
      super()
      bind(exchange_class.new, routing_key: routing_key)
    end
  end
end

.legion_queue_classObject



59
60
61
62
63
# File 'lib/legion/extensions/llm/transport/fleet_lane.rb', line 59

def legion_queue_class
  return nil unless defined?(::Legion::Transport::Queue)

  ::Legion::Transport::Queue
end

.queue_arguments(config) ⇒ Object



26
27
28
29
30
31
32
33
34
35
36
37
# File 'lib/legion/extensions/llm/transport/fleet_lane.rb', line 26

def queue_arguments(config)
  {
    'x-queue-type' => 'quorum',
    'x-queue-leader-locator' => 'balanced',
    'x-expires' => config.fetch(:queue_expires_ms),
    'x-message-ttl' => config.fetch(:message_ttl_ms),
    'x-overflow' => 'reject-publish',
    'x-max-length' => config.fetch(:queue_max_length),
    'x-delivery-limit' => config.fetch(:delivery_limit),
    'x-consumer-timeout' => config.fetch(:consumer_ack_timeout_ms)
  }
end

.queue_defaultsObject



13
14
15
# File 'lib/legion/extensions/llm/transport/fleet_lane.rb', line 13

def queue_defaults
  Legion::Extensions::Llm.default_settings.dig(:fleet, :consumer) || {}
end

.queue_options(settings = {}) ⇒ Object



17
18
19
20
21
22
23
24
# File 'lib/legion/extensions/llm/transport/fleet_lane.rb', line 17

def queue_options(settings = {})
  config = queue_defaults.merge((settings || {}).compact.transform_keys(&:to_sym))
  {
    durable: true,
    auto_delete: false,
    arguments: queue_arguments(config)
  }
end