Class: Legion::Extensions::Llm::Transport::Messages::FleetRequest

Inherits:
Transport::Message
  • Object
show all
Includes:
Fleet::EnvelopeValidation, Fleet::PublishSafety
Defined in:
lib/legion/extensions/llm/transport/messages/fleet_request.rb

Overview

Strict protocol-v3 request envelope for outbound fleet work. E2: required fields come from the one Fleet::Protocol::REQUIRED_FIELDS list.

Constant Summary collapse

PRIORITY_MAP =
{ critical: 9, high: 7, normal: 5, low: 2 }.freeze
DEFAULT_PUBLISH_OPTIONS =
{
  mandatory: true,
  publisher_confirm: true,
  spool: false,
  return_result: true
}.freeze

Instance Method Summary collapse

Instance Method Details

#app_idObject



31
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 31

def app_id = @options[:app_id] || 'lex-llm'

#correlation_idObject



33
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 33

def correlation_id = @options[:correlation_id]

#encrypt?Boolean

Returns:

  • (Boolean)


30
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 30

def encrypt? = Fleet::Settings.value(:fleet, :compliance, :encrypt_fleet, default: true) == true

#exchangeObject



28
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 28

def exchange = Exchanges::Fleet

#expirationObject



44
45
46
47
48
49
50
51
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 44

def expiration
  ttl = @options[:ttl] || @options[:timeout_seconds]
  return super unless ttl

  (Float(ttl) * 1000).ceil.to_s
rescue ArgumentError, TypeError
  super
end

#messageObject



82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 82

def message
  super.merge(
    protocol_version: @options[:protocol_version],
    request_id: @options[:request_id],
    correlation_id: correlation_id,
    idempotency_key: @options[:idempotency_key],
    operation: @options[:operation],
    provider: @options[:provider],
    provider_instance: @options[:provider_instance],
    model: @options[:model],
    params: @options[:params] || {},
    reply_to: reply_to,
    message_context: @options[:message_context],
    caller: @options[:caller],
    trace_context: @options[:trace_context],
    signed_token: @options[:signed_token],
    timeout_seconds: @options[:timeout_seconds],
    expires_at: @options[:expires_at],
    execution_contract: @options[:execution_contract],
    offering_id: @options[:offering_id]
  ).compact
end

#message_idObject



34
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 34

def message_id = @options[:message_id] ||= "llm_fleet_req_#{SecureRandom.uuid}"

#priorityObject



36
37
38
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 36

def priority
  PRIORITY_MAP.fetch(@options[:priority].to_sym, 5) if @options[:priority]
end

#publish(options = nil) ⇒ Object



53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 53

def publish(options = nil)
  raise unless @valid

  requested_options = DEFAULT_PUBLISH_OPTIONS.merge(@options).merge(options || {})
  return_result = return_publish_result?(requested_options)
  publish_options = request_publish_options(requested_options)
  validate_payload_size
  exchange_dest = fleet_exchange
  return_state = {}
  install_return_listener(exchange_dest, requested_options, return_state)
  prepare_publisher_confirms(exchange_dest, requested_options)
  exchange_dest.publish(encode_message, **publish_options)
  return nil unless return_result

  publish_result(exchange_dest, requested_options.merge(publish_options), return_state)
rescue Bunny::ConnectionClosedError, Bunny::ChannelAlreadyClosed, Bunny::ChannelError,
       Bunny::NetworkErrorWrapper, IOError, Timeout::Error => e
  handle_exception(e, level: :warn, handled: true, operation: 'llm.fleet.request.publish')
  publish_failure_result(:failed, e, publish_options || requested_options || @options)
end

#reply_toObject



32
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 32

def reply_to = @options[:reply_to]

#routing_keyObject



40
41
42
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 40

def routing_key
  @options[:routing_key] || raise(ArgumentError, 'routing_key is required')
end

#typeObject



29
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 29

def type = Fleet::Protocol::REQUEST_TYPE

#validateObject



74
75
76
77
78
79
80
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 74

def validate
  reject_legacy_options!
  require_option!(:routing_key)
  Fleet::Protocol::REQUIRED_FIELDS.each { |key| require_option!(key) }
  require_protocol_version!
  @valid = true
end