Class: Legion::Extensions::Llm::Transport::Messages::FleetRequest
- Inherits:
-
Transport::Message
- Object
- Transport::Message
- Legion::Extensions::Llm::Transport::Messages::FleetRequest
- 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
- #app_id ⇒ Object
- #correlation_id ⇒ Object
- #encrypt? ⇒ Boolean
- #exchange ⇒ Object
- #expiration ⇒ Object
- #message ⇒ Object
- #message_id ⇒ Object
- #priority ⇒ Object
- #publish(options = nil) ⇒ Object
- #reply_to ⇒ Object
- #routing_key ⇒ Object
- #type ⇒ Object
- #validate ⇒ Object
Instance Method Details
#app_id ⇒ Object
31 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 31 def app_id = @options[:app_id] || 'lex-llm' |
#correlation_id ⇒ Object
33 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 33 def correlation_id = @options[:correlation_id] |
#encrypt? ⇒ 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 |
#exchange ⇒ Object
28 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 28 def exchange = Exchanges::Fleet |
#expiration ⇒ Object
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 |
#message ⇒ Object
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 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_id ⇒ Object
34 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 34 def = @options[:message_id] ||= "llm_fleet_req_#{SecureRandom.uuid}" |
#priority ⇒ Object
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( = nil) raise unless @valid = DEFAULT_PUBLISH_OPTIONS.merge(@options).merge( || {}) return_result = return_publish_result?() = () validate_payload_size exchange_dest = fleet_exchange return_state = {} install_return_listener(exchange_dest, , return_state) prepare_publisher_confirms(exchange_dest, ) exchange_dest.publish(, **) return nil unless return_result publish_result(exchange_dest, .merge(), 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, || || @options) end |
#reply_to ⇒ Object
32 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 32 def reply_to = @options[:reply_to] |
#routing_key ⇒ Object
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 |
#type ⇒ Object
29 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 29 def type = Fleet::Protocol::REQUEST_TYPE |
#validate ⇒ Object
74 75 76 77 78 79 80 |
# File 'lib/legion/extensions/llm/transport/messages/fleet_request.rb', line 74 def validate require_option!(:routing_key) Fleet::Protocol::REQUIRED_FIELDS.each { |key| require_option!(key) } require_protocol_version! @valid = true end |